1use crate::payload::Payload;
16
17#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
19pub enum Category {
20 Workflow,
22 WorkflowTask,
25 Activity,
26 Timer,
27 ChildWorkflow,
28 ExternalWorkflow,
30 Update,
31 Nexus,
32 Marker,
33 SearchAttributes,
34}
35
36impl Category {
37 pub fn is_plumbing(self) -> bool {
40 matches!(self, Category::WorkflowTask)
41 }
42}
43
44#[derive(Debug, Clone, Copy, PartialEq, Eq)]
46pub enum Role {
47 Opens,
49 Continues,
51 Closes,
53}
54
55#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
57pub enum Outcome {
58 #[default]
59 Pending,
60 Completed,
61 Failed,
62 Canceled,
63 TimedOut,
64 Terminated,
65 ContinuedAsNew,
66 Rejected,
67}
68
69impl Outcome {
70 pub fn is_failure(self) -> bool {
73 matches!(
74 self,
75 Outcome::Failed | Outcome::TimedOut | Outcome::Terminated | Outcome::Rejected
76 )
77 }
78
79 pub fn label(self) -> &'static str {
80 match self {
81 Outcome::Pending => "Pending",
82 Outcome::Completed => "Completed",
83 Outcome::Failed => "Failed",
84 Outcome::Canceled => "Canceled",
85 Outcome::TimedOut => "TimedOut",
86 Outcome::Terminated => "Terminated",
87 Outcome::ContinuedAsNew => "ContinuedAsNew",
88 Outcome::Rejected => "Rejected",
89 }
90 }
91}
92
93#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
100pub enum GroupRef {
101 Workflow,
103 Opened(i64),
105}
106
107#[derive(Debug, Clone, PartialEq, Eq)]
114pub struct NormalizedEvent {
115 pub id: i64,
116 pub time: Option<i64>,
118 pub name: &'static str,
121 pub category: Category,
122 pub group: GroupRef,
123 pub role: Role,
124 pub outcome: Outcome,
125 pub subject: String,
127 pub attempt: Option<i32>,
130 pub failure: Option<String>,
132 pub fields: Vec<(&'static str, String)>,
134 pub payloads: Vec<(String, Payload)>,
137}
138
139impl NormalizedEvent {
140 pub fn new(
142 id: i64,
143 name: &'static str,
144 category: Category,
145 group: GroupRef,
146 role: Role,
147 ) -> Self {
148 Self {
149 id,
150 time: None,
151 name,
152 category,
153 group,
154 role,
155 outcome: Outcome::Pending,
156 subject: String::new(),
157 attempt: None,
158 failure: None,
159 fields: Vec::new(),
160 payloads: Vec::new(),
161 }
162 }
163
164 pub fn with_time(mut self, time: Option<i64>) -> Self {
165 self.time = time;
166 self
167 }
168
169 pub fn with_subject(mut self, subject: impl Into<String>) -> Self {
170 self.subject = subject.into();
171 self
172 }
173
174 pub fn with_outcome(mut self, outcome: Outcome) -> Self {
175 self.outcome = outcome;
176 self
177 }
178}
179
180#[derive(Debug, Clone, PartialEq, Eq)]
182pub struct Group {
183 pub key: GroupRef,
184 pub category: Category,
185 pub subject: String,
187 pub events: Vec<i64>,
189 pub started_at: Option<i64>,
190 pub ended_at: Option<i64>,
192 pub outcome: Outcome,
193 pub attempts: i32,
195 pub failure: Option<String>,
196}
197
198impl Group {
199 pub fn is_open(&self) -> bool {
201 self.ended_at.is_none() && self.outcome == Outcome::Pending
202 }
203
204 pub fn duration_ms(&self) -> Option<i64> {
206 Some(self.ended_at? - self.started_at?)
207 }
208
209 pub fn payload_ends(&self) -> Vec<i64> {
215 let mut ends: Vec<i64> = [self.events.first(), self.events.last()]
216 .into_iter()
217 .flatten()
218 .copied()
219 .collect();
220 ends.dedup();
221 ends
222 }
223
224 pub fn first_event(&self) -> Option<i64> {
226 self.events.first().copied()
227 }
228}
229
230pub fn group_events(events: &[NormalizedEvent]) -> Vec<Group> {
241 let mut groups: Vec<Group> = Vec::new();
242 let mut index: Vec<GroupRef> = Vec::new();
244
245 for ev in events {
246 let at = match index.iter().position(|k| *k == ev.group) {
247 Some(at) => at,
248 None => {
249 groups.push(Group {
250 key: ev.group,
251 category: ev.category,
252 subject: ev.subject.clone(),
253 events: Vec::new(),
254 started_at: ev.time,
255 ended_at: None,
256 outcome: Outcome::Pending,
257 attempts: 1,
258 failure: None,
259 });
260 index.push(ev.group);
261 groups.len() - 1
262 }
263 };
264 let g = &mut groups[at];
265
266 g.events.push(ev.id);
267 if let Some(n) = ev.attempt {
268 g.attempts = g.attempts.max(n);
269 }
270 if ev.role == Role::Opens && !ev.subject.is_empty() && g.subject.is_empty() {
273 g.subject = ev.subject.clone();
274 }
275 if ev.failure.is_some() {
276 g.failure = ev.failure.clone();
277 }
278 if ev.role == Role::Closes {
279 g.ended_at = ev.time;
280 g.outcome = ev.outcome;
281 }
282 }
283
284 groups
285}
286
287pub fn merge_events(existing: &mut Vec<NormalizedEvent>, incoming: Vec<NormalizedEvent>) -> usize {
296 let highest = existing.last().map(|e| e.id).unwrap_or(i64::MIN);
297 let before = existing.len();
298 existing.extend(incoming.into_iter().filter(|e| e.id > highest));
299 existing.len() - before
300}
301
302pub fn reset_point(events: &[NormalizedEvent], at: i64) -> Option<i64> {
313 events
314 .iter()
315 .filter(|e| e.id <= at)
316 .rfind(|e| e.category == Category::WorkflowTask && e.outcome == Outcome::Completed)
317 .map(|e| e.id)
318}
319
320pub fn failures(groups: &[Group]) -> Vec<&Group> {
325 groups.iter().filter(|g| g.outcome.is_failure()).collect()
326}
327
328#[cfg(test)]
329mod tests {
330 use super::*;
331
332 #[test]
333 fn an_open_group_lists_its_single_event_once() {
334 let g = Group {
337 key: GroupRef::Workflow,
338 category: Category::Workflow,
339 subject: "PayloadProbe".into(),
340 events: vec![1],
341 started_at: None,
342 ended_at: None,
343 outcome: Outcome::Pending,
344 attempts: 1,
345 failure: None,
346 };
347 assert_eq!(g.payload_ends(), vec![1]);
348 }
349
350 #[test]
351 fn a_closed_group_lists_the_event_that_opened_and_the_one_that_closed_it() {
352 let g = Group {
353 key: GroupRef::Opened(5),
354 category: Category::Activity,
355 subject: "ChargeCard".into(),
356 events: vec![5, 6, 7],
357 started_at: None,
358 ended_at: None,
359 outcome: Outcome::Completed,
360 attempts: 1,
361 failure: None,
362 };
363 assert_eq!(g.payload_ends(), vec![5, 7], "the middle carries nothing");
364 }
365
366 fn ev(id: i64, name: &'static str, group: GroupRef, role: Role, time: i64) -> NormalizedEvent {
367 NormalizedEvent::new(id, name, Category::Activity, group, role).with_time(Some(time))
368 }
369
370 fn retried_activity() -> Vec<NormalizedEvent> {
373 let mut scheduled = ev(
374 5,
375 "ActivityTaskScheduled",
376 GroupRef::Opened(5),
377 Role::Opens,
378 1_000,
379 )
380 .with_subject("ChargeCard");
381 scheduled.fields.push(("activityId", "charge".into()));
382
383 let mut started = ev(
384 6,
385 "ActivityTaskStarted",
386 GroupRef::Opened(5),
387 Role::Continues,
388 2_000,
389 );
390 started.attempt = Some(2);
392 started.failure = Some("card declined".into());
393
394 let completed = ev(
395 7,
396 "ActivityTaskCompleted",
397 GroupRef::Opened(5),
398 Role::Closes,
399 41_000,
400 )
401 .with_outcome(Outcome::Completed);
402
403 vec![scheduled, started, completed]
404 }
405
406 #[test]
407 fn three_events_become_one_group() {
408 let groups = group_events(&retried_activity());
409 assert_eq!(groups.len(), 1);
410
411 let g = &groups[0];
412 assert_eq!(g.key, GroupRef::Opened(5));
413 assert_eq!(g.subject, "ChargeCard");
414 assert_eq!(g.events, [5, 6, 7]);
415 assert_eq!(g.outcome, Outcome::Completed);
416 assert_eq!(g.attempts, 2, "the retry must be visible on the group");
417 assert_eq!(g.started_at, Some(1_000));
418 assert_eq!(g.ended_at, Some(41_000));
419 assert_eq!(g.duration_ms(), Some(40_000));
420 assert!(!g.is_open());
421 }
422
423 #[test]
424 fn a_group_with_no_closing_event_is_still_running() {
425 let events = &retried_activity()[..2];
426 let groups = group_events(events);
427 assert!(groups[0].is_open());
428 assert_eq!(groups[0].outcome, Outcome::Pending);
429 assert_eq!(
430 groups[0].duration_ms(),
431 None,
432 "a running group has no duration"
433 );
434 }
435
436 #[test]
437 fn interleaved_groups_do_not_bleed_into_each_other() {
438 let events = vec![
442 ev(
443 5,
444 "ActivityTaskScheduled",
445 GroupRef::Opened(5),
446 Role::Opens,
447 100,
448 )
449 .with_subject("A"),
450 ev(
451 6,
452 "ActivityTaskScheduled",
453 GroupRef::Opened(6),
454 Role::Opens,
455 110,
456 )
457 .with_subject("B"),
458 ev(
459 7,
460 "ActivityTaskStarted",
461 GroupRef::Opened(6),
462 Role::Continues,
463 120,
464 ),
465 ev(
466 8,
467 "ActivityTaskStarted",
468 GroupRef::Opened(5),
469 Role::Continues,
470 130,
471 ),
472 ev(
473 9,
474 "ActivityTaskFailed",
475 GroupRef::Opened(6),
476 Role::Closes,
477 140,
478 )
479 .with_outcome(Outcome::Failed),
480 ev(
481 10,
482 "ActivityTaskCompleted",
483 GroupRef::Opened(5),
484 Role::Closes,
485 150,
486 )
487 .with_outcome(Outcome::Completed),
488 ];
489 let groups = group_events(&events);
490
491 assert_eq!(groups.len(), 2);
492 assert_eq!(groups[0].subject, "A");
494 assert_eq!(groups[0].events, [5, 8, 10]);
495 assert_eq!(groups[0].outcome, Outcome::Completed);
496 assert_eq!(groups[1].subject, "B");
497 assert_eq!(groups[1].events, [6, 7, 9]);
498 assert_eq!(groups[1].outcome, Outcome::Failed);
499 }
500
501 #[test]
502 fn workflow_level_events_share_one_group() {
503 let events = vec![
504 NormalizedEvent::new(
505 1,
506 "WorkflowExecutionStarted",
507 Category::Workflow,
508 GroupRef::Workflow,
509 Role::Opens,
510 )
511 .with_time(Some(10))
512 .with_subject("OrderWorkflow"),
513 NormalizedEvent::new(
514 2,
515 "WorkflowExecutionSignaled",
516 Category::Workflow,
517 GroupRef::Workflow,
518 Role::Continues,
519 )
520 .with_time(Some(20)),
521 NormalizedEvent::new(
522 3,
523 "WorkflowExecutionCompleted",
524 Category::Workflow,
525 GroupRef::Workflow,
526 Role::Closes,
527 )
528 .with_time(Some(30))
529 .with_outcome(Outcome::Completed),
530 ];
531 let groups = group_events(&events);
532 assert_eq!(groups.len(), 1);
533 assert_eq!(groups[0].key, GroupRef::Workflow);
534 assert_eq!(groups[0].subject, "OrderWorkflow");
535 assert_eq!(groups[0].outcome, Outcome::Completed);
536 }
537
538 #[test]
539 fn an_orphaned_event_opens_its_own_group_rather_than_vanishing() {
540 let events = vec![
543 ev(
544 42,
545 "ActivityTaskCompleted",
546 GroupRef::Opened(5),
547 Role::Closes,
548 900,
549 )
550 .with_outcome(Outcome::Completed),
551 ];
552 let groups = group_events(&events);
553 assert_eq!(groups.len(), 1);
554 assert_eq!(groups[0].events, [42]);
555 assert_eq!(groups[0].outcome, Outcome::Completed);
556 }
557
558 #[test]
559 fn a_later_event_does_not_rename_its_group() {
560 let events = vec![
561 ev(
562 5,
563 "ActivityTaskScheduled",
564 GroupRef::Opened(5),
565 Role::Opens,
566 10,
567 )
568 .with_subject("real"),
569 ev(
570 6,
571 "ActivityTaskStarted",
572 GroupRef::Opened(5),
573 Role::Continues,
574 20,
575 )
576 .with_subject("other"),
577 ];
578 assert_eq!(group_events(&events)[0].subject, "real");
579 }
580
581 #[test]
582 fn the_last_failure_on_a_group_is_the_one_kept() {
583 let mut first = ev(
584 6,
585 "ActivityTaskStarted",
586 GroupRef::Opened(5),
587 Role::Continues,
588 20,
589 );
590 first.failure = Some("first".into());
591 let mut last = ev(
592 7,
593 "ActivityTaskFailed",
594 GroupRef::Opened(5),
595 Role::Closes,
596 30,
597 );
598 last.failure = Some("final".into());
599 last.outcome = Outcome::Failed;
600
601 let groups = group_events(&[
602 ev(
603 5,
604 "ActivityTaskScheduled",
605 GroupRef::Opened(5),
606 Role::Opens,
607 10,
608 ),
609 first,
610 last,
611 ]);
612 assert_eq!(groups[0].failure.as_deref(), Some("final"));
613 }
614
615 #[test]
616 fn failures_are_findable_without_scrolling() {
617 let events = vec![
618 ev(
619 1,
620 "ActivityTaskScheduled",
621 GroupRef::Opened(1),
622 Role::Opens,
623 10,
624 )
625 .with_subject("ok"),
626 ev(
627 2,
628 "ActivityTaskCompleted",
629 GroupRef::Opened(1),
630 Role::Closes,
631 20,
632 )
633 .with_outcome(Outcome::Completed),
634 ev(
635 3,
636 "ActivityTaskScheduled",
637 GroupRef::Opened(3),
638 Role::Opens,
639 30,
640 )
641 .with_subject("bad"),
642 ev(
643 4,
644 "ActivityTaskTimedOut",
645 GroupRef::Opened(3),
646 Role::Closes,
647 40,
648 )
649 .with_outcome(Outcome::TimedOut),
650 ];
651 let groups = group_events(&events);
652 let bad = failures(&groups);
653 assert_eq!(bad.len(), 1);
654 assert_eq!(bad[0].subject, "bad");
655 }
656
657 #[test]
658 fn every_outcome_agrees_with_itself_about_being_a_failure() {
659 for (o, fail) in [
660 (Outcome::Pending, false),
661 (Outcome::Completed, false),
662 (Outcome::Canceled, false),
663 (Outcome::ContinuedAsNew, false),
664 (Outcome::Failed, true),
665 (Outcome::TimedOut, true),
666 (Outcome::Terminated, true),
667 (Outcome::Rejected, true),
668 ] {
669 assert_eq!(o.is_failure(), fail, "{} classified wrongly", o.label());
670 }
671 }
672
673 #[test]
674 fn replayed_events_are_not_appended_twice() {
675 let mut held: Vec<NormalizedEvent> = retried_activity();
679 assert_eq!(held.len(), 3);
680
681 let replay = retried_activity();
682 assert_eq!(merge_events(&mut held, replay), 0, "nothing was new");
683 assert_eq!(held.len(), 3);
684
685 let fresh = vec![ev(
687 8,
688 "TimerStarted",
689 GroupRef::Opened(8),
690 Role::Opens,
691 50_000,
692 )];
693 assert_eq!(merge_events(&mut held, fresh), 1);
694 assert_eq!(held.len(), 4);
695 }
696
697 #[test]
698 fn a_partial_replay_keeps_only_the_tail() {
699 let mut held: Vec<NormalizedEvent> = retried_activity();
700 let mut incoming = retried_activity()[1..].to_vec();
702 incoming.push(ev(
703 8,
704 "TimerStarted",
705 GroupRef::Opened(8),
706 Role::Opens,
707 50_000,
708 ));
709 incoming.push(ev(
710 9,
711 "TimerFired",
712 GroupRef::Opened(8),
713 Role::Closes,
714 60_000,
715 ));
716
717 assert_eq!(merge_events(&mut held, incoming), 2);
718 let ids: Vec<i64> = held.iter().map(|e| e.id).collect();
719 assert_eq!(ids, [5, 6, 7, 8, 9]);
720 }
721
722 #[test]
723 fn merging_into_an_empty_history_keeps_everything() {
724 let mut held = Vec::new();
725 assert_eq!(merge_events(&mut held, retried_activity()), 3);
726 assert_eq!(held.len(), 3);
727 }
728
729 #[test]
730 fn a_reset_resolves_back_to_the_last_completed_workflow_task() {
731 let events = vec![
734 NormalizedEvent::new(1, "S", Category::Workflow, GroupRef::Workflow, Role::Opens),
735 NormalizedEvent::new(
736 2,
737 "WTS",
738 Category::WorkflowTask,
739 GroupRef::Opened(2),
740 Role::Opens,
741 ),
742 NormalizedEvent::new(
743 3,
744 "WTC",
745 Category::WorkflowTask,
746 GroupRef::Opened(2),
747 Role::Closes,
748 )
749 .with_outcome(Outcome::Completed),
750 NormalizedEvent::new(
751 4,
752 "ATS",
753 Category::Activity,
754 GroupRef::Opened(4),
755 Role::Opens,
756 ),
757 NormalizedEvent::new(
758 5,
759 "ATC",
760 Category::Activity,
761 GroupRef::Opened(4),
762 Role::Closes,
763 )
764 .with_outcome(Outcome::Completed),
765 ];
766 assert_eq!(
767 reset_point(&events, 5),
768 Some(3),
769 "back to the workflow task"
770 );
771 assert_eq!(reset_point(&events, 3), Some(3), "already on one");
772 assert_eq!(reset_point(&events, 2), None, "nothing completed yet");
773 }
774
775 #[test]
776 fn a_failed_workflow_task_is_not_a_reset_point() {
777 let events = vec![
780 NormalizedEvent::new(
781 2,
782 "WTS",
783 Category::WorkflowTask,
784 GroupRef::Opened(2),
785 Role::Opens,
786 ),
787 NormalizedEvent::new(
788 3,
789 "WTF",
790 Category::WorkflowTask,
791 GroupRef::Opened(2),
792 Role::Closes,
793 )
794 .with_outcome(Outcome::Failed),
795 ];
796 assert_eq!(reset_point(&events, 3), None);
797 }
798
799 #[test]
800 fn a_reset_takes_the_latest_valid_point_not_the_first() {
801 let events = vec![
802 NormalizedEvent::new(
803 3,
804 "WTC",
805 Category::WorkflowTask,
806 GroupRef::Opened(2),
807 Role::Closes,
808 )
809 .with_outcome(Outcome::Completed),
810 NormalizedEvent::new(
811 9,
812 "WTC",
813 Category::WorkflowTask,
814 GroupRef::Opened(8),
815 Role::Closes,
816 )
817 .with_outcome(Outcome::Completed),
818 NormalizedEvent::new(
819 12,
820 "ATC",
821 Category::Activity,
822 GroupRef::Opened(10),
823 Role::Closes,
824 )
825 .with_outcome(Outcome::Completed),
826 ];
827 assert_eq!(reset_point(&events, 12), Some(9));
828 assert_eq!(reset_point(&events, 8), Some(3));
829 }
830
831 #[test]
832 fn an_empty_history_groups_to_nothing() {
833 assert!(group_events(&[]).is_empty());
834 assert!(failures(&[]).is_empty());
835 }
836}