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
23pub 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 pub async fn collection(&self, id: &CollectionId) -> Result<StorageCollectionWrapper, RetrievalError> {
86 self.wait_loaded().await;
87 self.0.collectionset.get(id).await
91 }
92
93 pub fn is_system_ready(&self) -> bool { *self.0.system_ready.read().unwrap() }
95
96 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 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 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 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 let event_getter = LocalEventGetter::new(storage.clone(), true);
133 event_getter.stage_event(event.clone());
134
135 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 let attested_state: Attested<EntityState> = system_entity.to_entity_state()?.into();
141 storage.set_state(attested_state.clone()).await?;
142
143 let mut items = self.0.items.write().unwrap();
145 items.push(system_entity);
146 *self.0.root.write().unwrap() = Some(attested_state);
147
148 *self.0.system_ready.write().unwrap() = true;
150 self.0.system_ready_notify.notify_waiters();
151
152 Ok(())
153 }
154
155 pub async fn join_system(&self, state: Attested<EntityState>) -> Result<(), MutationError> {
157 self.wait_loaded().await;
159
160 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 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 tracing::info!("Resetting storage to replace mismatched root");
180 {
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 storage.set_state(state.clone()).await?;
193
194 {
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 pub async fn hard_reset(&self) -> Result<()> {
209 self.0.collectionset.delete_all_collections().await?;
211
212 {
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 self.0.reactor.system_reset();
232
233 Ok(())
234 }
235
236 pub fn is_loaded(&self) -> bool { self.0.loaded.get().is_some() }
238
239 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 {
281 let mut items = self.0.items.write().unwrap();
282 items.extend(entities);
283 }
284
285 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 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 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}