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