durable_workflow/
cancellation_context.rs

1use super::*;
2use chrono::{SecondsFormat, Utc};
3use std::sync::Weak;
4
5/// One immutable local request in the ordered cancellation lineage.
6#[derive(Clone, Debug, PartialEq, Eq)]
7pub struct CancellationLineage {
8    request_id: String,
9    workflow_instance_id: String,
10    workflow_run_id: String,
11}
12
13impl CancellationLineage {
14    pub fn request_id(&self) -> &str {
15        &self.request_id
16    }
17
18    pub fn workflow_instance_id(&self) -> &str {
19        &self.workflow_instance_id
20    }
21
22    pub fn workflow_run_id(&self) -> &str {
23        &self.workflow_run_id
24    }
25
26    fn to_value(&self) -> Value {
27        json!({
28            "request_id": self.request_id,
29            "workflow_instance_id": self.workflow_instance_id,
30            "workflow_run_id": self.workflow_run_id,
31        })
32    }
33}
34
35/// Original cancellation metadata restored from canonical request history.
36///
37/// Fields and nested metadata have read-only accessors. Descendants retain
38/// the original root identity, request time and cleanup budget.
39#[derive(Clone, Debug)]
40pub struct CancellationContext {
41    request_id: String,
42    root_request_id: String,
43    root_workflow_instance_id: String,
44    root_workflow_run_id: String,
45    parent_request_id: Option<String>,
46    reason: Option<String>,
47    requester: BTreeMap<String, String>,
48    source: String,
49    requested_at: DateTime<Utc>,
50    cleanup_deadline_at: DateTime<Utc>,
51    lineage: Vec<CancellationLineage>,
52    scope_origin: Option<Box<ScopedCancellationContext>>,
53    replay: Option<Weak<Mutex<WorkflowState>>>,
54}
55
56impl PartialEq for CancellationContext {
57    fn eq(&self, other: &Self) -> bool {
58        self.request_id == other.request_id
59            && self.root_request_id == other.root_request_id
60            && self.root_workflow_instance_id == other.root_workflow_instance_id
61            && self.root_workflow_run_id == other.root_workflow_run_id
62            && self.parent_request_id == other.parent_request_id
63            && self.reason == other.reason
64            && self.requester == other.requester
65            && self.source == other.source
66            && self.requested_at == other.requested_at
67            && self.cleanup_deadline_at == other.cleanup_deadline_at
68            && self.lineage == other.lineage
69            && self.scope_origin == other.scope_origin
70    }
71}
72
73impl Eq for CancellationContext {}
74
75fn invalid_context(message: &str) -> Error {
76    Error::InvalidCooperativeCancellation(message.to_owned())
77}
78
79fn context_text<'a>(value: &'a Value, key: &str) -> Result<&'a str> {
80    value[key]
81        .as_str()
82        .filter(|text| !text.trim().is_empty())
83        .ok_or_else(|| invalid_context("cancellation context identity must be a non-empty string"))
84}
85
86fn nullable_text(value: &Value, key: &str, allow_empty: bool) -> Result<Option<String>> {
87    match value.get(key) {
88        None | Some(Value::Null) => Ok(None),
89        Some(Value::String(text)) if allow_empty || !text.is_empty() => Ok(Some(text.clone())),
90        _ => Err(invalid_context(
91            "cancellation parent identity or reason is invalid",
92        )),
93    }
94}
95
96impl CancellationContext {
97    /// Read legacy v1 and candidate scoped v2 without granting a new cleanup budget.
98    pub fn from_value(value: &Value) -> Result<Self> {
99        let scope_origin = match value["schema"].as_str() {
100            Some("durable-workflow.cancellation-context/v2") => Some(Box::new(
101                ScopedCancellationContext::from_value(&value["scope_origin"])?,
102            )),
103            Some("durable-workflow.cancellation-context/v1") => {
104                if value.get("scope_origin").is_some()
105                    || value.get("scope_authority_deadline_at").is_some()
106                {
107                    return Err(invalid_context(
108                        "legacy cancellation context cannot discard a scoped origin",
109                    ));
110                }
111                None
112            }
113            _ => return Err(invalid_context("unsupported cancellation context schema")),
114        };
115        let requester = value["requester"]
116            .as_object()
117            .filter(|requester| !requester.is_empty())
118            .ok_or_else(|| invalid_context("cancellation requester must identify its caller"))?;
119        let mut normalized_requester = BTreeMap::new();
120        for (key, value) in requester {
121            let text = value
122                .as_str()
123                .filter(|text| !text.is_empty())
124                .ok_or_else(|| {
125                    invalid_context("cancellation requester contains unsupported metadata")
126                })?;
127            if !matches!(key.as_str(), "type" | "id" | "label") {
128                return Err(invalid_context(
129                    "cancellation requester contains unsupported metadata",
130                ));
131            }
132            normalized_requester.insert(key.clone(), text.to_owned());
133        }
134        let lineage = value["lineage"]
135            .as_array()
136            .filter(|lineage| !lineage.is_empty())
137            .ok_or_else(|| invalid_context("cancellation lineage must contain the root request"))?;
138        let mut normalized = Vec::new();
139        for entry in lineage {
140            normalized.push(CancellationLineage {
141                request_id: context_text(entry, "request_id")?.to_owned(),
142                workflow_instance_id: context_text(entry, "workflow_instance_id")?.to_owned(),
143                workflow_run_id: context_text(entry, "workflow_run_id")?.to_owned(),
144            });
145        }
146        let request_id = context_text(value, "request_id")?.to_owned();
147        let root_request_id = context_text(value, "root_request_id")?.to_owned();
148        let root_workflow_instance_id =
149            context_text(value, "root_workflow_instance_id")?.to_owned();
150        let root_workflow_run_id = context_text(value, "root_workflow_run_id")?.to_owned();
151        let parent_request_id = nullable_text(value, "parent_request_id", false)?;
152        let reason = nullable_text(value, "reason", true)?;
153        let requests: BTreeSet<_> = normalized.iter().map(|entry| &entry.request_id).collect();
154        let runs: BTreeSet<_> = normalized
155            .iter()
156            .map(|entry| &entry.workflow_run_id)
157            .collect();
158        if requests.len() != normalized.len() || runs.len() != normalized.len() {
159            return Err(invalid_context(
160                "cancellation lineage cannot contain a cycle",
161            ));
162        }
163        let expected_parent = normalized
164            .iter()
165            .rev()
166            .nth(1)
167            .map(|entry| &entry.request_id);
168        if normalized[0].request_id != root_request_id
169            || normalized[0].workflow_instance_id != root_workflow_instance_id
170            || normalized[0].workflow_run_id != root_workflow_run_id
171            || normalized.last().unwrap().request_id != request_id
172            || (scope_origin.is_none() && parent_request_id.as_ref() != expected_parent)
173        {
174            return Err(invalid_context(
175                "cancellation lineage does not match its request identities",
176            ));
177        }
178        let requested_at = DateTime::parse_from_rfc3339(context_text(value, "requested_at")?)
179            .map_err(|_| invalid_context("cancellation request timestamp is invalid"))?
180            .with_timezone(&Utc);
181        let cleanup_deadline_at =
182            DateTime::parse_from_rfc3339(context_text(value, "cleanup_deadline_at")?)
183                .map_err(|_| invalid_context("cancellation deadline timestamp is invalid"))?
184                .with_timezone(&Utc);
185        if cleanup_deadline_at <= requested_at {
186            return Err(invalid_context(
187                "cancellation deadline must follow the original request",
188            ));
189        }
190        let context = Self {
191            request_id,
192            root_request_id,
193            root_workflow_instance_id,
194            root_workflow_run_id,
195            parent_request_id,
196            reason,
197            requester: normalized_requester,
198            source: context_text(value, "source")?.to_owned(),
199            requested_at,
200            cleanup_deadline_at,
201            lineage: normalized,
202            scope_origin,
203            replay: None,
204        };
205        context.assert_scope_origin(value)?;
206        Ok(context)
207    }
208
209    fn assert_scope_origin(&self, value: &Value) -> Result<()> {
210        let Some(origin) = self.scope_origin.as_deref() else {
211            return Ok(());
212        };
213        let root = &origin.root_context;
214        let last = self.lineage.last().unwrap();
215        let authority =
216            DateTime::parse_from_rfc3339(context_text(value, "scope_authority_deadline_at")?)
217                .map_err(|_| invalid_context("scoped cancellation authority timestamp is invalid"))?
218                .with_timezone(&Utc);
219        if self.parent_request_id.as_deref() != Some(origin.request_id())
220            || self.root_request_id != root.root_request_id
221            || self.root_workflow_instance_id != root.root_workflow_instance_id
222            || self.root_workflow_run_id != root.root_workflow_run_id
223            || self.reason != root.reason
224            || self.requester != root.requester
225            || self.source != root.source
226            || self.requested_at != root.requested_at
227            || self.cleanup_deadline_at != authority
228            || self.cleanup_deadline_at > origin.deadline()
229            || origin.lineage.iter().any(|entry| {
230                entry.request_id == last.request_id || entry.workflow_run_id == last.workflow_run_id
231            })
232        {
233            return Err(invalid_context(
234                "run cancellation does not preserve its original scope context",
235            ));
236        }
237        let mut expected = vec![root.lineage[0].clone()];
238        for entry in &origin.lineage {
239            if entry.workflow_run_id == root.root_workflow_run_id {
240                continue;
241            }
242            let address = CancellationLineage {
243                request_id: entry.request_id.clone(),
244                workflow_instance_id: entry.workflow_instance_id.clone(),
245                workflow_run_id: entry.workflow_run_id.clone(),
246            };
247            if expected.last().unwrap().workflow_run_id == entry.workflow_run_id {
248                *expected.last_mut().unwrap() = address;
249            } else {
250                expected.push(address);
251            }
252        }
253        expected.push(last.clone());
254        if self.lineage != expected {
255            return Err(invalid_context(
256                "run cancellation lineage discards or replaces its original scope ancestry",
257            ));
258        }
259        Ok(())
260    }
261
262    pub fn request_id(&self) -> &str {
263        &self.request_id
264    }
265    pub fn root_request_id(&self) -> &str {
266        &self.root_request_id
267    }
268    pub fn root_workflow_instance_id(&self) -> &str {
269        &self.root_workflow_instance_id
270    }
271    pub fn root_workflow_run_id(&self) -> &str {
272        &self.root_workflow_run_id
273    }
274    pub fn parent_request_id(&self) -> Option<&str> {
275        self.parent_request_id.as_deref()
276    }
277    pub fn reason(&self) -> Option<&str> {
278        self.reason.as_deref()
279    }
280    pub fn requester(&self) -> &BTreeMap<String, String> {
281        &self.requester
282    }
283    pub fn source(&self) -> &str {
284        &self.source
285    }
286    pub fn requested_at(&self) -> DateTime<Utc> {
287        self.requested_at
288    }
289    pub fn deadline(&self) -> DateTime<Utc> {
290        self.cleanup_deadline_at
291    }
292    pub fn scope_origin(&self) -> Option<&ScopedCancellationContext> {
293        self.scope_origin.as_deref()
294    }
295
296    /// Remaining cleanup budget at the last blocking result consumed by this replay.
297    ///
298    /// Never reads host time. Detached metadata, an ended replay, and a missing or
299    /// invalid committed timestamp return an explicit error. Expiry returns zero.
300    pub fn remaining(&self) -> Result<Duration> {
301        let state = self
302            .replay
303            .as_ref()
304            .and_then(Weak::upgrade)
305            .filter(cancellation_replay_clock::is_active)
306            .ok_or_else(|| {
307                invalid_context("remaining() is available only in its active workflow replay")
308            })?;
309        let time = state
310            .lock()
311            .map_err(|_| Error::WorkflowStatePoisoned)?
312            .cancellation_time()?;
313        Ok((self.cleanup_deadline_at - time)
314            .to_std()
315            .unwrap_or(Duration::ZERO))
316    }
317
318    pub(super) fn with_replay(mut self, replay: Option<Weak<Mutex<WorkflowState>>>) -> Self {
319        self.replay = replay;
320        self
321    }
322    pub fn lineage(&self) -> &[CancellationLineage] {
323        &self.lineage
324    }
325
326    /// Detached metadata in the portable context schema.
327    pub fn to_value(&self) -> Value {
328        let mut value = json!({
329            "schema": if self.scope_origin.is_none() { "durable-workflow.cancellation-context/v1" }
330                else { "durable-workflow.cancellation-context/v2" },
331            "request_id": self.request_id, "root_request_id": self.root_request_id,
332            "root_workflow_instance_id": self.root_workflow_instance_id,
333            "root_workflow_run_id": self.root_workflow_run_id,
334            "parent_request_id": self.parent_request_id, "reason": self.reason,
335            "requester": self.requester, "source": self.source,
336            "requested_at": self.requested_at.to_rfc3339_opts(SecondsFormat::Micros, true),
337            "cleanup_deadline_at": self.cleanup_deadline_at.to_rfc3339_opts(SecondsFormat::Micros, true),
338            "lineage": self.lineage.iter().map(CancellationLineage::to_value).collect::<Vec<_>>(),
339        });
340        if let Some(origin) = &self.scope_origin {
341            value["scope_origin"] = origin.to_value();
342            value["scope_authority_deadline_at"] = value["cleanup_deadline_at"].clone();
343        }
344        value
345    }
346}
347
348/// One immutable scope address and its bounded cleanup deadline.
349#[derive(Clone, Debug, PartialEq, Eq)]
350pub struct ScopedCancellationLineage {
351    request_id: String,
352    workflow_instance_id: String,
353    workflow_run_id: String,
354    scope_id: String,
355    cleanup_deadline_at: DateTime<Utc>,
356}
357
358impl ScopedCancellationLineage {
359    pub fn request_id(&self) -> &str {
360        &self.request_id
361    }
362    pub fn workflow_instance_id(&self) -> &str {
363        &self.workflow_instance_id
364    }
365    pub fn workflow_run_id(&self) -> &str {
366        &self.workflow_run_id
367    }
368    pub fn scope_id(&self) -> &str {
369        &self.scope_id
370    }
371    pub fn deadline(&self) -> DateTime<Utc> {
372        self.cleanup_deadline_at
373    }
374
375    fn to_value(&self) -> Value {
376        json!({
377            "request_id": self.request_id, "workflow_instance_id": self.workflow_instance_id,
378            "workflow_run_id": self.workflow_run_id, "scope_id": self.scope_id,
379            "cleanup_deadline_at": self.cleanup_deadline_at.to_rfc3339_opts(SecondsFormat::Micros, true),
380        })
381    }
382}
383
384/// Original scope ancestry carried by a candidate cooperative child request.
385/// Reading this immutable metadata does not authorize entering a scope body.
386#[derive(Clone, Debug, PartialEq, Eq)]
387pub struct ScopedCancellationContext {
388    root_context: CancellationContext,
389    lineage: Vec<ScopedCancellationLineage>,
390}
391
392fn assert_context_keys(value: &Value, keys: &[&str]) -> Result<()> {
393    let object = value
394        .as_object()
395        .ok_or_else(|| invalid_context("scoped context must be an object"))?;
396    if object.len() != keys.len() || keys.iter().any(|key| !object.contains_key(*key)) {
397        return Err(invalid_context(
398            "scoped cancellation context has missing or unsupported fields",
399        ));
400    }
401    Ok(())
402}
403
404impl ScopedCancellationContext {
405    /// Remaining original scope authority at the last consumed durable result.
406    #[doc(hidden)]
407    pub fn remaining(&self) -> Result<Duration> {
408        let state = self
409            .root_context
410            .replay
411            .as_ref()
412            .and_then(Weak::upgrade)
413            .filter(cancellation_replay_clock::is_active)
414            .ok_or_else(|| {
415                invalid_context("remaining() is available only in its active scope replay")
416            })?;
417        let state = state.lock().map_err(|_| Error::WorkflowStatePoisoned)?;
418        let (_, ceiling) = state
419            .scope_delivery
420            .as_ref()
421            .and_then(|replay| replay.contexts.get(self.scope_id()))
422            .filter(|(original, _)| original == self)
423            .ok_or_else(|| {
424                invalid_context("scope remaining() requires its original consumed delivery")
425            })?;
426        Ok((self.deadline().min(*ceiling) - state.cancellation_time()?)
427            .to_std()
428            .unwrap_or(Duration::ZERO))
429    }
430
431    pub(super) fn with_replay(mut self, replay: Option<Weak<Mutex<WorkflowState>>>) -> Self {
432        self.root_context = self.root_context.with_replay(replay);
433        self
434    }
435
436    pub(super) fn from_run_context(context: &CancellationContext) -> Result<Self> {
437        let mut root = context.to_value();
438        let mut lineage = Vec::new();
439        if let Some(origin) = context.scope_origin() {
440            root = origin.root_context().to_value();
441            lineage.extend(
442                origin
443                    .lineage()
444                    .iter()
445                    .map(ScopedCancellationLineage::to_value),
446            );
447            let local = context.lineage().last().unwrap();
448            lineage.push(json!({"request_id":context.request_id(),
449                "workflow_instance_id":local.workflow_instance_id(),
450                "workflow_run_id":local.workflow_run_id(), "scope_id":"root",
451                "cleanup_deadline_at":context.deadline().to_rfc3339_opts(SecondsFormat::Micros, true)}));
452        } else {
453            root["request_id"] = json!(context.root_request_id());
454            root["parent_request_id"] = Value::Null;
455            root["lineage"] = json!([context.lineage()[0].to_value()]);
456            lineage.extend(context.lineage().iter().map(|entry| json!({
457                "request_id":entry.request_id(), "workflow_instance_id":entry.workflow_instance_id(),
458                "workflow_run_id":entry.workflow_run_id(), "scope_id":"root",
459                "cleanup_deadline_at":context.deadline().to_rfc3339_opts(SecondsFormat::Micros, true)})));
460        }
461        Self::from_value(
462            &json!({"schema":"durable-workflow.scoped-cancellation-context/v1",
463            "root_context":root, "lineage":lineage}),
464        )
465    }
466
467    pub fn from_value(value: &Value) -> Result<Self> {
468        assert_context_keys(value, &["schema", "root_context", "lineage"])?;
469        if value["schema"] != "durable-workflow.scoped-cancellation-context/v1"
470            || value["root_context"]["schema"] != "durable-workflow.cancellation-context/v1"
471        {
472            return Err(invalid_context(
473                "scoped cancellation requires its original root context",
474            ));
475        }
476        let root = CancellationContext::from_value(&value["root_context"])?;
477        let root_value = root.to_value();
478        let root_keys: Vec<_> = root_value
479            .as_object()
480            .unwrap()
481            .keys()
482            .map(String::as_str)
483            .collect();
484        assert_context_keys(&value["root_context"], &root_keys)?;
485        if root.request_id != root.root_request_id
486            || root.parent_request_id.is_some()
487            || root.lineage.len() != 1
488        {
489            return Err(invalid_context(
490                "scoped cancellation root must contain the original root request",
491            ));
492        }
493        let lineage = value["lineage"]
494            .as_array()
495            .filter(|rows| !rows.is_empty())
496            .ok_or_else(|| {
497                invalid_context("scoped cancellation lineage must contain its root address")
498            })?;
499        let mut normalized = Vec::new();
500        let mut requests = BTreeSet::new();
501        let mut addresses = BTreeSet::new();
502        let mut instances_by_run = BTreeMap::new();
503        let mut last_run = String::new();
504        let mut deadline = root.deadline();
505        for (index, entry) in lineage.iter().enumerate() {
506            assert_context_keys(
507                entry,
508                &[
509                    "request_id",
510                    "workflow_instance_id",
511                    "workflow_run_id",
512                    "scope_id",
513                    "cleanup_deadline_at",
514                ],
515            )?;
516            let address = ScopedCancellationLineage {
517                request_id: context_text(entry, "request_id")?.to_owned(),
518                workflow_instance_id: context_text(entry, "workflow_instance_id")?.to_owned(),
519                workflow_run_id: context_text(entry, "workflow_run_id")?.to_owned(),
520                scope_id: context_text(entry, "scope_id")?.to_owned(),
521                cleanup_deadline_at: DateTime::parse_from_rfc3339(context_text(
522                    entry,
523                    "cleanup_deadline_at",
524                )?)
525                .map_err(|_| invalid_context("scoped cancellation deadline timestamp is invalid"))?
526                .with_timezone(&Utc),
527            };
528            if index == 0
529                && (address.request_id != root.request_id
530                    || address.workflow_instance_id != root.root_workflow_instance_id
531                    || address.workflow_run_id != root.root_workflow_run_id
532                    || address.cleanup_deadline_at != root.deadline())
533            {
534                return Err(invalid_context(
535                    "scoped cancellation root address does not match its request",
536                ));
537            }
538            if !requests.insert(address.request_id.clone())
539                || !addresses.insert((address.workflow_run_id.clone(), address.scope_id.clone()))
540            {
541                return Err(invalid_context(
542                    "scoped cancellation cannot repeat a request or address",
543                ));
544            }
545            if let Some(instance) = instances_by_run.get(&address.workflow_run_id) {
546                if instance != &address.workflow_instance_id || last_run != address.workflow_run_id
547                {
548                    return Err(invalid_context(
549                        "scoped cancellation cannot reenter or reassign an earlier run",
550                    ));
551                }
552            }
553            if address.cleanup_deadline_at <= root.requested_at
554                || address.cleanup_deadline_at > deadline
555            {
556                return Err(invalid_context(
557                    "scoped cancellation cannot extend a descendant budget",
558                ));
559            }
560            instances_by_run.insert(
561                address.workflow_run_id.clone(),
562                address.workflow_instance_id.clone(),
563            );
564            last_run = address.workflow_run_id.clone();
565            deadline = address.cleanup_deadline_at;
566            normalized.push(address);
567        }
568        Ok(Self {
569            root_context: root,
570            lineage: normalized,
571        })
572    }
573
574    pub fn root_context(&self) -> &CancellationContext {
575        &self.root_context
576    }
577    pub fn lineage(&self) -> &[ScopedCancellationLineage] {
578        &self.lineage
579    }
580    pub fn request_id(&self) -> &str {
581        &self.lineage.last().unwrap().request_id
582    }
583    pub fn parent_request_id(&self) -> Option<&str> {
584        self.lineage
585            .iter()
586            .rev()
587            .nth(1)
588            .map(|entry| entry.request_id.as_str())
589    }
590    pub fn workflow_instance_id(&self) -> &str {
591        &self.lineage.last().unwrap().workflow_instance_id
592    }
593    pub fn workflow_run_id(&self) -> &str {
594        &self.lineage.last().unwrap().workflow_run_id
595    }
596    pub fn scope_id(&self) -> &str {
597        &self.lineage.last().unwrap().scope_id
598    }
599    pub fn root_scope_id(&self) -> &str {
600        &self.lineage[0].scope_id
601    }
602    pub fn requested_at(&self) -> DateTime<Utc> {
603        self.root_context.requested_at()
604    }
605    pub fn root_deadline(&self) -> DateTime<Utc> {
606        self.root_context.deadline()
607    }
608    pub fn deadline(&self) -> DateTime<Utc> {
609        self.lineage.last().unwrap().cleanup_deadline_at
610    }
611    pub fn to_value(&self) -> Value {
612        json!({
613            "schema": "durable-workflow.scoped-cancellation-context/v1",
614            "root_context": self.root_context.to_value(),
615            "lineage": self.lineage.iter().map(ScopedCancellationLineage::to_value).collect::<Vec<_>>(),
616        })
617    }
618}