Skip to main content

backbone_core/
aggregate.rs

1//! DDD Aggregate Root Pattern
2//!
3//! Provides traits for implementing Aggregate Roots in Domain-Driven Design.
4//! Aggregate Roots are the entry point to a cluster of domain objects
5//! and are responsible for maintaining invariants.
6//!
7//! # Example
8//!
9//! ```ignore
10//! use backbone_core::{AggregateRoot, PersistentEntity};
11//!
12//! #[derive(Clone)]
13//! pub struct Order {
14//!     id: String,
15//!     items: Vec<OrderItem>,
16//!     status: OrderStatus,
17//!     events: Vec<OrderEvent>,
18//!     // ... other fields
19//! }
20//!
21//! impl AggregateRoot for Order {
22//!     type Event = OrderEvent;
23//!
24//!     fn uncommitted_events(&self) -> &[Self::Event] {
25//!         &self.events
26//!     }
27//!
28//!     fn clear_events(&mut self) {
29//!         self.events.clear();
30//!     }
31//!
32//!     fn apply_event(&mut self, event: Self::Event) {
33//!         match event {
34//!             OrderEvent::ItemAdded { item } => self.items.push(item),
35//!             OrderEvent::Confirmed => self.status = OrderStatus::Confirmed,
36//!             // ...
37//!         }
38//!         self.events.push(event);
39//!     }
40//! }
41//! ```
42
43use crate::persistence::PersistentEntity;
44
45/// DDD Aggregate Root trait.
46///
47/// Aggregate Roots are characterized by:
48/// - Being the only entry point to the aggregate
49/// - Maintaining consistency boundaries
50/// - Publishing domain events for state changes
51/// - Having a unique identity (via `PersistentEntity`)
52pub trait AggregateRoot: PersistentEntity {
53    /// The type of domain events this aggregate can produce.
54    type Event: Send + Sync + Clone;
55
56    /// Get all uncommitted domain events.
57    ///
58    /// These events should be persisted/published after the aggregate
59    /// is saved to the repository.
60    fn uncommitted_events(&self) -> &[Self::Event];
61
62    /// Clear all uncommitted events.
63    ///
64    /// This should be called after events have been persisted/published.
65    fn clear_events(&mut self);
66
67    /// Apply a domain event to update the aggregate state.
68    ///
69    /// This method should:
70    /// 1. Update the aggregate's internal state based on the event
71    /// 2. Store the event in the uncommitted events list
72    fn apply_event(&mut self, event: Self::Event);
73
74    /// Get the aggregate version for optimistic concurrency.
75    ///
76    /// Returns `None` if versioning is not supported.
77    fn version(&self) -> Option<u64> {
78        None
79    }
80
81    /// Increment the aggregate version.
82    ///
83    /// Called after successful persistence.
84    fn increment_version(&mut self) {
85        // Default implementation does nothing
86    }
87
88    /// Check if there are any uncommitted events.
89    fn has_uncommitted_events(&self) -> bool {
90        !self.uncommitted_events().is_empty()
91    }
92
93    /// Get the number of uncommitted events.
94    fn uncommitted_event_count(&self) -> usize {
95        self.uncommitted_events().len()
96    }
97}
98
99/// Extension trait for aggregates that support event sourcing.
100pub trait EventSourcedAggregate: AggregateRoot {
101    /// Reconstruct the aggregate from a stream of events.
102    fn from_events(id: String, events: impl IntoIterator<Item = Self::Event>) -> Self
103    where
104        Self: Sized;
105
106    /// Get all events (including committed ones) for event sourcing.
107    fn all_events(&self) -> Vec<Self::Event>;
108}
109
110/// Trait for aggregates that enforce invariants.
111pub trait InvariantAggregate: AggregateRoot {
112    /// Error type for invariant violations.
113    type InvariantError: std::error::Error + Send + Sync;
114
115    /// Check all aggregate invariants.
116    ///
117    /// This should be called before persisting the aggregate.
118    fn check_invariants(&self) -> Result<(), Self::InvariantError>;
119}
120
121/// Helper struct for tracking aggregate metadata.
122#[derive(Debug, Clone, Default)]
123pub struct AggregateMetadata {
124    /// Current version of the aggregate.
125    pub version: u64,
126    /// Timestamp of last modification.
127    pub last_modified: Option<chrono::DateTime<chrono::Utc>>,
128    /// ID of the user who last modified the aggregate.
129    pub last_modified_by: Option<String>,
130}
131
132impl AggregateMetadata {
133    /// Create new metadata with version 0.
134    pub fn new() -> Self {
135        Self::default()
136    }
137
138    /// Create metadata with a specific version.
139    pub fn with_version(version: u64) -> Self {
140        Self {
141            version,
142            ..Default::default()
143        }
144    }
145
146    /// Increment the version.
147    pub fn increment(&mut self) {
148        self.version += 1;
149        self.last_modified = Some(chrono::Utc::now());
150    }
151}
152
153#[cfg(test)]
154mod tests {
155    use super::*;
156    use chrono::{DateTime, Utc};
157    use serde::{Deserialize, Serialize};
158
159    #[derive(Clone, Debug, Serialize, Deserialize)]
160    struct TestEvent {
161        data: String,
162    }
163
164    #[derive(Clone, Debug, Serialize, Deserialize)]
165    struct TestAggregate {
166        id: String,
167        data: String,
168        events: Vec<TestEvent>,
169        created_at: DateTime<Utc>,
170        updated_at: DateTime<Utc>,
171        deleted_at: Option<DateTime<Utc>>,
172    }
173
174    impl PersistentEntity for TestAggregate {
175        fn entity_id(&self) -> String {
176            self.id.clone()
177        }
178
179        fn set_entity_id(&mut self, id: String) {
180            self.id = id;
181        }
182
183        fn created_at(&self) -> Option<DateTime<Utc>> {
184            Some(self.created_at)
185        }
186
187        fn set_created_at(&mut self, at: DateTime<Utc>) {
188            self.created_at = at;
189        }
190
191        fn updated_at(&self) -> Option<DateTime<Utc>> {
192            Some(self.updated_at)
193        }
194
195        fn set_updated_at(&mut self, at: DateTime<Utc>) {
196            self.updated_at = at;
197        }
198
199        fn deleted_at(&self) -> Option<DateTime<Utc>> {
200            self.deleted_at
201        }
202
203        fn set_deleted_at(&mut self, at: Option<DateTime<Utc>>) {
204            self.deleted_at = at;
205        }
206    }
207
208    impl AggregateRoot for TestAggregate {
209        type Event = TestEvent;
210
211        fn uncommitted_events(&self) -> &[Self::Event] {
212            &self.events
213        }
214
215        fn clear_events(&mut self) {
216            self.events.clear();
217        }
218
219        fn apply_event(&mut self, event: Self::Event) {
220            self.data = event.data.clone();
221            self.events.push(event);
222        }
223    }
224
225    #[test]
226    fn test_aggregate_events() {
227        let mut aggregate = TestAggregate {
228            id: "1".to_string(),
229            data: "initial".to_string(),
230            events: vec![],
231            created_at: Utc::now(),
232            updated_at: Utc::now(),
233            deleted_at: None,
234        };
235
236        assert!(!aggregate.has_uncommitted_events());
237
238        aggregate.apply_event(TestEvent {
239            data: "updated".to_string(),
240        });
241
242        assert!(aggregate.has_uncommitted_events());
243        assert_eq!(aggregate.uncommitted_event_count(), 1);
244        assert_eq!(aggregate.data, "updated");
245
246        aggregate.clear_events();
247        assert!(!aggregate.has_uncommitted_events());
248    }
249
250    #[test]
251    fn test_aggregate_metadata() {
252        let mut metadata = AggregateMetadata::new();
253        assert_eq!(metadata.version, 0);
254
255        metadata.increment();
256        assert_eq!(metadata.version, 1);
257        assert!(metadata.last_modified.is_some());
258    }
259}