Skip to main content

ankurah_core/
system.rs

1use ankurah_proto::{self as proto, Attested, CollectionId, EntityState, Event};
2use anyhow::{anyhow, Result};
3use std::collections::BTreeMap;
4use std::marker::PhantomData;
5use std::sync::{Arc, OnceLock, RwLock};
6use tokio::sync::Notify;
7use tracing::{error, warn};
8
9use crate::collectionset::CollectionSet;
10use crate::entity::{Entity, WeakEntitySet};
11use crate::error::MutationError;
12use crate::error::RetrievalError;
13use crate::notice_info;
14use crate::policy::PolicyAgent;
15use crate::property::{Property, PropertyError};
16use crate::reactor::Reactor;
17use crate::retrieval::{LocalEventGetter, LocalStateGetter, SuspenseEvents};
18use crate::storage::{StorageCollectionWrapper, StorageEngine};
19use crate::{property::backend::LWWBackend, value::Value};
20pub const SYSTEM_COLLECTION_ID: &str = "_ankurah_system";
21pub const PROTECTED_COLLECTIONS: &[&str] = &[SYSTEM_COLLECTION_ID];
22
23/// System catalog manager for storing various metadata about the system
24/// * root clock
25/// * valid collections (TODO)
26/// * property definitions (TODO)
27
28pub struct SystemManager<SE, PA>(Arc<Inner<SE, PA>>);
29impl<SE, PA> Clone for SystemManager<SE, PA> {
30    fn clone(&self) -> Self { Self(self.0.clone()) }
31}
32
33struct Inner<SE, PA> {
34    collectionset: CollectionSet<SE>,
35    collection_map: RwLock<BTreeMap<CollectionId, Entity>>,
36    entities: WeakEntitySet,
37    durable: bool,
38    root: RwLock<Option<Attested<EntityState>>>,
39    items: RwLock<Vec<Entity>>,
40    loaded: OnceLock<()>,
41    loading: Notify,
42    system_ready: RwLock<bool>,
43    system_ready_notify: Notify,
44    reactor: Reactor,
45    _phantom: PhantomData<PA>,
46}
47
48impl<SE, PA> SystemManager<SE, PA>
49where
50    SE: StorageEngine + Send + Sync + 'static,
51    PA: PolicyAgent + Send + Sync + 'static,
52{
53    pub(crate) fn new(collections: CollectionSet<SE>, entities: WeakEntitySet, reactor: Reactor, durable: bool) -> Self {
54        let me = Self(Arc::new(Inner {
55            collectionset: collections,
56            entities,
57            durable,
58            items: RwLock::new(Vec::new()),
59            root: RwLock::new(None),
60            loaded: OnceLock::new(),
61            loading: Notify::new(),
62            collection_map: RwLock::new(BTreeMap::new()),
63            system_ready: RwLock::new(false),
64            system_ready_notify: Notify::new(),
65            reactor,
66            _phantom: PhantomData,
67        }));
68        {
69            let me = me.clone();
70            crate::task::spawn(async move {
71                if let Err(e) = me.load_system_catalog().await {
72                    error!("Failed to load system catalog: {}", e);
73                }
74            });
75        }
76        me
77    }
78
79    pub fn root(&self) -> Option<Attested<EntityState>> { self.0.root.read().unwrap().as_ref().map(|r| r.clone()) }
80
81    pub fn items(&self) -> Vec<Entity> { self.0.items.read().unwrap().clone() }
82
83    /// get an existing collection if it's defined in the system catalog, else insert a SysItem::Collection
84    /// then return collections.get to get the StorageCollectionWrapper
85    pub async fn collection(&self, id: &CollectionId) -> Result<StorageCollectionWrapper, RetrievalError> {
86        self.wait_loaded().await;
87        // TODO - update the system catalog to create an entity for this collection
88
89        // Return the collection wrapper
90        self.0.collectionset.get(id).await
91    }
92
93    /// Returns true if we've successfully initialized or joined a system
94    pub fn is_system_ready(&self) -> bool { *self.0.system_ready.read().unwrap() }
95
96    /// Waits until we've successfully initialized or joined a system
97    pub async fn wait_system_ready(&self) {
98        if !self.is_system_ready() {
99            self.0.system_ready_notify.notified().await;
100        }
101    }
102
103    /// Creates a new system root. This should only be called once per system by durable nodes
104    /// The rest of the nodes must "join" this system.
105    pub async fn create(&self) -> Result<()> {
106        if !self.0.durable {
107            return Err(anyhow!("Only durable nodes can create a new system"));
108        }
109
110        // Wait for local system catalog to be loaded
111        self.wait_loaded().await;
112
113        {
114            let items = self.0.items.read().unwrap();
115            if !items.is_empty() {
116                return Err(anyhow!("System root already exists"));
117            }
118        }
119
120        // TODO - see if we can use the Model derive macro for a SysCatalogItem model rather than doing this manually
121        let collection_id = CollectionId::fixed_name(SYSTEM_COLLECTION_ID);
122        let storage = self.0.collectionset.get(&collection_id).await?;
123
124        let system_entity = self.0.entities.create(collection_id.clone());
125
126        let lww_backend = system_entity.get_backend::<LWWBackend>().expect("LWW Backend should exist");
127        lww_backend.set("item".into(), proto::sys::Item::SysRoot.into_value()?);
128
129        let event = system_entity.generate_commit_event()?.ok_or(anyhow!("Expected event"))?;
130
131        // Stage the event, apply, then commit
132        let event_getter = LocalEventGetter::new(storage.clone(), true);
133        event_getter.stage_event(event.clone());
134
135        // Apply the creation event so LWW values are tagged with event_id before serialization.
136        system_entity.apply_event(&event_getter, &event).await?;
137        let attested_event: Attested<Event> = event.clone().into();
138        event_getter.commit_event(&attested_event).await?;
139        // Now get the entity state after the head is updated
140        let attested_state: Attested<EntityState> = system_entity.to_entity_state()?.into();
141        storage.set_state(attested_state.clone()).await?;
142
143        // Update our system state
144        let mut items = self.0.items.write().unwrap();
145        items.push(system_entity);
146        *self.0.root.write().unwrap() = Some(attested_state);
147
148        // Mark system as ready and notify waiters
149        *self.0.system_ready.write().unwrap() = true;
150        self.0.system_ready_notify.notify_waiters();
151
152        Ok(())
153    }
154
155    /// Joins an existing system. This should only be called by ephemeral nodes.
156    pub async fn join_system(&self, state: Attested<EntityState>) -> Result<(), MutationError> {
157        // Wait for catalog to be loaded before proceeding
158        self.wait_loaded().await;
159
160        // If node is durable, fail - durable nodes should not join an existing system
161        if self.0.durable {
162            warn!("Durable node attempted to join system - this is not allowed");
163            return Err(MutationError::General(Box::new(std::io::Error::other("Durable nodes cannot join an existing system"))));
164        }
165
166        let root_state = self.root();
167
168        // If we have a matching root, we're already in sync - just mark ready and return
169        if let Some(root) = root_state {
170            if root.payload.state.head == state.payload.state.head {
171                notice_info!("Found matching root - Node is part of the same system");
172                *self.0.system_ready.write().unwrap() = true;
173                self.0.system_ready_notify.notify_waiters();
174                return Ok(());
175            }
176            tracing::warn!("Mismatched root state during join: local={:?}, remote={:?}", root, state.payload.state.head);
177
178            // Only reset storage if we have a root that needs to be replaced
179            tracing::info!("Resetting storage to replace mismatched root");
180            // Drop locks before reset
181            {
182                let mut root = self.0.root.write().expect("Root lock poisoned");
183                *root = None;
184            }
185            self.hard_reset().await.map_err(|e| MutationError::General(Box::new(std::io::Error::other(e.to_string()))))?;
186        }
187
188        let collection_id = CollectionId::fixed_name(SYSTEM_COLLECTION_ID);
189        let storage = self.0.collectionset.get(&collection_id).await?;
190
191        // Set the state
192        storage.set_state(state.clone()).await?;
193
194        // Set root and mark system as ready
195        {
196            let mut root = self.0.root.write().expect("Root lock poisoned");
197            *root = Some(state);
198        }
199        *self.0.system_ready.write().unwrap() = true;
200        self.0.system_ready_notify.notify_waiters();
201
202        Ok(())
203    }
204
205    /// Resets all storage by deleting all collections, including the system collection.
206    /// This is used when an ephemeral node needs to join a system with a different root.
207    /// **This is a destructive operation and should be used with extreme caution.**
208    pub async fn hard_reset(&self) -> Result<()> {
209        // Delete all collections from storage
210        self.0.collectionset.delete_all_collections().await?;
211
212        // Reset our state
213        {
214            let mut items = self.0.items.write().unwrap();
215            items.clear();
216        }
217        {
218            let mut root = self.0.root.write().unwrap();
219            *root = None;
220        }
221        {
222            let mut collection_map = self.0.collection_map.write().unwrap();
223            collection_map.clear();
224        }
225        {
226            let mut system_ready = self.0.system_ready.write().unwrap();
227            *system_ready = false;
228        }
229
230        // Reset the reactor state to notify subscriptions
231        self.0.reactor.system_reset();
232
233        Ok(())
234    }
235
236    /// Returns true if the local system catalog is loaded
237    pub fn is_loaded(&self) -> bool { self.0.loaded.get().is_some() }
238
239    /// Waits for the local system catalog to be loaded
240    pub async fn wait_loaded(&self) {
241        if !self.is_loaded() {
242            self.0.loading.notified().await;
243        }
244    }
245
246    async fn load_system_catalog(&self) -> Result<()> {
247        if self.is_loaded() {
248            return Err(anyhow!("System catalog already loaded"));
249        }
250
251        let collection_id = CollectionId::fixed_name(SYSTEM_COLLECTION_ID);
252        let storage = self.0.collectionset.get(&collection_id).await?;
253
254        let mut entities = Vec::new();
255        let mut root_state = None;
256
257        let state_getter = LocalStateGetter::new(storage.clone());
258        let event_getter = LocalEventGetter::new(storage.clone(), self.0.durable);
259
260        for state in
261            storage.fetch_states(&ankql::ast::Selection { predicate: ankql::ast::Predicate::True, order_by: None, limit: None }).await?
262        {
263            let (_entity_changed, entity) = self
264                .0
265                .entities
266                .with_state(&state_getter, &event_getter, state.payload.entity_id, collection_id.clone(), state.payload.state.clone())
267                .await?;
268            let lww_backend = entity.get_backend::<LWWBackend>().expect("LWW Backend should exist");
269            if let Some(value) = lww_backend.get(&"item".to_string()) {
270                let item = proto::sys::Item::from_value(Some(value)).expect("Invalid sys item");
271
272                if let proto::sys::Item::SysRoot = &item {
273                    root_state = Some(state);
274                }
275                entities.push(entity);
276            }
277        }
278
279        // Update our system state
280        {
281            let mut items = self.0.items.write().unwrap();
282            items.extend(entities);
283        }
284
285        // If we loaded a system root and we're a durable node, we're ready
286        let has_root = root_state.is_some();
287        {
288            let mut root = self.0.root.write().expect("Root lock poisoned");
289            *root = root_state;
290        }
291
292        // Only mark ready if we're a durable node and found a root
293        // Ephemeral nodes must explicitly join via join_system()
294        if has_root && self.0.durable {
295            *self.0.system_ready.write().unwrap() = true;
296            self.0.system_ready_notify.notify_waiters();
297        }
298
299        // Set loaded state and notify waiters
300        self.0.loaded.set(()).expect("Loading flag already set");
301        self.0.loading.notify_waiters();
302        Ok(())
303    }
304}
305
306impl Property for proto::sys::Item {
307    fn into_value(&self) -> std::result::Result<Option<Value>, crate::property::PropertyError> {
308        Ok(Some(Value::String(
309            serde_json::to_string(self).map_err(|_| PropertyError::InvalidValue { value: "".to_string(), ty: "sys::Item".to_string() })?,
310        )))
311    }
312
313    fn from_value(value: Option<Value>) -> std::result::Result<Self, crate::property::PropertyError> {
314        if let Some(Value::String(string)) = value {
315            let item: proto::sys::Item = serde_json::from_str(&string)
316                .map_err(|_| PropertyError::InvalidValue { value: "".to_string(), ty: "sys::Item".to_string() })?;
317            Ok(item)
318        } else {
319            Err(PropertyError::InvalidValue { value: "".to_string(), ty: "sys::Item".to_string() })
320        }
321    }
322}