Skip to main content

xaeroflux_actors/
lib.rs

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};
23// ================================================================================================
24// GLOBAL RING BUFFERS - MAIN (EventBus writes to these, AOF/VectorSearch read from these)
25// ================================================================================================
26
27pub 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
33// ================================================================================================
34// GLOBAL RING BUFFERS - P2P (P2P actors write to these - for future use)
35// ================================================================================================
36
37pub 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
43// ================================================================================================
44// EVENT BUS - JUST HOUSES WRITERS
45// ================================================================================================
46
47pub 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    /// Create new EventBus with writers to main ring buffers
65    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    /// Write XS event
82    pub fn write_xs(&mut self, event: PooledEvent<XS_TSHIRT_SIZE>) {
83        self.xs_writer.add(event);
84    }
85
86    /// Write S event
87    pub fn write_s(&mut self, event: PooledEvent<S_TSHIRT_SIZE>) {
88        self.s_writer.add(event);
89    }
90
91    /// Write M event
92    pub fn write_m(&mut self, event: PooledEvent<M_TSHIRT_SIZE>) {
93        self.m_writer.add(event);
94    }
95
96    /// Write L event
97    pub fn write_l(&mut self, event: PooledEvent<L_TSHIRT_SIZE>) {
98        self.l_writer.add(event);
99    }
100
101    /// Write XL event
102    pub fn write_xl(&mut self, event: PooledEvent<XL_TSHIRT_SIZE>) {
103        self.xl_writer.add(event);
104    }
105
106    /// Helper to write data to optimal ring buffer
107    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    /// Helper to create PooledEvent
133    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    /// Get writer for P2P XS ring
153    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    /// Get writer for P2P S ring
159    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    /// Get writer for P2P M ring
165    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    /// Get writer for P2P L ring
171    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    /// Get writer for P2P XL ring
177    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
183// ================================================================================================
184// VECTOR SEARCH TYPES
185// ================================================================================================
186
187pub 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    /// Create a new XaeroFlux instance
219    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    /// Start the AOF actor
234    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    /// Start P2P networking with XaeroID
242    pub fn start_p2p(&mut self, xaero_id: XaeroID) -> Result<(), Box<dyn std::error::Error>> {
243        // Ensure AOF is started first
244        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                // Get static reference to S ring buffer
251                let s_ring: &'static RingBuffer<S_TSHIRT_SIZE, S_CAPACITY> = S_RING.get_or_init(RingBuffer::new);
252
253                // Create a simple AofState for P2P actor
254                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    /// Start vector search with extractors
273    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    /// Write event data to optimal ring buffer
288    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    /// Send text message via P2P (convenience method)
293    pub fn send_text(&mut self, text: &str) -> Result<(), XaeroFluxError> {
294        self.write_event(text.as_bytes(), 1) // event_type 1 for text
295    }
296
297    /// Send file via P2P (convenience method)
298    pub fn send_file_data(&mut self, file_data: &[u8]) -> Result<(), XaeroFluxError> {
299        self.write_event(file_data, 2) // event_type 2 for files
300    }
301
302    // ================================================================================================
303    // VECTOR SEARCH API
304    // ================================================================================================
305
306    /// Search using a single vector
307    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    /// Search using multiple vectors
337    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    /// Get vector search statistics
349    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            // Initialize if needed (lazy init per thread)
372            if bus_ref.is_none() {
373                *bus_ref = Some(EventBus::new());
374            }
375
376            // Write using the thread's local EventBus
377            bus_ref.as_mut().unwrap().write_optimal(data, event_type)
378        })
379    }
380}
381
382// ================================================================================================
383// ERROR TYPES
384// ================================================================================================
385
386#[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// ================================================================================================
408// TESTS
409// ================================================================================================
410
411#[cfg(test)]
412mod tests {
413    use super::*;
414
415    #[test]
416    fn test_event_bus_creation() {
417        let bus = EventBus::new();
418        // Basic creation test - writers should be ready
419        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        // Test that P2P actors can get writers
436        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]; // Larger than XL
481        let result = xf.write_event(&oversized_data, 42);
482        assert!(matches!(result, Err(XaeroFluxError::DataTooLarge(_))));
483
484        println!("✅ Oversized data properly rejected");
485    }
486}