Skip to content

Commit afc7671

Browse files
maxholmanclaude
andcommitted
refactor(transport): deduplicate QUIC/WS server task setup
Extract spawn_server_tasks() in server.rs — resolves shared resources from ServerOptions, spawns the control loop with a Handler, and builds the AcceptResult. Replaces ~35 lines of identical setup in both QUIC and WS accept paths. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
1 parent 624bc9c commit afc7671

3 files changed

Lines changed: 91 additions & 120 deletions

File tree

crates/core/src/server/quic/mod.rs

Lines changed: 6 additions & 60 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,6 @@ use wallhack_wire::{
99

1010
use crate::{
1111
NodeRole, SocketAddrExt as _,
12-
control::{handler::Handler, metrics::Metrics, peers::Registry, routes::RouteTable},
1312
psk::HandshakeExt,
1413
server::tls::{ALPN_QUIC_HTTP, configure_crypto},
1514
transport::{
@@ -185,66 +184,13 @@ impl Server for QuicServer {
185184
}
186185
}
187186

188-
// Get or create shared metrics
189-
let metrics = self
190-
.options
191-
.metrics
192-
.clone()
193-
.unwrap_or_else(|| Arc::new(Metrics::default()));
194-
195-
let channels = super::server::DataChannels::new();
196-
197-
// Create control channel for injecting outgoing control messages
198-
let (control_tx, control_rx) = tokio::sync::mpsc::channel::<ControlMessage>(64);
199-
200-
// Spawn control stream task with handler
201-
let handler_config = self.options.handler_config.clone();
202-
let peers = self
203-
.options
204-
.peers
205-
.clone()
206-
.unwrap_or_else(|| Arc::new(Registry::new()));
207-
let routes = self
208-
.options
209-
.routes
210-
.clone()
211-
.unwrap_or_else(RouteTable::shared);
212-
213-
let peer_name = peer_handshake.as_ref().map(|hs| hs.name.clone());
214-
let route_updates = self.options.route_updates.clone().unwrap_or_else(|| {
215-
let (tx, _) = tokio::sync::broadcast::channel(16);
216-
tx
217-
});
218-
219-
{
220-
let metrics = Arc::clone(&metrics);
221-
let peer_registry = Arc::clone(&peers);
222-
tokio::spawn(async move {
223-
let handler = Handler::new(handler_config, metrics, peers, routes, route_updates);
224-
let mut channels = protocol::ControlChannels {
225-
outgoing_rx: control_rx,
226-
handshake_tx: None, // Handshake already read above
227-
control_response_tx: None, // server doesn't issue ControlRequests
228-
peer_registry: Some(peer_registry),
229-
peer_name,
230-
};
231-
let mut control_stream =
232-
wallhack_transport::erased::BoxBiStream::new(control_stream);
233-
let exit = channels
234-
.run(&mut control_stream, Some(&handler), Duration::from_secs(30))
235-
.await;
236-
tracing::debug!("Control stream finished: {exit:?}");
237-
});
238-
}
239-
240-
// Data tasks are NOT spawned here — the caller does that after PSK validation.
241-
Ok(Some(AcceptResult::with_handshake(
242-
Arc::clone(&transport),
243-
channels,
244-
remote_addr,
245-
metrics,
187+
let control_stream = wallhack_transport::erased::BoxBiStream::new(control_stream);
188+
Ok(Some(super::server::spawn_server_tasks(
189+
transport,
190+
control_stream,
191+
&self.options,
246192
peer_handshake,
247-
control_tx,
193+
remote_addr,
248194
channel_binding,
249195
)))
250196
}

crates/core/src/server/server.rs

Lines changed: 79 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -194,6 +194,85 @@ pub struct ServerOptions {
194194
pub local_handshake: Option<Handshake>,
195195
}
196196

197+
/// Spawn the control task and build an `AcceptResult` — shared by all
198+
/// server transports.
199+
///
200+
/// Called after the transport handshake exchange. Resolves shared
201+
/// resources from `options`, spawns the control loop with a `Handler`,
202+
/// and returns a fully wired `AcceptResult`.
203+
pub fn spawn_server_tasks<T: Transport + 'static>(
204+
transport: Arc<T>,
205+
control_stream: wallhack_transport::erased::BoxBiStream,
206+
options: &ServerOptions,
207+
peer_handshake: Option<Handshake>,
208+
remote_addr: String,
209+
channel_binding: Option<[u8; crate::psk::CHANNEL_BINDING_LEN]>,
210+
) -> AcceptResult<T>
211+
where
212+
T::SendStream: 'static,
213+
T::RecvStream: 'static,
214+
T::BiStream: Send + 'static,
215+
{
216+
use crate::{
217+
control::{handler::Handler, metrics::Metrics, peers::Registry, routes::RouteTable},
218+
transport::protocol,
219+
};
220+
221+
let metrics = options
222+
.metrics
223+
.clone()
224+
.unwrap_or_else(|| Arc::new(Metrics::default()));
225+
226+
let channels = DataChannels::new();
227+
let (control_tx, control_rx) = mpsc::channel::<ControlMessage>(64);
228+
229+
let handler_config = options.handler_config.clone();
230+
let peers = options
231+
.peers
232+
.clone()
233+
.unwrap_or_else(|| Arc::new(Registry::new()));
234+
let routes = options.routes.clone().unwrap_or_else(RouteTable::shared);
235+
let peer_name = peer_handshake.as_ref().map(|hs| hs.name.clone());
236+
let route_updates = options.route_updates.clone().unwrap_or_else(|| {
237+
let (tx, _) = tokio::sync::broadcast::channel(16);
238+
tx
239+
});
240+
241+
{
242+
let metrics = Arc::clone(&metrics);
243+
let peer_registry = Arc::clone(&peers);
244+
tokio::spawn(async move {
245+
let handler = Handler::new(handler_config, metrics, peers, routes, route_updates);
246+
let mut channels = protocol::ControlChannels {
247+
outgoing_rx: control_rx,
248+
handshake_tx: None,
249+
control_response_tx: None,
250+
peer_registry: Some(peer_registry),
251+
peer_name,
252+
};
253+
let mut control_stream = control_stream;
254+
let exit = channels
255+
.run(
256+
&mut control_stream,
257+
Some(&handler),
258+
std::time::Duration::from_secs(30),
259+
)
260+
.await;
261+
tracing::debug!("Control stream finished: {exit:?}");
262+
});
263+
}
264+
265+
AcceptResult::with_handshake(
266+
transport,
267+
channels,
268+
remote_addr,
269+
metrics,
270+
peer_handshake,
271+
control_tx,
272+
channel_binding,
273+
)
274+
}
275+
197276
pub trait Server {
198277
type Error: std::error::Error + std::fmt::Debug + Send + Sync + 'static;
199278
type Transport: Transport;

crates/core/src/server/ws/mod.rs

Lines changed: 6 additions & 60 deletions
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,6 @@ use yamux::Mode;
2525

2626
use crate::{
2727
NodeRole, SocketAddrExt as _,
28-
control::{handler::Handler, metrics::Metrics, peers::Registry, routes::RouteTable},
2928
psk::HandshakeExt,
3029
transport::{
3130
protocol,
@@ -266,66 +265,13 @@ impl Server for WebSocketServer {
266265
}
267266
}
268267

269-
// Get or create shared metrics
270-
let metrics = self
271-
.options
272-
.metrics
273-
.clone()
274-
.unwrap_or_else(|| Arc::new(Metrics::default()));
275-
276-
let channels = super::server::DataChannels::new();
277-
278-
// Create control channel for injecting outgoing control messages
279-
let (control_tx, control_rx) = tokio::sync::mpsc::channel::<ControlMessage>(64);
280-
281-
// Spawn control stream task with handler
282-
let handler_config = self.options.handler_config.clone();
283-
let peers = self
284-
.options
285-
.peers
286-
.clone()
287-
.unwrap_or_else(|| Arc::new(Registry::new()));
288-
let routes = self
289-
.options
290-
.routes
291-
.clone()
292-
.unwrap_or_else(RouteTable::shared);
293-
294-
let peer_name = peer_handshake.as_ref().map(|hs| hs.name.clone());
295-
let route_updates = self.options.route_updates.clone().unwrap_or_else(|| {
296-
let (tx, _) = tokio::sync::broadcast::channel(16);
297-
tx
298-
});
299-
300-
{
301-
let metrics = Arc::clone(&metrics);
302-
let peer_registry = Arc::clone(&peers);
303-
tokio::spawn(async move {
304-
let handler = Handler::new(handler_config, metrics, peers, routes, route_updates);
305-
let mut channels = protocol::ControlChannels {
306-
outgoing_rx: control_rx,
307-
handshake_tx: None, // Handshake already read above
308-
control_response_tx: None, // server doesn't issue ControlRequests
309-
peer_registry: Some(peer_registry),
310-
peer_name,
311-
};
312-
let mut control_stream =
313-
wallhack_transport::erased::BoxBiStream::new(control_stream);
314-
let exit = channels
315-
.run(&mut control_stream, Some(&handler), Duration::from_secs(30))
316-
.await;
317-
tracing::debug!("Control stream finished: {exit:?}");
318-
});
319-
}
320-
321-
// Data tasks are NOT spawned here — the caller does that after PSK validation.
322-
Ok(Some(AcceptResult::with_handshake(
323-
Arc::clone(&transport),
324-
channels,
325-
peer_addr.to_string(),
326-
metrics,
268+
let control_stream = wallhack_transport::erased::BoxBiStream::new(control_stream);
269+
Ok(Some(super::server::spawn_server_tasks(
270+
transport,
271+
control_stream,
272+
&self.options,
327273
peer_handshake,
328-
control_tx,
274+
peer_addr.to_string(),
329275
channel_binding,
330276
)))
331277
}

0 commit comments

Comments
 (0)