1use std::collections::{BTreeMap, HashMap};
4use std::sync::{Mutex, MutexGuard, PoisonError};
5
6use aion_core::{
7 Event, TimerId, WorkflowFilter, WorkflowId, WorkflowStatus, WorkflowSummary, status_from_events,
8};
9use async_trait::async_trait;
10use chrono::{DateTime, Utc};
11
12use crate::namespace::{
13 MintOutcome, NamespaceOrigin, NamespacePlacement, NamespaceRecord, NamespaceState,
14 NamespaceStore,
15};
16use crate::package::{PackageRecord, PackageRouteRecord, PackageStore};
17use crate::visibility::{ListWorkflowsFilter, VisibilityRecord, VisibilityStore};
18use crate::{
19 ReadableEventStore, RunSummary, StoreError, TimerEntry, WritableEventStore, WriteToken,
20};
21
22#[derive(Debug, Default)]
24pub struct InMemoryStore {
25 state: Mutex<InMemoryState>,
26 namespaces: Mutex<BTreeMap<String, NamespaceRecord>>,
27}
28
29#[async_trait]
30impl VisibilityStore for InMemoryStore {
31 async fn record_visibility(&self, record: VisibilityRecord) -> Result<(), StoreError> {
32 let mut state = self.lock_state()?;
33 state
34 .visibility
35 .insert((record.workflow_id.clone(), record.run_id.clone()), record);
36 Ok(())
37 }
38
39 async fn list_workflows(
40 &self,
41 filter: ListWorkflowsFilter,
42 ) -> Result<Vec<crate::visibility::WorkflowSummary>, StoreError> {
43 let state = self.lock_state()?;
44 let mut summaries = state
45 .visibility
46 .values()
47 .cloned()
48 .map(crate::visibility::WorkflowSummary::from)
49 .filter(|summary| filter.matches(summary))
50 .collect::<Vec<_>>();
51 summaries.sort_by(|left, right| {
52 left.start_time.cmp(&right.start_time).then_with(|| {
53 left.workflow_id
54 .to_string()
55 .cmp(&right.workflow_id.to_string())
56 })
57 });
58 let offset = filter.offset.and_then(|value| usize::try_from(value).ok());
59 if let Some(offset) = offset {
60 summaries = summaries.into_iter().skip(offset).collect();
61 }
62 if let Some(limit) = filter.limit.and_then(|value| usize::try_from(value).ok()) {
63 summaries.truncate(limit);
64 }
65 Ok(summaries)
66 }
67
68 async fn count_workflows(&self, filter: ListWorkflowsFilter) -> Result<u64, StoreError> {
69 let state = self.lock_state()?;
70 Ok(state
71 .visibility
72 .values()
73 .cloned()
74 .map(crate::visibility::WorkflowSummary::from)
75 .filter(|summary| filter.matches(summary))
76 .count()
77 .try_into()
78 .unwrap_or(u64::MAX))
79 }
80}
81
82#[derive(Debug, Default)]
83struct InMemoryState {
84 histories: HashMap<WorkflowId, Vec<Event>>,
85 timers: HashMap<(WorkflowId, TimerId), TimerEntry>,
86 visibility: HashMap<(WorkflowId, aion_core::RunId), VisibilityRecord>,
87 packages: HashMap<(String, String), PackageRecord>,
88 package_routes: HashMap<String, String>,
89}
90
91impl InMemoryStore {
92 fn lock_state(&self) -> Result<MutexGuard<'_, InMemoryState>, StoreError> {
93 self.state
94 .lock()
95 .map_err(|error| StoreError::Backend(format!("in-memory store lock poisoned: {error}")))
96 }
97
98 fn lock_namespaces(&self) -> MutexGuard<'_, BTreeMap<String, NamespaceRecord>> {
104 self.namespaces
105 .lock()
106 .unwrap_or_else(PoisonError::into_inner)
107 }
108}
109
110fn history_head(history: &[Event]) -> u64 {
111 history.iter().map(Event::seq).max().unwrap_or_default()
112}
113
114fn history_in_sequence_order(history: &[Event]) -> Vec<Event> {
115 let mut ordered = history.to_vec();
116 ordered.sort_by_key(Event::seq);
117 ordered
118}
119
120#[async_trait]
121impl PackageStore for InMemoryStore {
122 async fn put_package(&self, record: PackageRecord) -> Result<(), StoreError> {
123 let primary = record.workflow_type.clone();
124 self.put_package_with_routes(record, &[primary]).await
125 }
126
127 async fn put_package_with_routes(
128 &self,
129 record: PackageRecord,
130 route_workflow_types: &[String],
131 ) -> Result<(), StoreError> {
132 let mut state = self.lock_state()?;
133 for workflow_type in route_workflow_types {
134 state
135 .package_routes
136 .insert(workflow_type.clone(), record.content_hash.clone());
137 }
138 state.packages.insert(
139 (record.workflow_type.clone(), record.content_hash.clone()),
140 record,
141 );
142 Ok(())
143 }
144
145 async fn list_packages(&self) -> Result<Vec<PackageRecord>, StoreError> {
146 let state = self.lock_state()?;
147 let mut records: Vec<PackageRecord> = state.packages.values().cloned().collect();
148 records.sort_by(|left, right| {
149 left.deployed_at
150 .cmp(&right.deployed_at)
151 .then_with(|| left.workflow_type.cmp(&right.workflow_type))
152 .then_with(|| left.content_hash.cmp(&right.content_hash))
153 });
154 Ok(records)
155 }
156
157 async fn delete_package(
158 &self,
159 workflow_type: &str,
160 content_hash: &str,
161 ) -> Result<(), StoreError> {
162 let mut state = self.lock_state()?;
163 state
164 .packages
165 .remove(&(workflow_type.to_owned(), content_hash.to_owned()));
166 Ok(())
167 }
168
169 async fn put_package_route(
170 &self,
171 workflow_type: &str,
172 content_hash: &str,
173 ) -> Result<(), StoreError> {
174 let mut state = self.lock_state()?;
175 state
176 .package_routes
177 .insert(workflow_type.to_owned(), content_hash.to_owned());
178 Ok(())
179 }
180
181 async fn list_package_routes(&self) -> Result<Vec<PackageRouteRecord>, StoreError> {
182 let state = self.lock_state()?;
183 let mut routes: Vec<PackageRouteRecord> = state
184 .package_routes
185 .iter()
186 .map(|(workflow_type, content_hash)| PackageRouteRecord {
187 workflow_type: workflow_type.clone(),
188 content_hash: content_hash.clone(),
189 })
190 .collect();
191 routes.sort_by(|left, right| left.workflow_type.cmp(&right.workflow_type));
192 Ok(routes)
193 }
194}
195
196#[async_trait]
197impl NamespaceStore for InMemoryStore {
198 async fn register_namespace(
199 &self,
200 name: &str,
201 origin: NamespaceOrigin,
202 ) -> Result<MintOutcome, StoreError> {
203 let now = Utc::now();
204 let mut namespaces = self.lock_namespaces();
205 if let Some(existing) = namespaces.get_mut(name) {
206 existing.bump_last_seen(now);
207 Ok(MintOutcome::AlreadyExisted)
208 } else {
209 namespaces.insert(
210 name.to_owned(),
211 NamespaceRecord::new_minted(name, origin, now),
212 );
213 Ok(MintOutcome::Created)
214 }
215 }
216
217 async fn put_namespace(&self, record: NamespaceRecord) -> Result<MintOutcome, StoreError> {
218 let now = Utc::now();
219 let mut namespaces = self.lock_namespaces();
220 if let Some(existing) = namespaces.get_mut(&record.name) {
221 existing.bump_last_seen(now);
225 Ok(MintOutcome::AlreadyExisted)
226 } else {
227 namespaces.insert(record.name.clone(), record);
228 Ok(MintOutcome::Created)
229 }
230 }
231
232 async fn list_namespaces(&self) -> Result<Vec<NamespaceRecord>, StoreError> {
233 let namespaces = self.lock_namespaces();
234 let mut records: Vec<NamespaceRecord> = namespaces.values().cloned().collect();
235 records.sort_by(|left, right| {
236 left.created_at
237 .cmp(&right.created_at)
238 .then_with(|| left.name.cmp(&right.name))
239 });
240 Ok(records)
241 }
242
243 async fn get_namespace(&self, name: &str) -> Result<Option<NamespaceRecord>, StoreError> {
244 let namespaces = self.lock_namespaces();
245 Ok(namespaces.get(name).cloned())
246 }
247
248 async fn set_namespace_placement(
249 &self,
250 name: &str,
251 placement: NamespacePlacement,
252 ) -> Result<Option<()>, StoreError> {
253 let now = Utc::now();
254 let mut namespaces = self.lock_namespaces();
255 let Some(existing) = namespaces.get_mut(name) else {
256 return Ok(None);
259 };
260 existing.placement = placement;
261 existing.bump_last_seen(now);
262 Ok(Some(()))
263 }
264
265 async fn deprecate_namespace(&self, name: &str) -> Result<(), StoreError> {
266 let mut namespaces = self.lock_namespaces();
267 if let Some(existing) = namespaces.get_mut(name) {
268 existing.state = NamespaceState::Deprecated;
269 }
270 Ok(())
273 }
274}
275
276#[async_trait]
277impl WritableEventStore for InMemoryStore {
278 async fn append(
279 &self,
280 _token: WriteToken,
281 workflow_id: &WorkflowId,
282 events: &[Event],
283 expected_seq: u64,
284 ) -> Result<(), StoreError> {
285 let mut state = self.lock_state()?;
286 let current_head = state
287 .histories
288 .get(workflow_id)
289 .map_or(0, |history| history_head(history));
290
291 if current_head != expected_seq {
292 return Err(StoreError::SequenceConflict {
293 expected: expected_seq,
294 found: current_head,
295 });
296 }
297
298 if events.is_empty() {
299 return Ok(());
300 }
301
302 for (next_seq, event) in (expected_seq + 1..).zip(events.iter()) {
303 if event.seq() != next_seq {
304 return Err(StoreError::Backend(format!(
305 "event sequence must be contiguous: expected {next_seq}, got {}",
306 event.seq()
307 )));
308 }
309 }
310
311 state
312 .histories
313 .entry(workflow_id.clone())
314 .or_default()
315 .extend(events.iter().cloned());
316 Ok(())
317 }
318}
319
320#[async_trait]
321impl ReadableEventStore for InMemoryStore {
322 async fn read_history(&self, workflow_id: &WorkflowId) -> Result<Vec<Event>, StoreError> {
323 let state = self.lock_state()?;
324 Ok(state
325 .histories
326 .get(workflow_id)
327 .map_or_else(Vec::new, |history| history_in_sequence_order(history)))
328 }
329
330 async fn read_history_from(
331 &self,
332 workflow_id: &WorkflowId,
333 from_seq: u64,
334 ) -> Result<Vec<Event>, StoreError> {
335 let state = self.lock_state()?;
336 Ok(state
337 .histories
338 .get(workflow_id)
339 .map_or_else(Vec::new, |history| {
340 let mut events = history
341 .iter()
342 .filter(|event| event.seq() >= from_seq)
343 .cloned()
344 .collect::<Vec<_>>();
345 events.sort_by_key(Event::seq);
346 events
347 }))
348 }
349
350 async fn read_run_chain(
351 &self,
352 workflow_id: &WorkflowId,
353 ) -> Result<Vec<RunSummary>, StoreError> {
354 let state = self.lock_state()?;
355 let Some(history) = state.histories.get(workflow_id) else {
356 return Ok(Vec::new());
357 };
358
359 crate::run_chain::run_chain_from_history(history)
360 }
361
362 async fn list_workflow_ids(&self) -> Result<Vec<WorkflowId>, StoreError> {
363 let state = self.lock_state()?;
364 let mut workflow_ids = state.histories.keys().cloned().collect::<Vec<_>>();
365 workflow_ids.sort_by_key(ToString::to_string);
366 Ok(workflow_ids)
367 }
368
369 async fn list_active(&self) -> Result<Vec<WorkflowId>, StoreError> {
370 let state = self.lock_state()?;
371 let mut active = state
372 .histories
373 .iter()
374 .filter(|(_, history)| {
375 matches!(
376 status_from_events(&history_in_sequence_order(history)),
377 WorkflowStatus::Running
378 )
379 })
380 .map(|(workflow_id, _)| workflow_id.clone())
381 .collect::<Vec<_>>();
382 active.sort_by_key(ToString::to_string);
383 Ok(active)
384 }
385
386 async fn list_paused(&self) -> Result<Vec<WorkflowId>, StoreError> {
387 let state = self.lock_state()?;
388 let mut paused = state
389 .histories
390 .iter()
391 .filter(|(_, history)| {
392 matches!(
393 status_from_events(&history_in_sequence_order(history)),
394 WorkflowStatus::Paused
395 )
396 })
397 .map(|(workflow_id, _)| workflow_id.clone())
398 .collect::<Vec<_>>();
399 paused.sort_by_key(ToString::to_string);
400 Ok(paused)
401 }
402
403 async fn query(&self, filter: &WorkflowFilter) -> Result<Vec<WorkflowSummary>, StoreError> {
404 let state = self.lock_state()?;
405 let mut summaries = state
406 .histories
407 .values()
408 .filter_map(|history| {
409 WorkflowSummary::from_history(&history_in_sequence_order(history))
410 })
411 .filter(|summary| filter.matches(summary))
412 .collect::<Vec<_>>();
413 summaries.sort_by(|left, right| {
414 left.started_at.cmp(&right.started_at).then_with(|| {
415 left.workflow_id
416 .to_string()
417 .cmp(&right.workflow_id.to_string())
418 })
419 });
420 Ok(summaries)
421 }
422
423 async fn schedule_timer(
424 &self,
425 workflow_id: &WorkflowId,
426 timer_id: &TimerId,
427 fire_at: DateTime<Utc>,
428 ) -> Result<(), StoreError> {
429 let mut state = self.lock_state()?;
430 state.timers.insert(
431 (workflow_id.clone(), timer_id.clone()),
432 TimerEntry {
433 workflow_id: workflow_id.clone(),
434 timer_id: timer_id.clone(),
435 fire_at,
436 },
437 );
438 Ok(())
439 }
440
441 async fn expired_timers(&self, as_of: DateTime<Utc>) -> Result<Vec<TimerEntry>, StoreError> {
442 let state = self.lock_state()?;
443 let mut timers = state
444 .timers
445 .values()
446 .filter(|entry| entry.fire_at <= as_of)
447 .cloned()
448 .collect::<Vec<_>>();
449 timers.sort_by(|left, right| {
450 left.fire_at
451 .cmp(&right.fire_at)
452 .then_with(|| {
453 left.workflow_id
454 .to_string()
455 .cmp(&right.workflow_id.to_string())
456 })
457 .then_with(|| left.timer_id.to_string().cmp(&right.timer_id.to_string()))
458 });
459 Ok(timers)
460 }
461}
462
463#[cfg(test)]
464mod tests {
465 use std::sync::Arc;
466
467 use aion_core::{
468 Event, EventEnvelope, Payload, TimerId, WorkflowError, WorkflowFilter, WorkflowId,
469 WorkflowStatus,
470 };
471 use chrono::{DateTime, Utc};
472 use serde_json::json;
473 use tokio::task;
474 use uuid::Uuid;
475
476 use super::InMemoryStore;
477 use crate::{ReadableEventStore, StoreError, TimerEntry, WritableEventStore, WriteToken};
478
479 fn write_token() -> WriteToken {
480 WriteToken::recorder()
481 }
482
483 fn recorded_at(offset_seconds: i64) -> DateTime<Utc> {
484 DateTime::from_timestamp(1_700_000_000 + offset_seconds, 0).unwrap_or_default()
485 }
486
487 fn workflow_id(value: u128) -> WorkflowId {
488 WorkflowId::new(Uuid::from_u128(value))
489 }
490
491 fn envelope(seq: u64, workflow_id: &WorkflowId) -> EventEnvelope {
492 EventEnvelope {
493 seq,
494 recorded_at: recorded_at(i64::try_from(seq).unwrap_or_default()),
495 workflow_id: workflow_id.clone(),
496 }
497 }
498
499 fn run_id(value: u128) -> aion_core::RunId {
500 aion_core::RunId::new(Uuid::from_u128(value))
501 }
502
503 fn payload(label: &str) -> Payload {
504 Payload::from_json(&json!({ "label": label })).unwrap_or_else(|error| {
505 Payload::new(
506 aion_core::ContentType::Json,
507 format!("{{\"payload_error\":\"{error}\"}}").into_bytes(),
508 )
509 })
510 }
511
512 fn workflow_started(seq: u64, workflow_id: &WorkflowId, workflow_type: &str) -> Event {
513 Event::WorkflowStarted {
514 envelope: envelope(seq, workflow_id),
515 workflow_type: workflow_type.to_owned(),
516 input: payload("input"),
517 run_id: aion_core::RunId::new(uuid::Uuid::from_u128(1)),
518 parent_run_id: None,
519 package_version: aion_core::PackageVersion::new("a".repeat(64)),
520 }
521 }
522
523 fn workflow_completed(seq: u64, workflow_id: &WorkflowId) -> Event {
524 Event::WorkflowCompleted {
525 envelope: envelope(seq, workflow_id),
526 result: payload("result"),
527 }
528 }
529
530 fn workflow_failed(seq: u64, workflow_id: &WorkflowId) -> Event {
531 Event::WorkflowFailed {
532 envelope: envelope(seq, workflow_id),
533 error: WorkflowError {
534 message: String::from("failed"),
535 details: None,
536 },
537 }
538 }
539
540 #[tokio::test]
541 async fn read_history_returns_empty_for_unknown_workflow() -> Result<(), StoreError> {
542 let store = InMemoryStore::default();
543
544 assert_eq!(store.read_history(&workflow_id(1)).await?, Vec::new());
545 Ok(())
546 }
547
548 #[tokio::test]
549 async fn append_preserves_sequence_order() -> Result<(), StoreError> {
550 let store = InMemoryStore::default();
551 let workflow_id = workflow_id(1);
552 let first = workflow_started(1, &workflow_id, "checkout");
553 let second = workflow_completed(2, &workflow_id);
554
555 store
556 .append(write_token(), &workflow_id, std::slice::from_ref(&first), 0)
557 .await?;
558 store
559 .append(
560 write_token(),
561 &workflow_id,
562 std::slice::from_ref(&second),
563 1,
564 )
565 .await?;
566
567 assert_eq!(store.read_history(&workflow_id).await?, vec![first, second]);
568 Ok(())
569 }
570
571 #[tokio::test]
572 async fn list_active_returns_only_running_workflows() -> Result<(), StoreError> {
573 let store = InMemoryStore::default();
574 let running = workflow_id(1);
575 let completed = workflow_id(2);
576
577 store
578 .append(
579 write_token(),
580 &running,
581 &[workflow_started(1, &running, "checkout")],
582 0,
583 )
584 .await?;
585 store
586 .append(
587 write_token(),
588 &completed,
589 &[
590 workflow_started(1, &completed, "checkout"),
591 workflow_completed(2, &completed),
592 ],
593 0,
594 )
595 .await?;
596
597 assert_eq!(store.list_active().await?, vec![running]);
598 Ok(())
599 }
600
601 fn workflow_paused(seq: u64, workflow_id: &WorkflowId) -> Event {
602 Event::WorkflowPaused {
603 envelope: envelope(seq, workflow_id),
604 run_id: run_id(1),
605 reason: None,
606 operator: None,
607 }
608 }
609
610 #[tokio::test]
615 async fn list_paused_and_list_active_partition_by_projected_status() -> Result<(), StoreError> {
616 let store = InMemoryStore::default();
617 let running = workflow_id(1);
618 let paused = workflow_id(2);
619
620 store
621 .append(
622 write_token(),
623 &running,
624 &[workflow_started(1, &running, "checkout")],
625 0,
626 )
627 .await?;
628 store
629 .append(
630 write_token(),
631 &paused,
632 &[
633 workflow_started(1, &paused, "checkout"),
634 workflow_paused(2, &paused),
635 ],
636 0,
637 )
638 .await?;
639
640 assert_eq!(
641 store.list_active().await?,
642 vec![running],
643 "a paused run is excluded from list_active (not respawned)"
644 );
645 assert_eq!(
646 store.list_paused().await?,
647 vec![paused],
648 "list_paused returns exactly the paused run (the hold rebuild source)"
649 );
650 Ok(())
651 }
652
653 #[tokio::test]
654 async fn list_workflow_ids_returns_running_and_terminal_histories() -> Result<(), StoreError> {
655 let store = InMemoryStore::default();
656 let running = workflow_id(2);
657 let completed = workflow_id(1);
658
659 store
660 .append(
661 write_token(),
662 &running,
663 &[workflow_started(1, &running, "checkout")],
664 0,
665 )
666 .await?;
667 store
668 .append(
669 write_token(),
670 &completed,
671 &[
672 workflow_started(1, &completed, "checkout"),
673 workflow_completed(2, &completed),
674 ],
675 0,
676 )
677 .await?;
678
679 assert_eq!(store.list_workflow_ids().await?, vec![completed, running]);
680 Ok(())
681 }
682
683 #[tokio::test]
684 async fn read_run_chain_projects_run_id_from_started_event() -> Result<(), StoreError> {
685 let store = InMemoryStore::default();
686 let workflow_id = workflow_id(1);
687
688 store
689 .append(
690 write_token(),
691 &workflow_id,
692 &[
693 workflow_started(1, &workflow_id, "checkout"),
694 workflow_completed(2, &workflow_id),
695 ],
696 0,
697 )
698 .await?;
699
700 let chain = store.read_run_chain(&workflow_id).await?;
701
702 assert_eq!(chain.len(), 1);
703 assert_eq!(chain[0].run_id, run_id(1));
705 assert_eq!(chain[0].status, WorkflowStatus::Completed);
706 assert_eq!(chain[0].closed_at, Some(recorded_at(2)));
707 Ok(())
708 }
709
710 #[tokio::test]
711 async fn query_uses_core_filter_semantics() -> Result<(), StoreError> {
712 let store = InMemoryStore::default();
713 let running_checkout = workflow_id(1);
714 let completed_checkout = workflow_id(2);
715 let failed_billing = workflow_id(3);
716
717 store
718 .append(
719 write_token(),
720 &running_checkout,
721 &[workflow_started(1, &running_checkout, "checkout")],
722 0,
723 )
724 .await?;
725 store
726 .append(
727 write_token(),
728 &completed_checkout,
729 &[
730 workflow_started(1, &completed_checkout, "checkout"),
731 workflow_completed(2, &completed_checkout),
732 ],
733 0,
734 )
735 .await?;
736 store
737 .append(
738 write_token(),
739 &failed_billing,
740 &[
741 workflow_started(1, &failed_billing, "billing"),
742 workflow_failed(2, &failed_billing),
743 ],
744 0,
745 )
746 .await?;
747
748 let filter = WorkflowFilter {
749 workflow_type: Some(String::from("checkout")),
750 status: Some(WorkflowStatus::Completed),
751 started_after: Some(recorded_at(1)),
752 started_before: Some(recorded_at(1)),
753 parent: None,
754 };
755 let summaries = store.query(&filter).await?;
756
757 assert_eq!(summaries.len(), 1);
758 assert_eq!(summaries[0].workflow_id, completed_checkout);
759 assert_eq!(summaries[0].status, WorkflowStatus::Completed);
760 Ok(())
761 }
762
763 #[tokio::test]
764 async fn stale_expected_sequence_writes_nothing() -> Result<(), StoreError> {
765 let store = InMemoryStore::default();
766 let workflow_id = workflow_id(1);
767 let first = workflow_started(1, &workflow_id, "checkout");
768
769 store
770 .append(write_token(), &workflow_id, std::slice::from_ref(&first), 0)
771 .await?;
772 let conflict = store
773 .append(
774 write_token(),
775 &workflow_id,
776 &[workflow_completed(2, &workflow_id)],
777 0,
778 )
779 .await;
780
781 assert_eq!(
782 conflict,
783 Err(StoreError::SequenceConflict {
784 expected: 0,
785 found: 1,
786 })
787 );
788 assert_eq!(store.read_history(&workflow_id).await?, vec![first]);
789 Ok(())
790 }
791
792 #[tokio::test]
793 async fn append_rejects_non_contiguous_event_sequences() -> Result<(), StoreError> {
794 let store = InMemoryStore::default();
795 let wf = workflow_id(1);
796
797 let result = store
798 .append(
799 write_token(),
800 &wf,
801 &[
802 workflow_started(1, &wf, "checkout"),
803 workflow_completed(5, &wf),
804 ],
805 0,
806 )
807 .await;
808
809 assert!(result.is_err());
810 assert!(matches!(result, Err(StoreError::Backend(_))));
811 assert_eq!(store.read_history(&wf).await?, Vec::new());
812 Ok(())
813 }
814
815 #[tokio::test]
816 async fn concurrent_appends_on_same_expected_sequence_conflict_once() -> Result<(), StoreError>
817 {
818 let store = Arc::new(InMemoryStore::default());
819 let workflow_id = workflow_id(1);
820 let first_store = Arc::clone(&store);
821 let first_workflow = workflow_id.clone();
822 let second_store = Arc::clone(&store);
823 let second_workflow = workflow_id.clone();
824
825 let first = task::spawn(async move {
826 first_store
827 .append(
828 write_token(),
829 &first_workflow,
830 &[workflow_started(1, &first_workflow, "checkout")],
831 0,
832 )
833 .await
834 });
835 let second = task::spawn(async move {
836 second_store
837 .append(
838 write_token(),
839 &second_workflow,
840 &[workflow_completed(1, &second_workflow)],
841 0,
842 )
843 .await
844 });
845
846 let results = [
847 first
848 .await
849 .map_err(|error| StoreError::Backend(format!("append task failed: {error}")))?,
850 second
851 .await
852 .map_err(|error| StoreError::Backend(format!("append task failed: {error}")))?,
853 ];
854
855 assert_eq!(results.iter().filter(|result| result.is_ok()).count(), 1);
856 assert_eq!(
857 results
858 .iter()
859 .filter(|result| matches!(
860 result,
861 Err(StoreError::SequenceConflict {
862 expected: 0,
863 found: 1
864 })
865 ))
866 .count(),
867 1
868 );
869 assert_eq!(store.read_history(&workflow_id).await?.len(), 1);
870 Ok(())
871 }
872
873 #[tokio::test]
874 async fn rescheduling_same_timer_replaces_prior_fire_at() -> Result<(), StoreError> {
875 let store = InMemoryStore::default();
876 let workflow_id = workflow_id(1);
877 let timer_id = TimerId::anonymous(1);
878 let first_fire_at = recorded_at(10);
879 let replacement_fire_at = recorded_at(30);
880
881 store
882 .schedule_timer(&workflow_id, &timer_id, first_fire_at)
883 .await?;
884 store
885 .schedule_timer(&workflow_id, &timer_id, replacement_fire_at)
886 .await?;
887
888 assert_eq!(store.expired_timers(first_fire_at).await?, Vec::new());
889 assert_eq!(
890 store.expired_timers(replacement_fire_at).await?,
891 vec![TimerEntry {
892 workflow_id,
893 timer_id,
894 fire_at: replacement_fire_at,
895 }]
896 );
897 Ok(())
898 }
899
900 #[tokio::test]
901 async fn expired_timers_include_boundary_and_exclude_future() -> Result<(), StoreError> {
902 let store = InMemoryStore::default();
903 let workflow_id = workflow_id(1);
904 let past_timer = TimerId::anonymous(1);
905 let boundary_timer = TimerId::anonymous(2);
906 let future_timer = TimerId::anonymous(3);
907 let as_of = recorded_at(20);
908
909 store
910 .schedule_timer(&workflow_id, &future_timer, recorded_at(30))
911 .await?;
912 store
913 .schedule_timer(&workflow_id, &boundary_timer, as_of)
914 .await?;
915 store
916 .schedule_timer(&workflow_id, &past_timer, recorded_at(10))
917 .await?;
918
919 assert_eq!(
920 store.expired_timers(as_of).await?,
921 vec![
922 TimerEntry {
923 workflow_id: workflow_id.clone(),
924 timer_id: past_timer,
925 fire_at: recorded_at(10),
926 },
927 TimerEntry {
928 workflow_id,
929 timer_id: boundary_timer,
930 fire_at: as_of,
931 },
932 ]
933 );
934 Ok(())
935 }
936}
937
938#[cfg(test)]
939mod namespace_tests {
940 #![allow(clippy::expect_used)]
941
942 use super::InMemoryStore;
943 use crate::namespace::{
944 MintOutcome, NamespaceOrigin, NamespacePlacement, NamespaceRecord, NamespaceState,
945 };
946 use crate::{NamespaceStore, StoreError};
947 use chrono::{TimeZone, Utc};
948 use std::collections::BTreeSet;
949
950 fn labels(values: &[&str]) -> BTreeSet<String> {
951 values.iter().map(|v| (*v).to_owned()).collect()
952 }
953
954 #[tokio::test]
958 async fn set_placement_updates_only_placement_and_reports_not_found() -> Result<(), StoreError>
959 {
960 let store = InMemoryStore::default();
961 store
962 .register_namespace("orders", NamespaceOrigin::Explicit)
963 .await?;
964 let original = store
965 .get_namespace("orders")
966 .await?
967 .expect("namespace must persist");
968 assert_eq!(original.placement, NamespacePlacement::Unplaced);
969
970 let placement = NamespacePlacement::Prefer {
971 nodes: labels(&["n1", "n2"]),
972 };
973 assert_eq!(
974 store
975 .set_namespace_placement("orders", placement.clone())
976 .await?,
977 Some(())
978 );
979 let updated = store
980 .get_namespace("orders")
981 .await?
982 .expect("namespace must persist");
983 assert_eq!(updated.placement, placement);
984 assert_eq!(updated.origin, original.origin);
986 assert_eq!(updated.created_at, original.created_at);
987 assert_eq!(updated.state, original.state);
988
989 assert_eq!(
991 store.set_namespace_placement("orders", placement).await?,
992 Some(())
993 );
994
995 assert_eq!(
997 store
998 .set_namespace_placement("ghost", NamespacePlacement::Unplaced)
999 .await?,
1000 None
1001 );
1002 assert!(store.get_namespace("ghost").await?.is_none());
1003 Ok(())
1004 }
1005
1006 #[tokio::test]
1007 async fn register_creates_if_absent_and_persists() -> Result<(), StoreError> {
1008 let store = InMemoryStore::default();
1009
1010 let outcome = store
1011 .register_namespace("orders", NamespaceOrigin::WorkerMint)
1012 .await?;
1013
1014 assert_eq!(outcome, MintOutcome::Created);
1015 let record = store
1016 .get_namespace("orders")
1017 .await?
1018 .expect("namespace must persist");
1019 assert_eq!(record.name, "orders");
1020 assert_eq!(record.origin, NamespaceOrigin::WorkerMint);
1021 assert_eq!(record.state, NamespaceState::Active);
1022 assert_eq!(record.created_at, record.last_seen);
1023 Ok(())
1024 }
1025
1026 #[tokio::test]
1027 async fn second_register_already_existed_bumps_last_seen_only() -> Result<(), StoreError> {
1028 let store = InMemoryStore::default();
1029
1030 let first = store
1031 .register_namespace("orders", NamespaceOrigin::WorkerMint)
1032 .await?;
1033 assert_eq!(first, MintOutcome::Created);
1034 let original = store
1035 .get_namespace("orders")
1036 .await?
1037 .expect("namespace must persist");
1038
1039 let second = store
1041 .register_namespace("orders", NamespaceOrigin::Explicit)
1042 .await?;
1043 assert_eq!(second, MintOutcome::AlreadyExisted);
1044
1045 let touched = store
1046 .get_namespace("orders")
1047 .await?
1048 .expect("namespace must persist");
1049 assert_eq!(touched.created_at, original.created_at);
1050 assert_eq!(touched.origin, NamespaceOrigin::WorkerMint);
1051 assert!(touched.last_seen >= original.last_seen);
1052 Ok(())
1053 }
1054
1055 #[tokio::test]
1056 async fn put_namespace_is_idempotent_on_existing_name() -> Result<(), StoreError> {
1057 let store = InMemoryStore::default();
1058 let now = Utc
1059 .with_ymd_and_hms(2026, 6, 30, 12, 0, 0)
1060 .single()
1061 .expect("valid instant");
1062
1063 let mut record = NamespaceRecord::new_minted("billing", NamespaceOrigin::Explicit, now);
1064 record.config.kind = Some("tenant".to_owned());
1065
1066 let created = store.put_namespace(record.clone()).await?;
1067 assert_eq!(created, MintOutcome::Created);
1068
1069 let mut replacement =
1072 NamespaceRecord::new_minted("billing", NamespaceOrigin::WorkerMint, now);
1073 replacement.config.kind = None;
1074 let again = store.put_namespace(replacement).await?;
1075 assert_eq!(again, MintOutcome::AlreadyExisted);
1076
1077 let stored = store
1078 .get_namespace("billing")
1079 .await?
1080 .expect("namespace must persist");
1081 assert_eq!(stored.origin, NamespaceOrigin::Explicit);
1082 assert_eq!(stored.config.kind.as_deref(), Some("tenant"));
1083 Ok(())
1084 }
1085
1086 #[tokio::test]
1087 async fn list_orders_by_created_at_then_name() -> Result<(), StoreError> {
1088 let store = InMemoryStore::default();
1089 let earlier = Utc
1090 .with_ymd_and_hms(2026, 6, 30, 12, 0, 0)
1091 .single()
1092 .expect("valid instant");
1093 let later = Utc
1094 .with_ymd_and_hms(2026, 6, 30, 13, 0, 0)
1095 .single()
1096 .expect("valid instant");
1097
1098 store
1100 .put_namespace(NamespaceRecord::new_minted(
1101 "zeta",
1102 NamespaceOrigin::Explicit,
1103 earlier,
1104 ))
1105 .await?;
1106 store
1107 .put_namespace(NamespaceRecord::new_minted(
1108 "alpha",
1109 NamespaceOrigin::Explicit,
1110 earlier,
1111 ))
1112 .await?;
1113 store
1114 .put_namespace(NamespaceRecord::new_minted(
1115 "beta",
1116 NamespaceOrigin::Explicit,
1117 later,
1118 ))
1119 .await?;
1120
1121 let listed: Vec<String> = store
1122 .list_namespaces()
1123 .await?
1124 .into_iter()
1125 .map(|record| record.name)
1126 .collect();
1127
1128 assert_eq!(listed, vec!["alpha", "zeta", "beta"]);
1129 Ok(())
1130 }
1131
1132 #[tokio::test]
1133 async fn get_returns_none_for_absent_name() -> Result<(), StoreError> {
1134 let store = InMemoryStore::default();
1135 assert!(store.get_namespace("missing").await?.is_none());
1136 Ok(())
1137 }
1138
1139 #[tokio::test]
1140 async fn deprecate_sets_state_and_is_idempotent() -> Result<(), StoreError> {
1141 let store = InMemoryStore::default();
1142 store
1143 .register_namespace("orders", NamespaceOrigin::WorkerMint)
1144 .await?;
1145
1146 store.deprecate_namespace("orders").await?;
1147 let deprecated = store
1148 .get_namespace("orders")
1149 .await?
1150 .expect("namespace must persist");
1151 assert_eq!(deprecated.state, NamespaceState::Deprecated);
1152
1153 store.deprecate_namespace("orders").await?;
1155 let still = store
1156 .get_namespace("orders")
1157 .await?
1158 .expect("namespace must persist");
1159 assert_eq!(still.state, NamespaceState::Deprecated);
1160
1161 store.deprecate_namespace("never-seen").await?;
1163 assert!(store.get_namespace("never-seen").await?.is_none());
1164 Ok(())
1165 }
1166}