MCPcopy Create free account
hub / github.com/NodeDB-Lab/nodedb / handle_sync_session

Function handle_sync_session

nodedb/src/control/server/sync/session_handler.rs:26–576  ·  view source on GitHub ↗

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>>,
)

Source from the content-addressed store, hash-verified

24
25/// Handle one sync session with full RLS, audit, DLQ wired in.
26pub(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;

Callers 1

accept_loopFunction · 0.85

Calls 15

dispatch_array_frameFunction · 0.85
to_stringMethod · 0.80
with_observerMethod · 0.80
try_recvMethod · 0.80
subscribe_to_channelMethod · 0.80
handle_updateMethod · 0.80
sendersMethod · 0.80
send_allMethod · 0.80

Tested by

no test coverage detected