Skip to main content

myko_server/
lib.rs

1//! Myko server runtime — WebSocket, durable event backends, peer federation.
2//!
3//! This crate contains the tokio-dependent parts of the Myko server:
4//! - `CellServer` — server lifecycle (durable catch-up init, WS accept loop)
5//! - `postgres` — PostgreSQL producer/consumer (event-table + LISTEN/NOTIFY)
6//! - `ws_handler` — WebSocket connection handling
7//! - `peer_registry` — federation with other servers
8//! - `mcp` — Model Context Protocol server
9//!
10//! Tokio-free server types (CellServerCtx, HandlerRegistry, etc.) live in `myko::server`.
11
12pub 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
21// Re-export all tokio-free server types from myko
22use 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/// Cell-based Myko server configuration.
48#[derive(Clone)]
49pub struct CellServerConfig {
50    /// Address to bind the WebSocket server
51    pub bind_addr: SocketAddr,
52    /// Disable Nagle's algorithm (set `TCP_NODELAY`) on accepted connections.
53    /// Myko's traffic is small, frequent, latency-sensitive messages (e.g.
54    /// ~60Hz pulses); with Nagle on, TCP coalesces successive small writes
55    /// into fewer segments that arrive together, so an even send cadence is
56    /// delivered as bursts. Defaults to `true`.
57    pub tcp_nodelay: bool,
58    /// Optional Postgres configuration for event persistence/distribution
59    pub postgres: Option<PostgresConfig>,
60    /// Server host ID (auto-generated if not provided)
61    pub host_id: Option<Uuid>,
62    /// Optional peer registry configuration for federation
63    pub peer_registry: Option<peer_registry::PeerRegistryConfig>,
64    /// Default persister override
65    pub default_persister: Option<Arc<dyn Persister>>,
66    /// Per-entity persister overrides keyed by entity type name
67    pub persister_overrides: HashMap<String, Arc<dyn Persister>>,
68    /// Optional pre-constructed peer-client map. When provided, it will be
69    /// used as-is (so any `PeerPersister` built against the same `Arc`
70    /// shares the live map). If `None`, the server creates its own.
71    pub peer_clients: Option<Arc<dashmap::DashMap<Arc<str>, Arc<MykoClient>>>>,
72}
73
74/// Builder for creating a CellServer.
75#[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    /// Optional pre-constructed peer-client map — useful when a
85    /// `PeerPersister` must reference the same map the server will use.
86    /// Defaults to a fresh empty map if not provided.
87    peer_clients: Option<Arc<dashmap::DashMap<Arc<str>, Arc<MykoClient>>>>,
88    after_init: Option<AfterInitCallback>,
89    /// Optional MCP `ServerInfo`. Defaults to `ServerInfo::default()` if not
90    /// set; binaries override this to advertise their own name / version /
91    /// instructions on the `/myko/mcp` endpoint.
92    server_info: Option<mcp::dispatch::ServerInfo>,
93}
94
95type AfterInitCallback = Box<dyn FnOnce(&CellServer) + Send>;
96
97impl CellServerBuilder {
98    /// Create a new server builder.
99    pub fn new() -> Self {
100        Self::default()
101    }
102
103    /// Set the WebSocket bind address.
104    pub fn with_bind_addr(mut self, addr: SocketAddr) -> Self {
105        self.bind_addr = Some(addr);
106        self
107    }
108
109    /// Set whether to disable Nagle's algorithm (`TCP_NODELAY`) on accepted
110    /// connections. Defaults to `true` (Nagle off) — recommended for myko's
111    /// small, frequent, latency-sensitive messages so an even send cadence
112    /// isn't delivered as coalesced bursts.
113    pub fn with_tcp_nodelay(mut self, enabled: bool) -> Self {
114        self.tcp_nodelay = Some(enabled);
115        self
116    }
117
118    /// Set the server host ID (auto-generated if not set).
119    pub fn with_host_id(mut self, id: Uuid) -> Self {
120        self.host_id = Some(id);
121        self
122    }
123
124    /// Configure Postgres for event persistence/distribution.
125    pub fn with_postgres(mut self, config: PostgresConfig) -> Self {
126        self.postgres = Some(config);
127        self
128    }
129
130    /// Configure peer registry for federation.
131    pub fn with_peer_registry(mut self, config: peer_registry::PeerRegistryConfig) -> Self {
132        self.peer_registry = Some(config);
133        self
134    }
135
136    /// Set the default persister used for all entity types without explicit overrides.
137    pub fn with_default_persister(mut self, persister: Arc<dyn Persister>) -> Self {
138        self.default_persister = Some(persister);
139        self
140    }
141
142    /// Override persister for a specific entity type (e.g. "Pulse").
143    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    /// Provide a pre-constructed peer-client map. The server's peer
154    /// registry will populate it as peers connect. Pass the same `Arc`
155    /// into `PeerPersister::new(...)` when you register a
156    /// `with_persister_override(..., PeerPersister)` so the persister
157    /// shares the live map.
158    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    /// Register a callback to run after initialization and relation establishment,
167    /// but before the WebSocket accept loop starts. Use this for starting subsystems
168    /// that need entity data (e.g., scene engine).
169    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    /// Set the MCP `ServerInfo` advertised on the `/myko/mcp` `initialize`
175    /// response. Defaults to `ServerInfo::default()` (`myko-mcp` /
176    /// `CARGO_PKG_VERSION` / no instructions).
177    pub fn with_server_info(mut self, info: mcp::dispatch::ServerInfo) -> Self {
178        self.server_info = Some(info);
179        self
180    }
181
182    /// Build the server.
183    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
206/// Cell-based Myko server.
207///
208/// Uses hyphae cells for reactive queries and reports instead of actors.
209pub struct CellServer {
210    /// Central entity store registry
211    pub registry: Arc<StoreRegistry>,
212    /// Handler registry for items, queries, and reports
213    pub handler_registry: Arc<HandlerRegistry>,
214    /// Relationship manager for cascade operations
215    pub relationship_manager: Arc<RelationshipManager>,
216    /// Optional Postgres producer handle
217    pub postgres_producer: Option<PostgresProducerHandle>,
218    /// Full-text search index
219    pub search_index: Arc<SearchIndex>,
220    /// Persister routing (default + per-entity overrides)
221    pub persisters: Arc<PersisterRouter>,
222    /// Server host ID
223    pub host_id: Uuid,
224    /// Server configuration
225    config: CellServerConfig,
226    /// Postgres producer (kept alive)
227    _postgres_producer_owner: Option<CellPostgresProducer>,
228    /// Postgres consumer (kept alive)
229    postgres_consumer: Option<CellPostgresConsumer>,
230    /// Whether the server is ready to accept connections
231    ready: Arc<AtomicBool>,
232    /// Peer registry for federation (initialized after catch-up)
233    peer_registry_instance: RwLock<Option<peer_registry::PeerRegistry>>,
234    /// Live peer clients shared with report context.
235    peer_clients: Arc<dashmap::DashMap<Arc<str>, Arc<MykoClient>>>,
236    /// Callback to run after init (catch-up + relations) but before WS loop
237    after_init: std::sync::Mutex<Option<AfterInitCallback>>,
238    /// MCP `ServerInfo` advertised on the `/myko/mcp` `initialize` response.
239    /// Set via [`CellServerBuilder::with_server_info`]; defaults to
240    /// `ServerInfo::default()`.
241    server_info: Arc<mcp::dispatch::ServerInfo>,
242    /// Sender for local+replicated event fan-out to saga runtime.
243    saga_event_tx: flume::Sender<MEvent>,
244    /// Receiver consumed when saga runtime starts.
245    saga_event_rx: std::sync::Mutex<Option<flume::Receiver<MEvent>>>,
246    /// Saga tasks kept alive for server lifetime.
247    saga_tasks: std::sync::Mutex<Vec<tokio::task::JoinHandle<()>>>,
248    /// Server ownership death-watch guard (kept alive for server lifetime).
249    _server_ownership_guard: std::sync::Mutex<Option<hyphae::SubscriptionGuard>>,
250    /// Hyphae cell inspector server (kept alive for the lifetime of the server)
251    #[cfg(feature = "inspector")]
252    _inspector: hyphae::server::InspectorServer,
253}
254
255impl CellServer {
256    /// Create a new server builder.
257    pub fn builder() -> CellServerBuilder {
258        CellServerBuilder::new()
259    }
260
261    /// Create a new cell-based server.
262    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        // Initialize the client registry for WebSocket client message dispatch
269        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        // If no durable consumer, server is immediately ready
301        let ready = Arc::new(AtomicBool::new(postgres_consumer.is_none()));
302
303        // Initialize full-text search index
304        let search_index = Arc::new(SearchIndex::new());
305
306        // Build persister routing:
307        // - explicit default from config if provided
308        // - otherwise Postgres producer handle when available
309        // - explicit per-entity overrides always win
310        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        // Start the hyphae cell inspector server
322        #[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    /// Start the peer registry for federation.
358    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    /// Check if peer registry is running.
369    pub fn has_peer_registry(&self) -> bool {
370        self.peer_registry_instance.read().unwrap().is_some()
371    }
372
373    /// Get the store registry.
374    pub fn registry(&self) -> Arc<StoreRegistry> {
375        self.registry.clone()
376    }
377
378    /// Get the handler registry.
379    pub fn handler_registry(&self) -> Arc<HandlerRegistry> {
380        self.handler_registry.clone()
381    }
382
383    /// Get the MCP `ServerInfo` advertised on the `/myko/mcp` `initialize`
384    /// response.
385    pub fn server_info(&self) -> Arc<mcp::dispatch::ServerInfo> {
386        self.server_info.clone()
387    }
388
389    /// Get a server context for module use.
390    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        // NOTE(ts): One unbounded flume channel per saga, with dispatch-side filtering
426        // so sagas only receive events matching their entity type and change type.
427        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        // NOTE(ts): Dispatcher fans out events to saga channels, filtering by
514        // entity type and change type so each saga only receives relevant events.
515        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    /// Create a Postgres-backed history store for replay/windback operations.
531    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    /// Initialize Postgres replay/listener and wait for catch-up.
540    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    /// Establish relationship invariants.
555    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    /// Check if the server is ready to accept connections.
562    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    /// Run the server with full initialization.
574    pub async fn run(&self) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
575        use tokio::net::TcpListener;
576
577        // Persisters can veto startup via startup healthchecks.
578        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        // Wait for Postgres catch-up if configured
592        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        // Build search index from store data (after catch-up)
601        log::info!("Building search index...");
602        self.search_index.build_from_registry(&self.registry);
603
604        // Establish relations (cleanup orphans, ensure required entities)
605        log::info!("Establishing relations...");
606        self.establish_relations();
607
608        // Claim orphaned server-owned items and start death watch
609        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        // Bind WebSocket listener first so peer publication only happens once
620        // the gateway is actually available.
621        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        // Start peer registry if configured
629        if self.config.peer_registry.is_some() {
630            self.start_peer_registry(None);
631        }
632
633        // Run after_init hook (e.g., scene engine startup)
634        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        // WS message-throughput summary thread. Emits a single log line every
646        // 250ms with inbound/outbound counts per message kind. Used for
647        // diagnosing server-vs-client pacing during slow loads.
648        crate::ws_timing::start_periodic_logger();
649
650        // Report-cache hit/miss summary thread. Replaces the per-call debug
651        // log spam that was dominating I/O during loads.
652        myko::server::report_cache_stats::start_periodic_logger();
653
654        // Entity-SET summary thread. Replaces the per-`set` "[entity] SET ..."
655        // debug spam (Pulse SETs dominate under pulse-heavy workloads).
656        myko::server::entity_set_stats::start_periodic_logger();
657
658        // Per-search summary thread. One log line per window listing each
659        // search that completed (entity_type, result count, elapsed).
660        myko::search::search_stats::start_periodic_logger();
661
662        log::info!("Server started");
663        self.run_ws_accept_loop(listener).await
664    }
665
666    /// Run just the accept loop (no Postgres / relations / saga startup).
667    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            // Disable Nagle (unless configured off): our writes are small,
689            // frequent, latency-sensitive messages (e.g. ~60Hz pulses). With
690            // Nagle on (tokio's default), TCP coalesces successive small writes
691            // into fewer segments that land together, so an even 60Hz send
692            // arrives at the client as ~15-20Hz bursts of 3-4 — which downstream
693            // latest-wins consumers (e.g. the pulse-unreal transform apply) then
694            // collapse to one update per burst, producing visibly choppy motion.
695            // Ship each write promptly instead.
696            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            // Check if server is ready (durable backend caught up)
703            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}