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::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/// Correct non-durable [`crate::EventStore`] implementation for tests and backend equivalence.
28#[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    /// Poison-tolerant lock over the namespace registry map.
105    ///
106    /// The namespace operations have no fallible body beyond the lock, so a
107    /// poisoned guard is recovered in place (matching the recording test-double
108    /// pattern in `testing.rs`) rather than surfaced as a `StoreError`.
109    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    /// Write backend bytes directly for persistence conformance testing.
122    ///
123    /// This deliberately bypasses the entity codec so the shared suite can pin
124    /// poisoned-row listing, replacement, and deletion behavior on every backend.
125    ///
126    /// # Errors
127    ///
128    /// Returns [`StoreError::Serialization`] when `name` is empty or whitespace.
129    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            // Idempotent on an existing name: reconcile as already-existing
257            // rather than overwriting the durable record wholesale. Only the
258            // staleness signal is refreshed.
259            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            // Placement targets an already-minted namespace: an absent row is a
292            // not-found the caller surfaces, never a silent mint here.
293            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        // A missing row is an idempotent no-op: deprecation never strands
306        // durable history and an absent registry entry is not an error.
307        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    /// #204 GATE-2 recovery mechanism: a paused run is EXCLUDED from `list_active`
733    /// (== Running) — so it is not respawned on restart — and is the sole content
734    /// of `list_paused` (== Paused), the durable source the dispatch-hold is
735    /// rebuilt from.
736    #[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        // run_id comes from the WorkflowStarted event (hardcoded to from_u128(1) in the helper)
826        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    /// `set_namespace_placement` updates ONLY the placement (+ `last_seen`) of an
1077    /// existing record, leaves the other fields untouched, is idempotent, and is a
1078    /// not-found (`Ok(None)`) for an absent namespace — never a silent mint.
1079    #[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        // Only placement + last_seen changed; identity/lifecycle preserved.
1107        assert_eq!(updated.origin, original.origin);
1108        assert_eq!(updated.created_at, original.created_at);
1109        assert_eq!(updated.state, original.state);
1110
1111        // Idempotent: re-applying the same placement is a successful no-op.
1112        assert_eq!(
1113            store.set_namespace_placement("orders", placement).await?,
1114            Some(())
1115        );
1116
1117        // Absent namespace: not-found, and nothing minted.
1118        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        // A different origin on re-register must NOT overwrite the recorded origin.
1162        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        // A second put with a DIFFERENT record body must reconcile as
1192        // AlreadyExisted and must not overwrite the stored record wholesale.
1193        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        // Two share `earlier` (tiebreak by name), one is `later`.
1221        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        // Deprecating again is a no-op, not an error.
1276        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        // Deprecating an absent namespace is also an idempotent no-op.
1284        store.deprecate_namespace("never-seen").await?;
1285        assert!(store.get_namespace("never-seen").await?.is_none());
1286        Ok(())
1287    }
1288}