1pub mod mcp;
13pub mod peer_persister;
14pub mod peer_registry;
15pub mod postgres;
16pub mod router;
17pub mod server_ownership;
18pub mod ws_handler;
19pub mod ws_timing;
20
21use std::{
23 collections::HashMap,
24 net::SocketAddr,
25 sync::{
26 Arc, RwLock,
27 atomic::{AtomicBool, Ordering},
28 },
29 time::Duration,
30};
31
32use futures_util::StreamExt;
33pub use myko::server::*;
34use myko::{
35 client::MykoClient, command::CommandContext, request::RequestContext, saga::SagaRegistration,
36 search::SearchIndex, store::StoreRegistry, wire::MEvent,
37};
38pub use peer_persister::PeerPersister;
39pub use server_ownership::ServerOwnershipManager;
40use uuid::Uuid;
41
42use crate::postgres::{
43 CellPostgresConsumer, CellPostgresProducer, PostgresConfig, PostgresHistoryReplayProvider,
44 PostgresHistoryStore, PostgresProducerHandle,
45};
46
47#[derive(Clone)]
49pub struct CellServerConfig {
50 pub bind_addr: SocketAddr,
52 pub tcp_nodelay: bool,
58 pub postgres: Option<PostgresConfig>,
60 pub host_id: Option<Uuid>,
62 pub peer_registry: Option<peer_registry::PeerRegistryConfig>,
64 pub default_persister: Option<Arc<dyn Persister>>,
66 pub persister_overrides: HashMap<String, Arc<dyn Persister>>,
68 pub peer_clients: Option<Arc<dashmap::DashMap<Arc<str>, Arc<MykoClient>>>>,
72}
73
74#[derive(Default)]
76pub struct CellServerBuilder {
77 bind_addr: Option<SocketAddr>,
78 tcp_nodelay: Option<bool>,
79 host_id: Option<Uuid>,
80 postgres: Option<PostgresConfig>,
81 peer_registry: Option<peer_registry::PeerRegistryConfig>,
82 default_persister: Option<Arc<dyn Persister>>,
83 persister_overrides: HashMap<String, Arc<dyn Persister>>,
84 peer_clients: Option<Arc<dashmap::DashMap<Arc<str>, Arc<MykoClient>>>>,
88 after_init: Option<AfterInitCallback>,
89 server_info: Option<mcp::dispatch::ServerInfo>,
93}
94
95type AfterInitCallback = Box<dyn FnOnce(&CellServer) + Send>;
96
97impl CellServerBuilder {
98 pub fn new() -> Self {
100 Self::default()
101 }
102
103 pub fn with_bind_addr(mut self, addr: SocketAddr) -> Self {
105 self.bind_addr = Some(addr);
106 self
107 }
108
109 pub fn with_tcp_nodelay(mut self, enabled: bool) -> Self {
114 self.tcp_nodelay = Some(enabled);
115 self
116 }
117
118 pub fn with_host_id(mut self, id: Uuid) -> Self {
120 self.host_id = Some(id);
121 self
122 }
123
124 pub fn with_postgres(mut self, config: PostgresConfig) -> Self {
126 self.postgres = Some(config);
127 self
128 }
129
130 pub fn with_peer_registry(mut self, config: peer_registry::PeerRegistryConfig) -> Self {
132 self.peer_registry = Some(config);
133 self
134 }
135
136 pub fn with_default_persister(mut self, persister: Arc<dyn Persister>) -> Self {
138 self.default_persister = Some(persister);
139 self
140 }
141
142 pub fn with_persister_override(
144 mut self,
145 entity_type: impl Into<String>,
146 persister: Arc<dyn Persister>,
147 ) -> Self {
148 self.persister_overrides
149 .insert(entity_type.into(), persister);
150 self
151 }
152
153 pub fn with_peer_clients(
159 mut self,
160 peer_clients: Arc<dashmap::DashMap<Arc<str>, Arc<MykoClient>>>,
161 ) -> Self {
162 self.peer_clients = Some(peer_clients);
163 self
164 }
165
166 pub fn after_init(mut self, f: impl FnOnce(&CellServer) + Send + 'static) -> Self {
170 self.after_init = Some(Box::new(f));
171 self
172 }
173
174 pub fn with_server_info(mut self, info: mcp::dispatch::ServerInfo) -> Self {
178 self.server_info = Some(info);
179 self
180 }
181
182 pub fn build(self) -> CellServer {
184 let bind_addr = self
185 .bind_addr
186 .unwrap_or_else(|| "127.0.0.1:5155".parse().unwrap());
187
188 let server_info = Arc::new(self.server_info.unwrap_or_default());
189
190 let mut server = CellServer::new(CellServerConfig {
191 bind_addr,
192 tcp_nodelay: self.tcp_nodelay.unwrap_or(true),
193 postgres: self.postgres,
194 host_id: self.host_id,
195 peer_registry: self.peer_registry,
196 default_persister: self.default_persister,
197 persister_overrides: self.persister_overrides,
198 peer_clients: self.peer_clients,
199 });
200 server.after_init = std::sync::Mutex::new(self.after_init);
201 server.server_info = server_info;
202 server
203 }
204}
205
206pub struct CellServer {
210 pub registry: Arc<StoreRegistry>,
212 pub handler_registry: Arc<HandlerRegistry>,
214 pub relationship_manager: Arc<RelationshipManager>,
216 pub postgres_producer: Option<PostgresProducerHandle>,
218 pub search_index: Arc<SearchIndex>,
220 pub persisters: Arc<PersisterRouter>,
222 pub host_id: Uuid,
224 config: CellServerConfig,
226 _postgres_producer_owner: Option<CellPostgresProducer>,
228 postgres_consumer: Option<CellPostgresConsumer>,
230 ready: Arc<AtomicBool>,
232 peer_registry_instance: RwLock<Option<peer_registry::PeerRegistry>>,
234 peer_clients: Arc<dashmap::DashMap<Arc<str>, Arc<MykoClient>>>,
236 after_init: std::sync::Mutex<Option<AfterInitCallback>>,
238 server_info: Arc<mcp::dispatch::ServerInfo>,
242 saga_event_tx: flume::Sender<MEvent>,
244 saga_event_rx: std::sync::Mutex<Option<flume::Receiver<MEvent>>>,
246 saga_tasks: std::sync::Mutex<Vec<tokio::task::JoinHandle<()>>>,
248 _server_ownership_guard: std::sync::Mutex<Option<hyphae::SubscriptionGuard>>,
250 #[cfg(feature = "inspector")]
252 _inspector: hyphae::server::InspectorServer,
253}
254
255impl CellServer {
256 pub fn builder() -> CellServerBuilder {
258 CellServerBuilder::new()
259 }
260
261 pub fn new(config: CellServerConfig) -> Self {
263 let host_id = config.host_id.unwrap_or_else(Uuid::new_v4);
264 let registry = Arc::new(StoreRegistry::new());
265 let handler_registry = Arc::new(HandlerRegistry::new());
266 let relationship_manager = Arc::new(RelationshipManager::new());
267
268 init_client_registry();
270
271 let (saga_event_tx, saga_event_rx) = flume::unbounded::<MEvent>();
272 let (postgres_producer_owner, postgres_producer, postgres_consumer) =
273 if let Some(ref postgres_config) = config.postgres {
274 match CellPostgresProducer::new(postgres_config, host_id) {
275 Ok(producer) => {
276 let handle = producer.handle();
277 let consumer = match CellPostgresConsumer::start(
278 postgres_config,
279 host_id,
280 handler_registry.clone(),
281 registry.clone(),
282 ) {
283 Ok(c) => Some(c),
284 Err(e) => {
285 log::error!("Failed to start Postgres consumer: {}", e);
286 None
287 }
288 };
289 (Some(producer), Some(handle), consumer)
290 }
291 Err(e) => {
292 log::error!("Failed to create Postgres producer: {}", e);
293 (None, None, None)
294 }
295 }
296 } else {
297 (None, None, None)
298 };
299
300 let ready = Arc::new(AtomicBool::new(postgres_consumer.is_none()));
302
303 let search_index = Arc::new(SearchIndex::new());
305
306 let mut persister_router = PersisterRouter::default();
311 if let Some(default_persister) = config.default_persister.clone() {
312 persister_router.set_default(Some(default_persister));
313 } else if let Some(handle) = postgres_producer.clone() {
314 persister_router.set_default(Some(Arc::new(handle) as Arc<dyn Persister>));
315 }
316 for (entity_type, persister) in &config.persister_overrides {
317 persister_router.set_override(entity_type.clone(), persister.clone());
318 }
319 let persisters = Arc::new(persister_router);
320
321 #[cfg(feature = "inspector")]
323 let inspector = hyphae::server::start_server("myko");
324 #[cfg(feature = "inspector")]
325 log::info!("Hyphae inspector on port {}", inspector.port());
326
327 let peer_clients = config
328 .peer_clients
329 .clone()
330 .unwrap_or_else(|| Arc::new(dashmap::DashMap::new()));
331
332 Self {
333 registry,
334 handler_registry,
335 relationship_manager,
336 postgres_producer,
337 search_index,
338 persisters,
339 host_id,
340 config,
341 _postgres_producer_owner: postgres_producer_owner,
342 postgres_consumer,
343 ready,
344 peer_registry_instance: RwLock::new(None),
345 peer_clients,
346 after_init: std::sync::Mutex::new(None),
347 server_info: Arc::new(mcp::dispatch::ServerInfo::default()),
348 saga_event_tx,
349 saga_event_rx: std::sync::Mutex::new(Some(saga_event_rx)),
350 saga_tasks: std::sync::Mutex::new(Vec::new()),
351 _server_ownership_guard: std::sync::Mutex::new(None),
352 #[cfg(feature = "inspector")]
353 _inspector: inspector,
354 }
355 }
356
357 pub fn start_peer_registry(&self, config: Option<peer_registry::PeerRegistryConfig>) {
359 let peer_config = config.or_else(|| self.config.peer_registry.clone());
360
361 if let Some(peer_config) = peer_config {
362 log::info!("Starting peer registry");
363 let pr = peer_registry::PeerRegistry::new(self.ctx(), peer_config);
364 *self.peer_registry_instance.write().unwrap() = Some(pr);
365 }
366 }
367
368 pub fn has_peer_registry(&self) -> bool {
370 self.peer_registry_instance.read().unwrap().is_some()
371 }
372
373 pub fn registry(&self) -> Arc<StoreRegistry> {
375 self.registry.clone()
376 }
377
378 pub fn handler_registry(&self) -> Arc<HandlerRegistry> {
380 self.handler_registry.clone()
381 }
382
383 pub fn server_info(&self) -> Arc<mcp::dispatch::ServerInfo> {
386 self.server_info.clone()
387 }
388
389 pub fn ctx(&self) -> CellServerCtx {
391 let history_replay: Option<Arc<dyn myko::server::HistoryReplayProvider>> =
392 self.config.postgres.as_ref().map(|pg| {
393 Arc::new(PostgresHistoryReplayProvider::new(pg.clone()))
394 as Arc<dyn myko::server::HistoryReplayProvider>
395 });
396 CellServerCtx::new(
397 self.host_id,
398 self.registry.clone(),
399 self.handler_registry.clone(),
400 self.relationship_manager.clone(),
401 self.persisters.clone(),
402 self.search_index.clone(),
403 self.peer_clients.clone(),
404 Some(self.saga_event_tx.clone()),
405 history_replay,
406 )
407 }
408
409 fn start_saga_runtime(&self) {
410 let registrations: Vec<_> = inventory::iter::<SagaRegistration>().collect();
411 if registrations.is_empty() {
412 return;
413 }
414 let Some(rx) = self
415 .saga_event_rx
416 .lock()
417 .expect("saga_event_rx mutex poisoned")
418 .take()
419 else {
420 return;
421 };
422
423 log::info!("Starting saga runtime with {} saga(s)", registrations.len());
424
425 struct SagaChannel {
428 tx: flume::Sender<MEvent>,
429 entity_type: &'static str,
430 change_type: myko::event::MEventType,
431 }
432 let mut saga_channels: Vec<SagaChannel> = Vec::new();
433
434 for registration in registrations {
435 let saga = (registration.create)();
436 let saga_name = saga.name().to_string();
437 let (saga_tx, saga_rx) = flume::unbounded::<MEvent>();
438 saga_channels.push(SagaChannel {
439 tx: saga_tx,
440 entity_type: registration.event_entity_type,
441 change_type: registration.event_change_type,
442 });
443 let events: myko::saga::EventStream = Box::pin(futures_util::stream::unfold(
444 saga_rx,
445 move |saga_rx| async move {
446 saga_rx
447 .recv_async()
448 .await
449 .ok()
450 .map(|event| (event, saga_rx))
451 },
452 ));
453
454 let saga_ctx = Arc::new(myko::saga::SagaContext::with_event_sink(
455 self.host_id,
456 self.registry.clone(),
457 self.saga_event_tx.clone(),
458 ));
459 let mut command_stream = saga.build_boxed(events, saga_ctx);
460
461 let host_id = self.host_id;
462 let registry = self.registry.clone();
463 let handler_registry = self.handler_registry.clone();
464 let relationship_manager = self.relationship_manager.clone();
465 let persisters = self.persisters.clone();
466 let search_index = self.search_index.clone();
467 let peer_clients = self.peer_clients.clone();
468 let saga_event_tx = self.saga_event_tx.clone();
469
470 let handle = tokio::spawn(async move {
471 while let Some(command) = command_stream.next().await {
472 let command_name = command.command_name();
473 log::debug!("Saga {} executing command {}", saga_name, command_name);
474 let req = Arc::new(RequestContext::internal(
475 Arc::from(Uuid::new_v4().to_string()),
476 host_id,
477 &format!("saga:{saga_name}"),
478 ));
479
480 let cmd_ctx = CommandContext::new(
481 Arc::from(command_name),
482 req,
483 Arc::new(CellServerCtx::new(
484 host_id,
485 registry.clone(),
486 handler_registry.clone(),
487 relationship_manager.clone(),
488 persisters.clone(),
489 search_index.clone(),
490 peer_clients.clone(),
491 Some(saga_event_tx.clone()),
492 None,
493 )),
494 );
495
496 if let Err(err) = command.execute_boxed(cmd_ctx) {
497 log::error!(
498 "Saga {} command {} failed: {}",
499 saga_name,
500 command_name,
501 err.message
502 );
503 }
504 }
505 });
506
507 self.saga_tasks
508 .lock()
509 .expect("saga_tasks mutex poisoned")
510 .push(handle);
511 }
512
513 let dispatcher = tokio::spawn(async move {
516 while let Ok(event) = rx.recv_async().await {
517 for ch in &saga_channels {
518 if event.item_type == ch.entity_type && event.change_type == ch.change_type {
519 let _ = ch.tx.send(event.clone());
520 }
521 }
522 }
523 });
524 self.saga_tasks
525 .lock()
526 .expect("saga_tasks mutex poisoned")
527 .push(dispatcher);
528 }
529
530 pub fn postgres_history_store(&self) -> Result<Option<PostgresHistoryStore>, String> {
532 self.config
533 .postgres
534 .clone()
535 .map(PostgresHistoryStore::new)
536 .transpose()
537 }
538
539 pub fn init_postgres_and_wait(&self, timeout: Duration) -> Result<(), String> {
541 if self.config.postgres.is_some() && self.postgres_consumer.is_none() {
542 return Err(
543 "Postgres is configured but the Postgres consumer is not running".to_string(),
544 );
545 }
546
547 if let Some(ref consumer) = self.postgres_consumer {
548 consumer.wait_until_caught_up(timeout)?;
549 self.ready.store(true, Ordering::SeqCst);
550 }
551 Ok(())
552 }
553
554 pub fn establish_relations(&self) {
556 if let Err(e) = self.relationship_manager.establish_relations(&self.ctx()) {
557 log::error!("Failed to establish relations: {e}");
558 }
559 }
560
561 pub fn is_ready(&self) -> bool {
563 if let Some(ref consumer) = self.postgres_consumer {
564 if consumer.is_caught_up() {
565 self.ready.store(true, Ordering::SeqCst);
566 return true;
567 }
568 return false;
569 }
570 true
571 }
572
573 pub async fn run(&self) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
575 use tokio::net::TcpListener;
576
577 let entity_types: Vec<&str> = self
579 .handler_registry
580 .entity_types()
581 .map(|t| t.as_ref())
582 .collect();
583 self.persisters
584 .startup_healthcheck(&entity_types)
585 .map_err(|reason| format!("Persister startup healthcheck failed: {reason}"))?;
586
587 if self.config.postgres.is_some() && self.postgres_consumer.is_none() {
588 return Err("Postgres is configured but the Postgres consumer failed to start".into());
589 }
590
591 if self.postgres_consumer.is_some() {
593 log::info!("Waiting for Postgres event consumer to catch up...");
594 let timeout = std::time::Duration::from_secs(300);
595 self.init_postgres_and_wait(timeout)
596 .map_err(|reason| format!("Postgres startup catch-up failed: {reason}"))?;
597 log::info!("Postgres caught up, ready to accept connections");
598 }
599
600 log::info!("Building search index...");
602 self.search_index.build_from_registry(&self.registry);
603
604 log::info!("Establishing relations...");
606 self.establish_relations();
607
608 log::info!("Checking server-owned item ownership...");
610 if let Err(e) = ServerOwnershipManager::claim_orphaned(&self.ctx()) {
611 log::error!("Failed to claim orphaned server-owned items: {}", e);
612 }
613 let ownership_guard = ServerOwnershipManager::watch_peer_deaths(&self.ctx());
614 *self
615 ._server_ownership_guard
616 .lock()
617 .expect("server_ownership_guard mutex poisoned") = Some(ownership_guard);
618
619 let listener = TcpListener::bind(&self.config.bind_addr).await?;
622 log::info!("CellServer listening on {}", self.config.bind_addr);
623 log::info!(
624 "Myko gateway: ws://{}/myko | MCP: /myko/mcp (POST + WS + SSE)",
625 self.config.bind_addr
626 );
627
628 if self.config.peer_registry.is_some() {
630 self.start_peer_registry(None);
631 }
632
633 if let Some(hook) = self
635 .after_init
636 .lock()
637 .expect("after_init mutex poisoned")
638 .take()
639 {
640 hook(self);
641 }
642
643 self.start_saga_runtime();
644
645 crate::ws_timing::start_periodic_logger();
649
650 myko::server::report_cache_stats::start_periodic_logger();
653
654 myko::server::entity_set_stats::start_periodic_logger();
657
658 myko::search::search_stats::start_periodic_logger();
661
662 log::info!("Server started");
663 self.run_ws_accept_loop(listener).await
664 }
665
666 pub async fn run_ws_loop(&self) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
668 use tokio::net::TcpListener;
669
670 let listener = TcpListener::bind(&self.config.bind_addr).await?;
671 log::info!("CellServer listening on {}", self.config.bind_addr);
672 log::info!(
673 "Myko gateway: ws://{}/myko | MCP: /myko/mcp (POST + WS + SSE)",
674 self.config.bind_addr
675 );
676 self.run_ws_accept_loop(listener).await
677 }
678
679 async fn run_ws_accept_loop(
680 &self,
681 listener: tokio::net::TcpListener,
682 ) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
683 let ready = self.ready.clone();
684
685 loop {
686 let (stream, addr) = listener.accept().await?;
687
688 if self.config.tcp_nodelay
697 && let Err(e) = stream.set_nodelay(true)
698 {
699 log::warn!("failed to set TCP_NODELAY on connection from {addr}: {e}");
700 }
701
702 if !ready.load(Ordering::SeqCst) {
704 if self.is_ready() {
705 log::info!("Server is now ready to accept connections");
706 } else {
707 log::warn!(
708 "Rejecting connection from {} - server not ready (durable backend catching up)",
709 addr
710 );
711 drop(stream);
712 continue;
713 }
714 }
715
716 log::debug!("New connection from {}", addr);
717
718 let ctx = Arc::new(self.ctx());
719 let server_info = self.server_info.clone();
720
721 tokio::spawn(async move {
722 if let Err(e) = router::route_connection(stream, addr, ctx, server_info).await {
723 log::error!("Connection error from {}: {}", addr, e);
724 }
725 });
726 }
727 }
728}
729
730#[cfg(test)]
731mod tests {
732 use super::*;
733
734 #[test]
735 fn test_server_creation() {
736 let config = CellServerConfig {
737 bind_addr: "127.0.0.1:0".parse().unwrap(),
738 tcp_nodelay: true,
739 postgres: None,
740 host_id: None,
741 peer_registry: None,
742 default_persister: None,
743 persister_overrides: HashMap::new(),
744 peer_clients: None,
745 };
746 let server = CellServer::new(config);
747 assert!(Arc::strong_count(&server.registry) >= 1);
748 }
749
750 #[test]
751 fn test_server_with_host_id() {
752 let host_id = Uuid::new_v4();
753 let config = CellServerConfig {
754 bind_addr: "127.0.0.1:0".parse().unwrap(),
755 tcp_nodelay: true,
756 postgres: None,
757 host_id: Some(host_id),
758 peer_registry: None,
759 default_persister: None,
760 persister_overrides: HashMap::new(),
761 peer_clients: None,
762 };
763 let server = CellServer::new(config);
764 assert_eq!(server.host_id, host_id);
765 }
766}