1extern crate core;
2
3pub mod aof;
4pub mod indexing;
5pub mod networking;
6pub mod read_api;
7
8use std::{
9 collections::HashMap,
10 sync::{Arc, Mutex, OnceLock},
11 thread::JoinHandle,
12};
13
14use rusted_ring::{
15 L_CAPACITY, L_TSHIRT_SIZE, M_CAPACITY, M_TSHIRT_SIZE, PooledEvent, RingBuffer, S_CAPACITY, S_TSHIRT_SIZE, Writer, XL_CAPACITY, XL_TSHIRT_SIZE, XS_CAPACITY, XS_TSHIRT_SIZE,
16};
17use xaeroid::XaeroID;
18
19use crate::{
20 aof::{ring_buffer_actor::AofActor, storage::lmdb::LmdbEnv},
21 networking::p2p::P2pActor,
22};
23pub static XS_RING: OnceLock<RingBuffer<XS_TSHIRT_SIZE, XS_CAPACITY>> = OnceLock::new();
28pub static S_RING: OnceLock<RingBuffer<S_TSHIRT_SIZE, S_CAPACITY>> = OnceLock::new();
29pub static M_RING: OnceLock<RingBuffer<M_TSHIRT_SIZE, M_CAPACITY>> = OnceLock::new();
30pub static L_RING: OnceLock<RingBuffer<L_TSHIRT_SIZE, L_CAPACITY>> = OnceLock::new();
31pub static XL_RING: OnceLock<RingBuffer<XL_TSHIRT_SIZE, XL_CAPACITY>> = OnceLock::new();
32
33pub static P2P_XS_RING: OnceLock<RingBuffer<XS_TSHIRT_SIZE, XS_CAPACITY>> = OnceLock::new();
38pub static P2P_S_RING: OnceLock<RingBuffer<S_TSHIRT_SIZE, S_CAPACITY>> = OnceLock::new();
39pub static P2P_M_RING: OnceLock<RingBuffer<M_TSHIRT_SIZE, M_CAPACITY>> = OnceLock::new();
40pub static P2P_L_RING: OnceLock<RingBuffer<L_TSHIRT_SIZE, L_CAPACITY>> = OnceLock::new();
41pub static P2P_XL_RING: OnceLock<RingBuffer<XL_TSHIRT_SIZE, XL_CAPACITY>> = OnceLock::new();
42
43pub struct EventBus {
48 xs_writer: Writer<XS_TSHIRT_SIZE, XS_CAPACITY>,
49 s_writer: Writer<S_TSHIRT_SIZE, S_CAPACITY>,
50 m_writer: Writer<M_TSHIRT_SIZE, M_CAPACITY>,
51 l_writer: Writer<L_TSHIRT_SIZE, L_CAPACITY>,
52 xl_writer: Writer<XL_TSHIRT_SIZE, XL_CAPACITY>,
53}
54
55impl Default for EventBus {
56 fn default() -> Self {
57 Self::new()
58 }
59}
60
61use std::cell::RefCell;
62
63impl EventBus {
64 pub fn new() -> Self {
66 let xs_ring = XS_RING.get_or_init(RingBuffer::new);
67 let s_ring = S_RING.get_or_init(RingBuffer::new);
68 let m_ring = M_RING.get_or_init(RingBuffer::new);
69 let l_ring = L_RING.get_or_init(RingBuffer::new);
70 let xl_ring = XL_RING.get_or_init(RingBuffer::new);
71
72 Self {
73 xs_writer: Writer::new(xs_ring),
74 s_writer: Writer::new(s_ring),
75 m_writer: Writer::new(m_ring),
76 l_writer: Writer::new(l_ring),
77 xl_writer: Writer::new(xl_ring),
78 }
79 }
80
81 pub fn write_xs(&mut self, event: PooledEvent<XS_TSHIRT_SIZE>) {
83 self.xs_writer.add(event);
84 }
85
86 pub fn write_s(&mut self, event: PooledEvent<S_TSHIRT_SIZE>) {
88 self.s_writer.add(event);
89 }
90
91 pub fn write_m(&mut self, event: PooledEvent<M_TSHIRT_SIZE>) {
93 self.m_writer.add(event);
94 }
95
96 pub fn write_l(&mut self, event: PooledEvent<L_TSHIRT_SIZE>) {
98 self.l_writer.add(event);
99 }
100
101 pub fn write_xl(&mut self, event: PooledEvent<XL_TSHIRT_SIZE>) {
103 self.xl_writer.add(event);
104 }
105
106 pub fn write_optimal(&mut self, data: &[u8], event_type: u32) -> Result<(), XaeroFluxError> {
108 let data_len = data.len();
109
110 if data_len <= XS_TSHIRT_SIZE {
111 let event = Self::create_pooled_event::<XS_TSHIRT_SIZE>(data, event_type)?;
112 self.write_xs(event);
113 } else if data_len <= S_TSHIRT_SIZE {
114 let event = Self::create_pooled_event::<S_TSHIRT_SIZE>(data, event_type)?;
115 self.write_s(event);
116 } else if data_len <= M_TSHIRT_SIZE {
117 let event = Self::create_pooled_event::<M_TSHIRT_SIZE>(data, event_type)?;
118 self.write_m(event);
119 } else if data_len <= L_TSHIRT_SIZE {
120 let event = Self::create_pooled_event::<L_TSHIRT_SIZE>(data, event_type)?;
121 self.write_l(event);
122 } else if data_len <= XL_TSHIRT_SIZE {
123 let event = Self::create_pooled_event::<XL_TSHIRT_SIZE>(data, event_type)?;
124 self.write_xl(event);
125 } else {
126 return Err(XaeroFluxError::DataTooLarge(data_len));
127 }
128
129 Ok(())
130 }
131
132 fn create_pooled_event<const SIZE: usize>(data: &[u8], event_type: u32) -> Result<PooledEvent<SIZE>, XaeroFluxError> {
134 if data.len() > SIZE {
135 return Err(XaeroFluxError::DataTooLarge(data.len()));
136 }
137
138 let mut event_data = [0u8; SIZE];
139 event_data[..data.len()].copy_from_slice(data);
140
141 Ok(PooledEvent {
142 data: event_data,
143 len: data.len() as u32,
144 event_type,
145 })
146 }
147}
148
149pub struct P2PRingAccess;
150
151impl P2PRingAccess {
152 pub fn xs_writer() -> Writer<XS_TSHIRT_SIZE, XS_CAPACITY> {
154 let ring = P2P_XS_RING.get_or_init(RingBuffer::new);
155 Writer::new(ring)
156 }
157
158 pub fn s_writer() -> Writer<S_TSHIRT_SIZE, S_CAPACITY> {
160 let ring = P2P_S_RING.get_or_init(RingBuffer::new);
161 Writer::new(ring)
162 }
163
164 pub fn m_writer() -> Writer<M_TSHIRT_SIZE, M_CAPACITY> {
166 let ring = P2P_M_RING.get_or_init(RingBuffer::new);
167 Writer::new(ring)
168 }
169
170 pub fn l_writer() -> Writer<L_TSHIRT_SIZE, L_CAPACITY> {
172 let ring = P2P_L_RING.get_or_init(RingBuffer::new);
173 Writer::new(ring)
174 }
175
176 pub fn xl_writer() -> Writer<XL_TSHIRT_SIZE, XL_CAPACITY> {
178 let ring = P2P_XL_RING.get_or_init(RingBuffer::new);
179 Writer::new(ring)
180 }
181}
182
183pub trait VectorExtractor: Send + Sync {
188 fn extract_vector(&self, event_data: &[u8]) -> Option<Vec<f32>>;
189}
190
191#[derive(Debug, Clone)]
192pub struct VectorSearchStats {
193 pub total_indexed: usize,
194 pub active_nodes: usize,
195}
196
197use crate::indexing::vec_search_actor::{VectorQueryRequest, VectorQueryResponse, VectorSearchActor};
198
199pub struct XaeroFlux {
200 pub event_bus: EventBus,
201 pub vector_search: Option<Arc<VectorSearchActor>>,
202 pub aof_handle: Option<JoinHandle<()>>,
203 pub p2p_handle: Option<JoinHandle<()>>,
204 pub read_handle: Option<Arc<Mutex<LmdbEnv>>>,
205}
206
207impl Default for XaeroFlux {
208 fn default() -> Self {
209 Self::new()
210 }
211}
212static XAERO_FLUX: OnceLock<XaeroFlux> = OnceLock::new();
213impl XaeroFlux {
214 pub fn instance() -> Option<&'static XaeroFlux> {
215 XAERO_FLUX.get()
216 }
217
218 fn new() -> Self {
220 Self {
221 event_bus: EventBus::new(),
222 vector_search: None,
223 aof_handle: None,
224 p2p_handle: None,
225 read_handle: None,
226 }
227 }
228
229 pub fn read_handle() -> Option<Arc<Mutex<LmdbEnv>>> {
230 XAERO_FLUX.get().and_then(|xf| xf.read_handle.clone())
231 }
232
233 pub fn start_aof(&mut self) -> Result<(), Box<dyn std::error::Error>> {
235 let aof_actor = AofActor::spin()?;
236 self.aof_handle = Some(aof_actor.jh);
237 self.read_handle = Some(aof_actor.env);
238 Ok(())
239 }
240
241 pub fn start_p2p(&mut self, xaero_id: XaeroID) -> Result<(), Box<dyn std::error::Error>> {
243 if self.aof_handle.is_none() {
245 return Err("AOF must be started before P2P".into());
246 }
247 let p2p_handle = std::thread::spawn(move || {
248 let rt = tokio::runtime::Handle::current();
249 let handle = rt.spawn(async move {
250 let s_ring: &'static RingBuffer<S_TSHIRT_SIZE, S_CAPACITY> = S_RING.get_or_init(RingBuffer::new);
252
253 let aof_state = Arc::new(crate::aof::ring_buffer_actor::AofState::new().expect("failed to create ring buffer actor"));
255
256 match P2pActor::<S_TSHIRT_SIZE, S_CAPACITY>::new(s_ring, xaero_id, aof_state).await {
257 Ok((mut actor, writer, reader)) =>
258 if let Err(e) = actor.start(writer, reader).await {
259 tracing::error!("P2P actor failed: {:?}", e);
260 },
261 Err(e) => {
262 tracing::error!("Failed to create P2P actor: {:?}", e);
263 }
264 }
265 });
266 });
267
268 self.p2p_handle = Some(p2p_handle);
269 Ok(())
270 }
271
272 pub fn start_vector_search(
274 &mut self,
275 extractors: HashMap<u32, Box<dyn crate::indexing::vec_search_actor::VectorExtractor>>,
276 vector_dimension: usize,
277 max_nb_connection: usize,
278 max_elements: usize,
279 max_layer: usize,
280 ef_construction: usize,
281 ) -> Result<(), Box<dyn std::error::Error>> {
282 let actor = VectorSearchActor::spin(max_nb_connection, max_elements, max_layer, ef_construction, extractors, vector_dimension)?;
283 self.vector_search = Some(Arc::new(actor));
284 Ok(())
285 }
286
287 pub fn write_event(&mut self, data: &[u8], event_type: u32) -> Result<(), XaeroFluxError> {
289 self.event_bus.write_optimal(data, event_type)
290 }
291
292 pub fn send_text(&mut self, text: &str) -> Result<(), XaeroFluxError> {
294 self.write_event(text.as_bytes(), 1) }
296
297 pub fn send_file_data(&mut self, file_data: &[u8]) -> Result<(), XaeroFluxError> {
299 self.write_event(file_data, 2) }
301
302 pub fn search_vector(&self, vector: Vec<f32>, k: u32, similarity_threshold: f32) -> Result<VectorQueryResponse<5>, XaeroFluxError> {
308 let vector_search = self.vector_search.as_ref().ok_or(XaeroFluxError::VectorSearchNotStarted)?;
309 let mut query_vector = [0.0f32; 256];
310 let copy_len = std::cmp::min(vector.len(), 256);
311 query_vector[..copy_len].copy_from_slice(&vector[..copy_len]);
312
313 let query = VectorQueryRequest {
314 query_id: xaeroflux_core::date_time::emit_secs(),
315 requester_id: [0; 32],
316 scope: crate::indexing::vec_search_actor::QueryScope {
317 group_id: None,
318 workspace_id: None,
319 object_id: None,
320 },
321 vector: query_vector,
322 k,
323 similarity_threshold,
324 time_window: xaeroflux_core::event::ScanWindow { start: 0, end: u64::MAX },
325 flags: crate::indexing::vec_search_actor::QueryFlags {
326 include_metadata: true,
327 include_operations: false,
328 fan_out: false,
329 use_lora_bias: false,
330 },
331 };
332
333 Ok(vector_search.search(&query))
334 }
335
336 pub fn search_vectors(&self, vectors: Vec<Vec<f32>>, k: u32, similarity_threshold: f32) -> Result<Vec<VectorQueryResponse<5>>, XaeroFluxError> {
338 let mut results = Vec::new();
339
340 for vector in vectors {
341 let result = self.search_vector(vector, k, similarity_threshold)?;
342 results.push(result);
343 }
344
345 Ok(results)
346 }
347
348 pub fn vector_search_stats(&self) -> Result<VectorSearchStats, XaeroFluxError> {
350 let vector_search = self.vector_search.as_ref().ok_or(XaeroFluxError::VectorSearchNotStarted)?;
351
352 let (total_indexed, active_nodes) = vector_search.get_stats();
353 Ok(VectorSearchStats { total_indexed, active_nodes })
354 }
355
356 pub fn initialize(xaero_id: XaeroID) -> Result<(), Box<dyn std::error::Error>> {
357 let mut xf = XaeroFlux::new();
358 xf.start_aof()?;
359 xf.start_p2p(xaero_id)?;
360 XAERO_FLUX.set(xf).map_err(|_| "Already initialized")?;
361 Ok(())
362 }
363
364 pub fn write_event_static(data: &[u8], event_type: u32) -> Result<(), XaeroFluxError> {
365 thread_local! {
366 static LOCAL_EVENT_BUS: RefCell<Option<EventBus>> = const { RefCell::new(None) };
367 }
368 LOCAL_EVENT_BUS.with(|bus_cell| {
369 let mut bus_ref = bus_cell.borrow_mut();
370
371 if bus_ref.is_none() {
373 *bus_ref = Some(EventBus::new());
374 }
375
376 bus_ref.as_mut().unwrap().write_optimal(data, event_type)
378 })
379 }
380}
381
382#[derive(Debug, Clone)]
387pub enum XaeroFluxError {
388 VectorSearchNotStarted,
389 DataTooLarge(usize),
390 InvalidData,
391 ActorError(String),
392}
393
394impl std::fmt::Display for XaeroFluxError {
395 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
396 match self {
397 XaeroFluxError::VectorSearchNotStarted => write!(f, "Vector search not started"),
398 XaeroFluxError::DataTooLarge(size) => write!(f, "Data too large: {} bytes", size),
399 XaeroFluxError::InvalidData => write!(f, "Invalid data"),
400 XaeroFluxError::ActorError(msg) => write!(f, "Actor error: {}", msg),
401 }
402 }
403}
404
405impl std::error::Error for XaeroFluxError {}
406
407#[cfg(test)]
412mod tests {
413 use super::*;
414
415 #[test]
416 fn test_event_bus_creation() {
417 let bus = EventBus::new();
418 println!("✅ EventBus created with writers");
420 }
421
422 #[test]
423 fn test_event_bus_write_optimal() {
424 let mut bus = EventBus::new();
425
426 let test_data = b"hello world";
427 let result = bus.write_optimal(test_data, 42);
428 assert!(result.is_ok());
429
430 println!("✅ EventBus write_optimal works");
431 }
432
433 #[test]
434 fn test_p2p_ring_access() {
435 let _xs_writer = P2PRingAccess::xs_writer();
437 let _s_writer = P2PRingAccess::s_writer();
438 let _m_writer = P2PRingAccess::m_writer();
439 let _l_writer = P2PRingAccess::l_writer();
440 let _xl_writer = P2PRingAccess::xl_writer();
441
442 println!("✅ P2P ring access works");
443 }
444
445 #[test]
446 fn test_xaeroflux_creation() {
447 let xf = XaeroFlux::new();
448 assert!(xf.vector_search.is_none());
449 assert!(xf.aof_handle.is_none());
450 assert!(xf.p2p_handle.is_none());
451
452 println!("✅ XaeroFlux created successfully");
453 }
454
455 #[test]
456 fn test_xaeroflux_write_event() {
457 let mut xf = XaeroFlux::new();
458
459 let test_data = b"test event data";
460 let result = xf.write_event(test_data, 42);
461 assert!(result.is_ok());
462
463 println!("✅ XaeroFlux write_event works");
464 }
465
466 #[test]
467 fn test_send_text() {
468 let mut xf = XaeroFlux::new();
469
470 let result = xf.send_text("Hello P2P world!");
471 assert!(result.is_ok());
472
473 println!("✅ XaeroFlux send_text works");
474 }
475
476 #[test]
477 fn test_oversized_data() {
478 let mut xf = XaeroFlux::new();
479
480 let oversized_data = vec![0u8; 20000]; let result = xf.write_event(&oversized_data, 42);
482 assert!(matches!(result, Err(XaeroFluxError::DataTooLarge(_))));
483
484 println!("✅ Oversized data properly rejected");
485 }
486}