Skip to main content

aion_store/
memory.rs

1//! `InMemoryStore` reference implementation and behavioural test suite.
2
3use 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/// Correct non-durable [`crate::EventStore`] implementation for tests and backend equivalence.
23#[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    /// Poison-tolerant lock over the namespace registry map.
99    ///
100    /// The namespace operations have no fallible body beyond the lock, so a
101    /// poisoned guard is recovered in place (matching the recording test-double
102    /// pattern in `testing.rs`) rather than surfaced as a `StoreError`.
103    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            // Idempotent on an existing name: reconcile as already-existing
222            // rather than overwriting the durable record wholesale. Only the
223            // staleness signal is refreshed.
224            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            // Placement targets an already-minted namespace: an absent row is a
257            // not-found the caller surfaces, never a silent mint here.
258            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        // A missing row is an idempotent no-op: deprecation never strands
271        // durable history and an absent registry entry is not an error.
272        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    /// #204 GATE-2 recovery mechanism: a paused run is EXCLUDED from `list_active`
611    /// (== Running) — so it is not respawned on restart — and is the sole content
612    /// of `list_paused` (== Paused), the durable source the dispatch-hold is
613    /// rebuilt from.
614    #[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        // run_id comes from the WorkflowStarted event (hardcoded to from_u128(1) in the helper)
704        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    /// `set_namespace_placement` updates ONLY the placement (+ `last_seen`) of an
955    /// existing record, leaves the other fields untouched, is idempotent, and is a
956    /// not-found (`Ok(None)`) for an absent namespace — never a silent mint.
957    #[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        // Only placement + last_seen changed; identity/lifecycle preserved.
985        assert_eq!(updated.origin, original.origin);
986        assert_eq!(updated.created_at, original.created_at);
987        assert_eq!(updated.state, original.state);
988
989        // Idempotent: re-applying the same placement is a successful no-op.
990        assert_eq!(
991            store.set_namespace_placement("orders", placement).await?,
992            Some(())
993        );
994
995        // Absent namespace: not-found, and nothing minted.
996        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        // A different origin on re-register must NOT overwrite the recorded origin.
1040        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        // A second put with a DIFFERENT record body must reconcile as
1070        // AlreadyExisted and must not overwrite the stored record wholesale.
1071        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        // Two share `earlier` (tiebreak by name), one is `later`.
1099        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        // Deprecating again is a no-op, not an error.
1154        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        // Deprecating an absent namespace is also an idempotent no-op.
1162        store.deprecate_namespace("never-seen").await?;
1163        assert!(store.get_namespace("never-seen").await?.is_none());
1164        Ok(())
1165    }
1166}