1use std::collections::HashMap;
4use std::collections::hash_map::Entry;
5use std::sync::{Arc, Mutex, MutexGuard};
6
7use aion_core::{ActivityId, RunId, WorkflowId};
8use sha2::{Digest, Sha256};
9use uuid::Uuid;
10
11use crate::error::{CompletionRejectionReason, ServerError};
12
13type ExecutionKey = (WorkflowId, ActivityId);
14
15#[derive(Clone, Debug, Eq, PartialEq)]
17pub struct CompletionToken(String);
18
19impl CompletionToken {
20 pub fn from_wire(
26 workflow_id: &WorkflowId,
27 activity_id: &ActivityId,
28 value: String,
29 ) -> Result<Self, ServerError> {
30 if value.is_empty() {
31 return Err(rejection(
32 workflow_id,
33 activity_id,
34 CompletionRejectionReason::MissingCompletionToken,
35 ));
36 }
37 Ok(Self(value))
38 }
39
40 #[must_use]
42 pub fn as_str(&self) -> &str {
43 &self.0
44 }
45
46 #[cfg(test)]
48 #[must_use]
49 pub(crate) fn for_test() -> Self {
50 Self("test-generation".to_owned())
51 }
52}
53
54#[derive(Clone, Debug, Eq, PartialEq)]
64struct SiteGeneration {
65 idempotency_key: String,
71 attempt: u32,
74 outstanding: Vec<CompletionToken>,
77}
78
79#[derive(Clone, Debug, Eq, PartialEq)]
87pub struct AcceptedGeneration(SiteGeneration);
88
89impl AcceptedGeneration {
90 #[must_use]
92 pub const fn attempt(&self) -> u32 {
93 self.0.attempt
94 }
95}
96
97#[derive(Clone, Debug, Default)]
103pub struct CompletionFences {
104 current: Arc<Mutex<HashMap<ExecutionKey, SiteGeneration>>>,
105}
106
107impl CompletionFences {
108 pub fn issue(
131 &self,
132 workflow_id: &WorkflowId,
133 run_id: &RunId,
134 activity_id: &ActivityId,
135 attempt: u32,
136 ) -> Result<CompletionToken, ServerError> {
137 let token = CompletionToken(Uuid::new_v4().to_string());
138 let issued_key = idempotency_key(workflow_id, run_id, activity_id);
139 let mut state = self.state()?;
140 match state.entry((workflow_id.clone(), activity_id.clone())) {
141 Entry::Vacant(slot) => {
142 slot.insert(SiteGeneration {
143 idempotency_key: issued_key,
144 attempt,
145 outstanding: vec![token.clone()],
146 });
147 }
148 Entry::Occupied(mut slot) => {
149 let current = slot.get_mut();
150 if current.idempotency_key != issued_key || attempt > current.attempt {
151 *current = SiteGeneration {
152 idempotency_key: issued_key,
153 attempt,
154 outstanding: vec![token.clone()],
155 };
156 } else if attempt == current.attempt {
157 current.outstanding.push(token.clone());
158 } else {
159 tracing::warn!(
160 workflow_id = %workflow_id,
161 activity_id = %activity_id,
162 superseded_attempt = attempt,
163 current_attempt = current.attempt,
164 "re-dispatch of an already superseded attempt registers no generation; \
165 its completion will be refused as a stale generation"
166 );
167 }
168 }
169 }
170 Ok(token)
171 }
172
173 pub fn accept(
187 &self,
188 workflow_id: &WorkflowId,
189 activity_id: &ActivityId,
190 submitted: &CompletionToken,
191 ) -> Result<AcceptedGeneration, ServerError> {
192 let mut state = self.state()?;
193 let Entry::Occupied(slot) = state.entry((workflow_id.clone(), activity_id.clone())) else {
194 return Err(rejection(
195 workflow_id,
196 activity_id,
197 CompletionRejectionReason::NoCurrentGeneration,
198 ));
199 };
200 if !slot.get().outstanding.contains(submitted) {
201 return Err(rejection(
202 workflow_id,
203 activity_id,
204 CompletionRejectionReason::StaleGeneration,
205 ));
206 }
207 Ok(AcceptedGeneration(slot.remove()))
208 }
209
210 pub fn restore_if_absent(
219 &self,
220 workflow_id: &WorkflowId,
221 activity_id: &ActivityId,
222 accepted: &AcceptedGeneration,
223 ) -> Result<(), ServerError> {
224 self.state()?
225 .entry((workflow_id.clone(), activity_id.clone()))
226 .or_insert_with(|| accepted.0.clone());
227 Ok(())
228 }
229
230 pub fn revoke(
243 &self,
244 workflow_id: &WorkflowId,
245 activity_id: &ActivityId,
246 token: &CompletionToken,
247 ) -> Result<(), ServerError> {
248 let mut state = self.state()?;
249 let Entry::Occupied(mut slot) = state.entry((workflow_id.clone(), activity_id.clone()))
250 else {
251 return Ok(());
252 };
253 slot.get_mut().outstanding.retain(|issued| issued != token);
254 if slot.get().outstanding.is_empty() {
255 slot.remove();
256 }
257 Ok(())
258 }
259
260 pub fn revoke_current(
269 &self,
270 workflow_id: &WorkflowId,
271 activity_id: &ActivityId,
272 ) -> Result<(), ServerError> {
273 self.state()?
274 .remove(&(workflow_id.clone(), activity_id.clone()));
275 Ok(())
276 }
277
278 fn state(&self) -> Result<MutexGuard<'_, HashMap<ExecutionKey, SiteGeneration>>, ServerError> {
279 self.current
280 .lock()
281 .map_err(|_| ServerError::lock_poisoned("activity completion fences"))
282 }
283}
284
285#[must_use]
291pub fn idempotency_key(
292 workflow_id: &WorkflowId,
293 run_id: &RunId,
294 activity_id: &ActivityId,
295) -> String {
296 let mut hasher = Sha256::new();
297 hasher.update(b"aion.activity.idempotency.v1\0");
298 hasher.update(workflow_id.as_uuid().as_bytes());
299 hasher.update(run_id.as_uuid().as_bytes());
300 hasher.update(activity_id.sequence_position().to_be_bytes());
301 encode_hex(&hasher.finalize())
302}
303
304fn encode_hex(bytes: &[u8]) -> String {
305 const DIGITS: &[u8; 16] = b"0123456789abcdef";
306 let mut encoded = String::with_capacity(bytes.len() * 2);
307 for byte in bytes {
308 encoded.push(char::from(DIGITS[usize::from(byte >> 4)]));
309 encoded.push(char::from(DIGITS[usize::from(byte & 0x0f)]));
310 }
311 encoded
312}
313
314fn rejection(
315 workflow_id: &WorkflowId,
316 activity_id: &ActivityId,
317 reason: CompletionRejectionReason,
318) -> ServerError {
319 ServerError::ActivityCompletionRejected {
320 workflow_id: workflow_id.clone(),
321 activity_id: activity_id.clone(),
322 reason,
323 }
324}
325
326#[cfg(test)]
327mod tests {
328 use super::{CompletionFences, CompletionToken, idempotency_key};
329 use crate::error::{CompletionRejectionReason, ServerError};
330 use aion_core::{ActivityId, RunId, WorkflowId};
331
332 type TestResult = Result<(), Box<dyn std::error::Error>>;
333
334 #[test]
335 fn idempotency_key_is_attempt_independent_and_site_run_scoped() {
336 let workflow = WorkflowId::new_v4();
337 let run_a = RunId::new_v4();
338 let run_b = RunId::new_v4();
339 let site_a = ActivityId::from_sequence_position(7);
340 let site_b = ActivityId::from_sequence_position(8);
341
342 let first_attempt = idempotency_key(&workflow, &run_a, &site_a);
343 let fifth_attempt = idempotency_key(&workflow, &run_a, &site_a);
344 assert_eq!(first_attempt, fifth_attempt);
345 assert_ne!(first_attempt, idempotency_key(&workflow, &run_a, &site_b));
346 assert_ne!(first_attempt, idempotency_key(&workflow, &run_b, &site_a));
347 }
348
349 #[test]
353 fn issuing_a_retry_rejects_the_stale_generation() -> TestResult {
354 let fences = CompletionFences::default();
355 let workflow = WorkflowId::new_v4();
356 let run = RunId::new_v4();
357 let activity = ActivityId::from_sequence_position(3);
358 let stale = fences.issue(&workflow, &run, &activity, 1)?;
359 let current = fences.issue(&workflow, &run, &activity, 2)?;
360
361 let rejected = fences.accept(&workflow, &activity, &stale);
362 assert!(matches!(
363 rejected,
364 Err(ServerError::ActivityCompletionRejected {
365 reason: CompletionRejectionReason::StaleGeneration,
366 ..
367 })
368 ));
369 fences.accept(&workflow, &activity, ¤t)?;
370 Ok(())
371 }
372
373 #[test]
374 fn accepted_generation_is_consumed_exactly_once() -> TestResult {
375 let fences = CompletionFences::default();
376 let workflow = WorkflowId::new_v4();
377 let run = RunId::new_v4();
378 let activity = ActivityId::from_sequence_position(4);
379 let token = fences.issue(&workflow, &run, &activity, 1)?;
380
381 fences.accept(&workflow, &activity, &token)?;
382 let duplicate = fences.accept(&workflow, &activity, &token);
383 assert!(matches!(
384 duplicate,
385 Err(ServerError::ActivityCompletionRejected {
386 reason: CompletionRejectionReason::NoCurrentGeneration,
387 ..
388 })
389 ));
390 Ok(())
391 }
392
393 #[test]
398 fn revoking_an_old_generation_does_not_remove_its_replacement() -> TestResult {
399 let fences = CompletionFences::default();
400 let workflow = WorkflowId::new_v4();
401 let run = RunId::new_v4();
402 let activity = ActivityId::from_sequence_position(6);
403 let old = fences.issue(&workflow, &run, &activity, 1)?;
404 let replacement = fences.issue(&workflow, &run, &activity, 2)?;
405
406 fences.revoke(&workflow, &activity, &old)?;
407 fences.accept(&workflow, &activity, &replacement)?;
408 Ok(())
409 }
410
411 #[test]
412 fn an_empty_wire_token_is_a_typed_compatibility_refusal() {
413 let workflow = WorkflowId::new_v4();
414 let activity = ActivityId::from_sequence_position(9);
415 let rejected = CompletionToken::from_wire(&workflow, &activity, String::new());
416 assert!(matches!(
417 rejected,
418 Err(ServerError::ActivityCompletionRejected {
419 reason: CompletionRejectionReason::MissingCompletionToken,
420 ..
421 })
422 ));
423 }
424
425 #[test]
426 fn a_pre_recovery_generation_is_rejected_after_recovery() -> TestResult {
427 let before_recovery = CompletionFences::default();
428 let workflow = WorkflowId::new_v4();
429 let run = RunId::new_v4();
430 let activity = ActivityId::from_sequence_position(5);
431 let stale = before_recovery.issue(&workflow, &run, &activity, 1)?;
432
433 let after_recovery = CompletionFences::default();
434 let current = after_recovery.issue(&workflow, &run, &activity, 1)?;
435 let rejected = after_recovery.accept(&workflow, &activity, &stale);
436 assert!(matches!(
437 rejected,
438 Err(ServerError::ActivityCompletionRejected {
439 reason: CompletionRejectionReason::StaleGeneration,
440 ..
441 })
442 ));
443 after_recovery.accept(&workflow, &activity, ¤t)?;
444 Ok(())
445 }
446
447 #[test]
461 fn a_redelivery_of_the_same_attempt_does_not_orphan_the_first_worker() -> TestResult {
462 let fences = CompletionFences::default();
463 let workflow = WorkflowId::new_v4();
464 let run = RunId::new_v4();
465 let activity = ActivityId::from_sequence_position(11);
466
467 let first = fences.issue(&workflow, &run, &activity, 1)?;
468 let second = fences.issue(&workflow, &run, &activity, 1)?;
470
471 let accepted = fences.accept(&workflow, &activity, &first);
472 assert!(
473 accepted.is_ok(),
474 "the worker that genuinely held the FIRST delivery finished the work; its result \
475 must be accepted, not thrown away: {accepted:?}"
476 );
477
478 let duplicate = fences.accept(&workflow, &activity, &second);
479 assert!(
480 matches!(
481 duplicate,
482 Err(ServerError::ActivityCompletionRejected {
483 reason: CompletionRejectionReason::NoCurrentGeneration,
484 ..
485 })
486 ),
487 "the first accepted completion must consume EVERY outstanding token for the site, \
488 so the redelivered worker's completion is the duplicate: {duplicate:?}"
489 );
490 Ok(())
491 }
492
493 #[test]
503 fn the_all_streams_closed_revoke_cannot_orphan_a_redelivered_sibling() -> TestResult {
504 let fences = CompletionFences::default();
505 let workflow = WorkflowId::new_v4();
506 let run = RunId::new_v4();
507 let activity = ActivityId::from_sequence_position(12);
508
509 let delivered = fences.issue(&workflow, &run, &activity, 1)?;
510 let redelivery = fences.issue(&workflow, &run, &activity, 1)?;
511
512 fences.revoke(&workflow, &activity, &redelivery)?;
515
516 let accepted = fences.accept(&workflow, &activity, &delivered);
517 assert!(
518 accepted.is_ok(),
519 "the sibling token the first worker still holds must survive the redelivery's own \
520 revoke: {accepted:?}"
521 );
522
523 let withdrawn = fences.accept(&workflow, &activity, &redelivery);
524 assert!(
525 matches!(
526 withdrawn,
527 Err(ServerError::ActivityCompletionRejected {
528 reason: CompletionRejectionReason::NoCurrentGeneration,
529 ..
530 })
531 ),
532 "a revoked token must never become truth: {withdrawn:?}"
533 );
534 Ok(())
535 }
536
537 #[test]
544 fn a_superseded_attempt_is_still_refused_after_a_retry_is_issued() -> TestResult {
545 let fences = CompletionFences::default();
546 let workflow = WorkflowId::new_v4();
547 let run = RunId::new_v4();
548 let activity = ActivityId::from_sequence_position(13);
549
550 let earlier = fences.issue(&workflow, &run, &activity, 1)?;
551 let current = fences.issue(&workflow, &run, &activity, 2)?;
552
553 let refused = fences.accept(&workflow, &activity, &earlier);
554 assert!(
555 matches!(
556 refused,
557 Err(ServerError::ActivityCompletionRejected {
558 reason: CompletionRejectionReason::StaleGeneration,
559 ..
560 })
561 ),
562 "a stale worker can never be recorded as truth: attempt 1's worker was superseded by \
563 attempt 2 and its completion must stay refused: {refused:?}"
564 );
565 fences.accept(&workflow, &activity, ¤t)?;
566 Ok(())
567 }
568
569 #[test]
577 fn a_superseded_execution_generation_is_still_refused_at_the_same_attempt() -> TestResult {
578 let fences = CompletionFences::default();
579 let workflow = WorkflowId::new_v4();
580 let run_a = RunId::new_v4();
581 let run_b = RunId::new_v4();
582 let activity = ActivityId::from_sequence_position(14);
583
584 let old_run = fences.issue(&workflow, &run_a, &activity, 1)?;
585 let new_run = fences.issue(&workflow, &run_b, &activity, 1)?;
586
587 let refused = fences.accept(&workflow, &activity, &old_run);
588 assert!(
589 matches!(
590 refused,
591 Err(ServerError::ActivityCompletionRejected {
592 reason: CompletionRejectionReason::StaleGeneration,
593 ..
594 })
595 ),
596 "a stale worker can never be recorded as truth: the superseded RUN's worker holds \
597 attempt 1 of a generation that no longer exists: {refused:?}"
598 );
599 fences.accept(&workflow, &activity, &new_run)?;
600 Ok(())
601 }
602
603 #[test]
607 fn a_restored_generation_carries_its_redelivered_sibling_back() -> TestResult {
608 let fences = CompletionFences::default();
609 let workflow = WorkflowId::new_v4();
610 let run = RunId::new_v4();
611 let activity = ActivityId::from_sequence_position(15);
612
613 let delivered = fences.issue(&workflow, &run, &activity, 1)?;
614 let redelivery = fences.issue(&workflow, &run, &activity, 1)?;
615
616 let accepted = fences.accept(&workflow, &activity, &delivered)?;
617 assert_eq!(accepted.attempt(), 1);
618 fences.restore_if_absent(&workflow, &activity, &accepted)?;
619
620 fences.accept(&workflow, &activity, &redelivery)?;
622 Ok(())
623 }
624
625 #[test]
628 fn a_restore_never_displaces_a_newer_generation() -> TestResult {
629 let fences = CompletionFences::default();
630 let workflow = WorkflowId::new_v4();
631 let run = RunId::new_v4();
632 let activity = ActivityId::from_sequence_position(16);
633
634 let first = fences.issue(&workflow, &run, &activity, 1)?;
635 let accepted = fences.accept(&workflow, &activity, &first)?;
636 let retry = fences.issue(&workflow, &run, &activity, 2)?;
637
638 fences.restore_if_absent(&workflow, &activity, &accepted)?;
639
640 let refused = fences.accept(&workflow, &activity, &first);
641 assert!(
642 matches!(
643 refused,
644 Err(ServerError::ActivityCompletionRejected {
645 reason: CompletionRejectionReason::StaleGeneration,
646 ..
647 })
648 ),
649 "the restore must not displace the retry that was issued meanwhile: {refused:?}"
650 );
651 fences.accept(&workflow, &activity, &retry)?;
652 Ok(())
653 }
654
655 #[test]
659 fn re_issuing_a_superseded_attempt_registers_no_generation() -> TestResult {
660 let fences = CompletionFences::default();
661 let workflow = WorkflowId::new_v4();
662 let run = RunId::new_v4();
663 let activity = ActivityId::from_sequence_position(17);
664
665 let current = fences.issue(&workflow, &run, &activity, 2)?;
666 let out_of_order = fences.issue(&workflow, &run, &activity, 1)?;
667
668 let refused = fences.accept(&workflow, &activity, &out_of_order);
669 assert!(
670 matches!(
671 refused,
672 Err(ServerError::ActivityCompletionRejected {
673 reason: CompletionRejectionReason::StaleGeneration,
674 ..
675 })
676 ),
677 "an out-of-order re-dispatch of a superseded attempt must not become acceptable: \
678 {refused:?}"
679 );
680 fences.accept(&workflow, &activity, ¤t)?;
681 Ok(())
682 }
683
684 #[test]
688 fn only_one_of_two_sibling_tokens_can_ever_become_truth() -> TestResult {
689 let fences = CompletionFences::default();
690 let workflow = WorkflowId::new_v4();
691 let run = RunId::new_v4();
692 let activity = ActivityId::from_sequence_position(18);
693
694 let first = fences.issue(&workflow, &run, &activity, 1)?;
695 let second = fences.issue(&workflow, &run, &activity, 1)?;
696
697 let accepted = [
698 fences.accept(&workflow, &activity, &second).is_ok(),
699 fences.accept(&workflow, &activity, &first).is_ok(),
700 ];
701 assert_eq!(
702 accepted.iter().filter(|ok| **ok).count(),
703 1,
704 "exactly one completion for a redelivered attempt may become truth"
705 );
706 Ok(())
707 }
708}