Handle one sync session with full RLS, audit, DLQ wired in.
(
mut ws: tokio_tungstenite::WebSocketStream<tokio::net::TcpStream>,
addr: SocketAddr,
state: &SyncListenerState,
shared: Option<Arc<SharedState>>,
)
| 24 | |
| 25 | /// Handle one sync session with full RLS, audit, DLQ wired in. |
| 26 | pub(super) async fn handle_sync_session( |
| 27 | mut ws: tokio_tungstenite::WebSocketStream<tokio::net::TcpStream>, |
| 28 | addr: SocketAddr, |
| 29 | state: &SyncListenerState, |
| 30 | shared: Option<Arc<SharedState>>, |
| 31 | ) { |
| 32 | use futures::{SinkExt, StreamExt}; |
| 33 | use tokio_tungstenite::tungstenite::Message; |
| 34 | |
| 35 | let session_id = format!( |
| 36 | "sync-{addr}-{}", |
| 37 | state.connections_accepted.load(Ordering::Relaxed) |
| 38 | ); |
| 39 | let mut session = |
| 40 | super::session::SyncSession::with_rate_limit(session_id.clone(), &state.config.rate_limit); |
| 41 | session.device_metadata.remote_addr = addr.to_string(); |
| 42 | |
| 43 | let jwt_validator = |
| 44 | crate::control::security::jwt::JwtValidator::new(state.config.jwt_config.clone()); |
| 45 | |
| 46 | let mut crdt_delivery_rx: Option< |
| 47 | tokio::sync::mpsc::Receiver<crate::event::crdt_sync::types::OutboundDelta>, |
| 48 | > = None; |
| 49 | let mut crdt_control_rx: Option< |
| 50 | tokio::sync::mpsc::Receiver<nodedb_types::sync::wire::SyncFrame>, |
| 51 | > = None; |
| 52 | let mut crdt_registered = false; |
| 53 | |
| 54 | let mut presence_rx: Option<tokio::sync::mpsc::Receiver<std::sync::Arc<Vec<u8>>>> = None; |
| 55 | let mut presence_registered = false; |
| 56 | |
| 57 | let array_inbound: Option<Arc<crate::control::array_sync::OriginArrayInbound>> = |
| 58 | shared.as_ref().map(|s| { |
| 59 | let engine = Arc::new(crate::control::array_sync::OriginApplyEngine::new( |
| 60 | Arc::clone(&s.array_sync_schemas), |
| 61 | Arc::clone(&s.array_sync_op_log), |
| 62 | )); |
| 63 | let fanout = Arc::new(crate::control::array_sync::ArrayFanout::new( |
| 64 | Arc::clone(&s.shape_registry), |
| 65 | Arc::clone(&s.array_delivery), |
| 66 | Arc::clone(&s.array_subscriber_cursors), |
| 67 | Arc::clone(&s.array_snapshot_hlcs), |
| 68 | Arc::clone(&s.array_merger_registry), |
| 69 | 0, |
| 70 | 0, |
| 71 | )); |
| 72 | let inbound = crate::control::array_sync::OriginArrayInbound::new( |
| 73 | engine, |
| 74 | Arc::clone(&s.array_sync_schemas), |
| 75 | Arc::clone(s), |
| 76 | crate::types::TenantId::new(0), |
| 77 | ) |
| 78 | .with_observer(fanout); |
| 79 | Arc::new(inbound) |
| 80 | }); |
| 81 | |
| 82 | let mut array_delivery_rx: Option<tokio::sync::mpsc::Receiver<Vec<u8>>> = None; |
| 83 | let mut array_delivery_registered = false; |
no test coverage detected