1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
use crate::metrics::ReverseMetrics;
use crate::{
redact_auth, relay_bidirectional_with_timeout, server_auth_handshake, ControlState,
ProtocolError,
};
use std::net::{IpAddr, SocketAddr};
use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::Arc;
use std::time::Instant;
use tokio::net::{TcpListener, TcpStream};
use tokio::sync::mpsc;
use tokio::task::JoinSet;
use tokio_util::sync::CancellationToken;
use tracing::{debug, error, info, warn};
/// Configuration for a reverse proxy server (acceptor side).
///
/// The server accepts control connections from remote clients and dispatches
/// externally-accepted connections back through the control channel.
#[derive(Debug, Clone)]
pub struct ReverseServerConfig {
/// Address to bind the control listener on.
pub control_bind: SocketAddr,
/// Address to bind the external listener on (for clients to connect to).
pub external_bind: Option<SocketAddr>,
/// Optional username for authentication.
pub auth_username: Option<String>,
/// Optional password for authentication.
pub auth_password: Option<String>,
/// Maximum concurrent control connections.
pub max_control_connections: u32,
/// Read timeout in milliseconds (for idle control connections).
pub read_timeout_ms: u64,
/// Optional list of allowed external bind addresses. When `Some` and
/// non-empty, the server rejects bind addresses not in the list.
/// When `None` or empty, no allowlist enforcement is applied.
pub allow_bind: Option<Vec<SocketAddr>>,
/// Maximum number of external listeners per control client. Currently
/// pproxy supports one external listener per control connection, so
/// defaults to 1.
pub max_listeners_per_client: u32,
/// Maximum concurrent streams per external listener.
pub max_streams_per_listener: u32,
/// Maximum number of concurrent external clients queued while waiting
/// for a control connection. Excess clients are dropped.
pub max_pending_external: u32,
}
impl Default for ReverseServerConfig {
fn default() -> Self {
Self {
control_bind: "127.0.0.1:0".parse().unwrap(),
external_bind: None,
auth_username: None,
auth_password: None,
max_control_connections: 256,
read_timeout_ms: 300_000,
allow_bind: None,
max_listeners_per_client: 1,
max_streams_per_listener: 1024,
max_pending_external: 1024,
}
}
}
impl ReverseServerConfig {
/// Returns true if the supplied external bind address is allowed by the
/// configured `allow_bind` policy. When `allow_bind` is `None` or empty,
/// all addresses are allowed.
pub fn is_bind_allowed(&self, addr: SocketAddr) -> bool {
match &self.allow_bind {
None => true,
Some(list) if list.is_empty() => true,
Some(list) => list.iter().any(|allowed| same_bind(allowed, &addr)),
}
}
/// Returns true if the address is loopback (127.0.0.0/8 or ::1).
pub fn is_loopback(addr: SocketAddr) -> bool {
match addr.ip() {
IpAddr::V4(v4) => v4.is_loopback(),
IpAddr::V6(v6) => v6.is_loopback(),
}
}
/// Validate this configuration. Returns an error if the configuration is
/// unsafe (e.g. external bind on a non-loopback address without
/// authentication and without an explicit `allow_bind` allowlist).
///
/// This is a defense-in-depth check: it catches misconfigurations that
/// would otherwise expose the reverse proxy to unauthenticated network
/// clients.
pub fn validate(&self) -> Result<(), ProtocolError> {
if let Some(external) = self.external_bind {
// Non-loopback external bind requires BOTH authentication
// credentials AND a non-empty `allow_bind` allowlist. This
// prevents accidentally exposing the reverse proxy to the
// local network without operator intent.
if !Self::is_loopback(external) {
let has_auth = self.auth_username.as_deref().is_some_and(|s| !s.is_empty())
&& self.auth_password.as_deref().is_some_and(|s| !s.is_empty());
let has_allowlist = matches!(&self.allow_bind, Some(list) if !list.is_empty());
if !has_auth {
return Err(ProtocolError::ConfigInvalid(format!(
"reverse server external_bind={external} is non-loopback but no \
authentication is configured; set auth_username/auth_password or \
bind to loopback"
)));
}
if !has_allowlist {
return Err(ProtocolError::ConfigInvalid(format!(
"reverse server external_bind={external} is non-loopback but \
allow_bind is empty; configure an explicit allowlist"
)));
}
}
}
Ok(())
}
}
fn same_bind(a: &SocketAddr, b: &SocketAddr) -> bool {
a.port() == b.port()
&& match (a.ip(), b.ip()) {
(IpAddr::V4(a4), IpAddr::V4(b4)) => a4 == b4,
(IpAddr::V6(a6), IpAddr::V6(b6)) => a6 == b6,
_ => false,
}
}
/// Active state of the reverse server, exposed for tests and admin hooks.
#[derive(Debug, Default)]
pub struct ReverseServerState {
/// Number of currently-active (accepted, awaiting use) control connections.
pub active_control: AtomicU32,
/// Number of currently-active external streams being relayed.
pub active_streams: AtomicU32,
/// Number of external clients waiting for a control connection.
pub pending_external: AtomicU32,
/// Number of listeners denied because of allow_bind.
pub denied_bind: AtomicU32,
/// Number of streams dropped because max_streams_per_listener was reached.
pub dropped_stream_limit: AtomicU32,
/// Number of external clients dropped because max_pending_external was reached.
pub dropped_pending_limit: AtomicU32,
}
impl ReverseServerState {
/// Snapshot of the counters for admin/log display.
pub fn snapshot(&self) -> ReverseServerStateSnapshot {
ReverseServerStateSnapshot {
active_control: self.active_control.load(Ordering::Relaxed),
active_streams: self.active_streams.load(Ordering::Relaxed),
pending_external: self.pending_external.load(Ordering::Relaxed),
denied_bind: self.denied_bind.load(Ordering::Relaxed),
dropped_stream_limit: self.dropped_stream_limit.load(Ordering::Relaxed),
dropped_pending_limit: self.dropped_pending_limit.load(Ordering::Relaxed),
}
}
}
/// Plain-data snapshot of [`ReverseServerState`].
#[derive(Debug, Clone, serde::Serialize)]
pub struct ReverseServerStateSnapshot {
pub active_control: u32,
pub active_streams: u32,
pub pending_external: u32,
pub denied_bind: u32,
pub dropped_stream_limit: u32,
pub dropped_pending_limit: u32,
}
/// The reverse proxy server (acceptor side).
///
/// Accepts control connections from reverse clients and external clients,
/// relaying traffic between them. Each control connection carries exactly
/// one proxy session (matching pproxy's backward model).
pub struct ReverseServer {
config: ReverseServerConfig,
cancel: CancellationToken,
metrics: Option<Arc<ReverseMetrics>>,
state: Arc<ReverseServerState>,
}
impl ReverseServer {
pub fn new(config: ReverseServerConfig) -> Self {
Self {
config,
cancel: CancellationToken::new(),
metrics: None,
state: Arc::new(ReverseServerState::default()),
}
}
/// Attach metrics to this server instance.
pub fn set_metrics(&mut self, metrics: Arc<ReverseMetrics>) {
self.metrics = Some(metrics);
}
/// Get a handle to the active server state.
pub fn state_handle(&self) -> Arc<ReverseServerState> {
self.state.clone()
}
/// Get a cancel token for external shutdown.
pub fn cancel_token(&self) -> CancellationToken {
self.cancel.clone()
}
/// Validate the configured bind address against `allow_bind` before
/// binding. Returns the resolved listener or an error.
async fn bind_external_listener(
config: &ReverseServerConfig,
state: &ReverseServerState,
) -> Result<Option<TcpListener>, ProtocolError> {
let external_bind = match config.external_bind {
Some(addr) => addr,
None => return Ok(None),
};
if !config.is_bind_allowed(external_bind) {
state.denied_bind.fetch_add(1, Ordering::Relaxed);
return Err(ProtocolError::BindDenied(external_bind));
}
let listener = TcpListener::bind(external_bind).await?;
let addr = listener.local_addr()?;
info!(addr = %addr, "reverse server listening for external clients");
Ok(Some(listener))
}
/// Start the reverse server.
pub async fn run(self) -> Result<(), ProtocolError> {
// Defense-in-depth validation: catch unsafe configurations (e.g.
// non-loopback external_bind without auth or allow_bind allowlist)
// before binding any sockets.
self.config.validate()?;
if self.config.auth_username.is_none() || self.config.auth_password.is_none() {
warn!(
control_bind = %self.config.control_bind,
"reverse server control channel has no authentication configured"
);
}
// Enforce the allow_bind policy up-front so misconfiguration is loud.
if let Some(external_bind) = self.config.external_bind {
if !self.config.is_bind_allowed(external_bind) {
self.state.denied_bind.fetch_add(1, Ordering::Relaxed);
return Err(ProtocolError::BindDenied(external_bind));
}
}
let control_listener = TcpListener::bind(&self.config.control_bind).await?;
let control_addr = control_listener.local_addr()?;
info!(addr = %control_addr, "reverse server listening for control connections");
let external_listener = Self::bind_external_listener(&self.config, &self.state).await?;
let config = Arc::new(self.config);
let cancel = self.cancel.clone();
let state = self.state.clone();
let metrics = self.metrics.clone();
// Channel for available control connections
let (control_tx, control_rx) = mpsc::unbounded_channel::<ControlStream>();
// Spawn control connection acceptor
let config_clone = config.clone();
let cancel_clone = cancel.clone();
let control_tx_clone = control_tx.clone();
let metrics_clone = metrics.clone();
let state_clone = state.clone();
let control_task = tokio::spawn(async move {
Self::accept_control_connections(
control_listener,
config_clone,
cancel_clone,
control_tx_clone,
metrics_clone,
state_clone,
)
.await;
});
// Spawn external client acceptor
let external_task = if let Some(external_listener) = external_listener {
let config_clone = config.clone();
let cancel_clone = cancel.clone();
let metrics_clone = metrics.clone();
let state_clone = state.clone();
Some(tokio::spawn(async move {
Self::accept_external_clients(
external_listener,
config_clone,
cancel_clone,
control_rx,
metrics_clone,
state_clone,
)
.await;
}))
} else {
// No external listener: drain the control channel so the
// counter accurately reflects connections that have not yet
// been paired with an external client. Each received stream
// is closed and the active_control counter is decremented.
let state_clone = state.clone();
let metrics_clone = metrics.clone();
let cancel_clone = cancel.clone();
Some(tokio::spawn(async move {
let mut control_rx = control_rx;
loop {
tokio::select! {
Some(ctrl) = control_rx.recv() => {
debug!(
control_peer = %ctrl.peer_addr,
"dropping control connection: no external listener"
);
drop(ctrl.stream);
state_clone.active_control.fetch_sub(1, Ordering::Relaxed);
if let Some(m) = metrics_clone.as_deref() {
m.record_control_closed();
}
}
_ = cancel_clone.cancelled() => break,
}
}
}))
};
// Wait for shutdown
cancel.cancelled().await;
let drain_start = Instant::now();
info!("reverse server shutting down, draining active streams");
// Stop the accept loops before returning. The external accept loop
// also aborts and joins every in-flight relay task it owns.
let _ = control_task.await;
if let Some(task) = external_task {
let _ = task.await;
}
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
let drain_ms = drain_start.elapsed().as_millis() as u64;
if let Some(ref m) = metrics {
m.record_drain(drain_ms);
}
info!(drain_ms, "reverse server drain complete");
Ok(())
}
/// Accept control connections, authenticate, and add to available pool.
async fn accept_control_connections(
listener: TcpListener,
config: Arc<ReverseServerConfig>,
cancel: CancellationToken,
control_tx: mpsc::UnboundedSender<ControlStream>,
metrics: Option<Arc<ReverseMetrics>>,
state: Arc<ReverseServerState>,
) {
loop {
tokio::select! {
result = listener.accept() => {
match result {
Ok((stream, peer_addr)) => {
// Enforce the per-server control connection cap.
// Atomically increment then check to avoid TOCTOU race.
let prev = state.active_control.fetch_add(1, Ordering::AcqRel);
if prev >= config.max_control_connections {
state.active_control.fetch_sub(1, Ordering::Relaxed);
warn!(
peer = %peer_addr,
max = config.max_control_connections,
"rejecting control connection: max reached"
);
if let Some(ref m) = metrics {
m.record_control_rejected(peer_addr, "max_control_connections");
}
drop(stream);
continue;
}
let config = config.clone();
let control_tx = control_tx.clone();
let metrics = metrics.clone();
let state = state.clone();
tokio::spawn(async move {
if let Err(e) = Self::handle_control_connection(
stream,
peer_addr,
config,
control_tx,
metrics.as_deref(),
state.clone(),
).await {
state.active_control.fetch_sub(1, Ordering::Relaxed);
debug!(peer = %peer_addr, error = %e, "control connection handler error");
}
});
}
Err(e) => {
error!(error = %e, "failed to accept control connection");
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
}
}
}
_ = cancel.cancelled() => {
break;
}
}
}
}
/// Handle a single control connection: authenticate and add to pool.
async fn handle_control_connection(
mut stream: TcpStream,
peer_addr: SocketAddr,
config: Arc<ReverseServerConfig>,
control_tx: mpsc::UnboundedSender<ControlStream>,
metrics: Option<&ReverseMetrics>,
state: Arc<ReverseServerState>,
) -> Result<(), ProtocolError> {
info!(peer = %peer_addr, state = ?ControlState::Connecting, "new control connection");
// Authenticate if configured
let redacted = if config.auth_username.is_some() && config.auth_password.is_some() {
let authenticating_start = Instant::now();
let result = server_auth_handshake(
&mut stream,
config.auth_username.as_deref(),
config.auth_password.as_deref(),
)
.await;
let elapsed = authenticating_start.elapsed().as_millis() as u64;
match result {
Ok(redacted) => {
info!(
peer = %peer_addr,
auth = %redacted,
duration_ms = elapsed,
state = ?ControlState::Authenticating,
"control connection authenticated"
);
if let Some(m) = metrics {
m.record_control_accepted(peer_addr);
m.record_state_duration(ControlState::Authenticating, elapsed);
}
Some(redacted)
}
Err(e) => {
warn!(
peer = %peer_addr,
error = %e,
duration_ms = elapsed,
state = ?ControlState::Authenticating,
"control connection auth failed"
);
if let Some(m) = metrics {
m.record_auth_failure(peer_addr, &e.to_string());
}
return Err(e);
}
}
} else {
// No auth configured: send accept handshake
crate::write_handshake_accept(&mut stream).await?;
info!(
peer = %peer_addr,
state = ?ControlState::Authenticating,
"control connection accepted (no auth)"
);
if let Some(m) = metrics {
m.record_control_accepted(peer_addr);
}
None
};
let ctrl = ControlStream {
stream,
peer_addr,
redacted_auth: redacted,
};
if control_tx.send(ctrl).is_err() {
state.active_control.fetch_sub(1, Ordering::Relaxed);
if let Some(m) = metrics {
m.record_control_closed();
}
warn!(peer = %peer_addr, "control channel closed, cannot add to pool");
}
Ok(())
}
/// Accept external clients and relay them through available control connections.
async fn accept_external_clients(
listener: TcpListener,
config: Arc<ReverseServerConfig>,
cancel: CancellationToken,
mut control_rx: mpsc::UnboundedReceiver<ControlStream>,
metrics: Option<Arc<ReverseMetrics>>,
state: Arc<ReverseServerState>,
) {
let mut relay_tasks = JoinSet::new();
loop {
tokio::select! {
result = listener.accept() => {
match result {
Ok((external_stream, peer_addr)) => {
match state.active_streams.fetch_update(
Ordering::AcqRel,
Ordering::Acquire,
|current| {
(current < config.max_streams_per_listener)
.then_some(current + 1)
},
) {
Ok(_) => {}
Err(current) => {
warn!(
peer = %peer_addr,
active = current,
max = config.max_streams_per_listener,
"dropping external client: max_streams_per_listener reached"
);
state.dropped_stream_limit.fetch_add(1, Ordering::Relaxed);
drop(external_stream);
continue;
}
}
match state.pending_external.fetch_update(
Ordering::AcqRel,
Ordering::Acquire,
|current| {
(current < config.max_pending_external)
.then_some(current + 1)
},
) {
Ok(_) => {}
Err(current) => {
state.active_streams.fetch_sub(1, Ordering::Release);
warn!(
peer = %peer_addr,
pending = current,
max = config.max_pending_external,
"dropping external client: max_pending_external reached"
);
state.dropped_pending_limit.fetch_add(1, Ordering::Relaxed);
drop(external_stream);
continue;
}
}
// Get an available control connection
let control = tokio::select! {
control = control_rx.recv() => control,
_ = cancel.cancelled() => {
state.pending_external.fetch_sub(1, Ordering::Release);
state.active_streams.fetch_sub(1, Ordering::Release);
drop(external_stream);
break;
}
};
match control {
Some(control) => {
state.pending_external.fetch_sub(1, Ordering::Release);
let metrics = metrics.clone();
let state = state.clone();
let idle_timeout = (config.read_timeout_ms > 0).then(|| {
std::time::Duration::from_millis(config.read_timeout_ms)
});
state.active_control.fetch_sub(1, Ordering::Relaxed);
relay_tasks.spawn(async move {
info!(
peer = %peer_addr,
control_peer = %control.peer_addr,
"relaying external client through control connection"
);
if let Some(m) = metrics.as_deref() {
m.record_stream_opened();
m.record_state_duration(ControlState::Ready, 0);
}
let relay_result = relay_bidirectional_with_timeout(
external_stream,
control.stream,
idle_timeout,
)
.await;
match relay_result {
Ok(()) => {
debug!(peer = %peer_addr, "relay finished cleanly");
}
Err(e) => {
debug!(peer = %peer_addr, error = %e, "relay ended");
}
}
if let Some(m) = metrics.as_deref() {
m.record_stream_closed(0);
m.record_control_closed();
}
state.active_streams.fetch_sub(1, Ordering::Release);
debug!(peer = %peer_addr, "relay finished");
});
}
None => {
state.pending_external.fetch_sub(1, Ordering::Release);
state.active_streams.fetch_sub(1, Ordering::Release);
warn!(peer = %peer_addr, "no control connections available, rejecting external client");
drop(external_stream);
}
}
}
Err(e) => {
error!(error = %e, "failed to accept external client");
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
}
}
}
_ = cancel.cancelled() => {
break;
}
}
}
relay_tasks.abort_all();
while relay_tasks.join_next().await.is_some() {}
}
/// Shut down the reverse server.
pub fn shutdown(&self) {
self.cancel.cancel();
}
}
/// A control stream paired with metadata, used when handing the stream off
/// from the auth phase to the relay phase.
pub struct ControlStream {
pub stream: TcpStream,
pub peer_addr: SocketAddr,
pub redacted_auth: Option<String>,
}
/// Helper that exposes the redacted auth form for tests and admin code.
pub fn format_auth_redacted(auth: &str) -> String {
redact_auth(auth)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn is_bind_allowed_with_none() {
let cfg = ReverseServerConfig {
allow_bind: None,
..Default::default()
};
assert!(cfg.is_bind_allowed("127.0.0.1:8080".parse().unwrap()));
}
#[test]
fn is_bind_allowed_with_empty() {
let cfg = ReverseServerConfig {
allow_bind: Some(vec![]),
..Default::default()
};
assert!(cfg.is_bind_allowed("127.0.0.1:8080".parse().unwrap()));
}
#[test]
fn is_bind_allowed_match() {
let cfg = ReverseServerConfig {
allow_bind: Some(vec!["127.0.0.1:8080".parse().unwrap()]),
..Default::default()
};
assert!(cfg.is_bind_allowed("127.0.0.1:8080".parse().unwrap()));
}
#[test]
fn is_bind_allowed_mismatch() {
let cfg = ReverseServerConfig {
allow_bind: Some(vec!["127.0.0.1:8080".parse().unwrap()]),
..Default::default()
};
assert!(!cfg.is_bind_allowed("0.0.0.0:8080".parse().unwrap()));
assert!(!cfg.is_bind_allowed("127.0.0.1:9090".parse().unwrap()));
}
#[test]
fn state_snapshot_round_trip() {
let s = ReverseServerState::default();
s.active_control.fetch_add(3, Ordering::Relaxed);
s.active_streams.fetch_add(2, Ordering::Relaxed);
s.pending_external.fetch_add(1, Ordering::Relaxed);
s.denied_bind.fetch_add(1, Ordering::Relaxed);
s.dropped_stream_limit.fetch_add(4, Ordering::Relaxed);
s.dropped_pending_limit.fetch_add(5, Ordering::Relaxed);
let snap = s.snapshot();
assert_eq!(snap.active_control, 3);
assert_eq!(snap.active_streams, 2);
assert_eq!(snap.pending_external, 1);
assert_eq!(snap.denied_bind, 1);
assert_eq!(snap.dropped_stream_limit, 4);
assert_eq!(snap.dropped_pending_limit, 5);
}
#[test]
fn format_auth_redacted_basic() {
assert_eq!(format_auth_redacted("user:pass"), "user:****");
}
#[test]
fn same_bind_v4() {
let a: SocketAddr = "127.0.0.1:8080".parse().unwrap();
let b: SocketAddr = "127.0.0.1:8080".parse().unwrap();
assert!(same_bind(&a, &b));
}
#[test]
fn same_bind_different_port() {
let a: SocketAddr = "127.0.0.1:8080".parse().unwrap();
let b: SocketAddr = "127.0.0.1:9090".parse().unwrap();
assert!(!same_bind(&a, &b));
}
#[test]
fn validate_loopback_ok() {
let cfg = ReverseServerConfig {
control_bind: "127.0.0.1:0".parse().unwrap(),
external_bind: Some("127.0.0.1:0".parse().unwrap()),
..Default::default()
};
assert!(cfg.validate().is_ok());
}
#[test]
fn validate_no_external_bind_ok() {
let cfg = ReverseServerConfig {
control_bind: "127.0.0.1:0".parse().unwrap(),
external_bind: None,
..Default::default()
};
assert!(cfg.validate().is_ok());
}
#[test]
fn validate_non_loopback_without_auth_rejected() {
let cfg = ReverseServerConfig {
control_bind: "127.0.0.1:0".parse().unwrap(),
external_bind: Some("0.0.0.0:9000".parse().unwrap()),
auth_username: None,
auth_password: None,
..Default::default()
};
let err = cfg.validate().unwrap_err();
assert!(
matches!(err, ProtocolError::ConfigInvalid(_)),
"got: {err:?}"
);
}
#[test]
fn validate_non_loopback_with_auth_but_no_allowlist_rejected() {
let cfg = ReverseServerConfig {
control_bind: "127.0.0.1:0".parse().unwrap(),
external_bind: Some("0.0.0.0:9000".parse().unwrap()),
auth_username: Some("user".to_string()),
auth_password: Some("pass".to_string()),
allow_bind: None,
..Default::default()
};
let err = cfg.validate().unwrap_err();
assert!(
matches!(err, ProtocolError::ConfigInvalid(_)),
"got: {err:?}"
);
}
#[test]
fn validate_non_loopback_with_auth_and_allowlist_ok() {
let cfg = ReverseServerConfig {
control_bind: "127.0.0.1:0".parse().unwrap(),
external_bind: Some("0.0.0.0:9000".parse().unwrap()),
auth_username: Some("user".to_string()),
auth_password: Some("pass".to_string()),
allow_bind: Some(vec!["0.0.0.0:9000".parse().unwrap()]),
..Default::default()
};
assert!(cfg.validate().is_ok());
}
#[test]
fn validate_ipv6_loopback_ok() {
let cfg = ReverseServerConfig {
control_bind: "127.0.0.1:0".parse().unwrap(),
external_bind: Some("[::1]:9000".parse().unwrap()),
..Default::default()
};
assert!(cfg.validate().is_ok());
}
#[test]
fn validate_ipv6_non_loopback_without_auth_rejected() {
let cfg = ReverseServerConfig {
control_bind: "127.0.0.1:0".parse().unwrap(),
external_bind: Some("[2001:db8::1]:9000".parse().unwrap()),
..Default::default()
};
let err = cfg.validate().unwrap_err();
assert!(
matches!(err, ProtocolError::ConfigInvalid(_)),
"got: {err:?}"
);
}
}