1use super::*;
2use chrono::{SecondsFormat, Utc};
3use std::sync::Weak;
4
5#[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#[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 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 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 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#[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#[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 #[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}