Skip to content

Commit 9d30cfe

Browse files
committed
vey-proxy: add nested connect escaper methods
1 parent a2a2405 commit 9d30cfe

71 files changed

Lines changed: 1676 additions & 220 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

lib/vey-io-ext/src/lib.rs

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,9 @@
33
* SPDX-FileCopyrightText: 2023-2025 ByteDance and/or its affiliates.
44
*/
55

6+
#[derive(Default)]
7+
pub struct NilLimitedStats(());
8+
69
mod cache;
710
mod limit;
811
mod listen;

lib/vey-io-ext/src/stream/limited/mod.rs

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -4,12 +4,10 @@
44
*/
55

66
mod read;
7-
pub use read::{
8-
ArcLimitedReaderStats, LimitedReader, LimitedReaderStats, NilLimitedReaderStats, SizedReader,
9-
};
7+
pub use read::{ArcLimitedReaderStats, LimitedReader, LimitedReaderStats, SizedReader};
108

119
mod stream;
1210
pub use stream::LimitedStream;
1311

1412
mod write;
15-
pub use write::{ArcLimitedWriterStats, LimitedWriter, LimitedWriterStats, NilLimitedWriterStats};
13+
pub use write::{ArcLimitedWriterStats, LimitedWriter, LimitedWriterStats};

lib/vey-io-ext/src/stream/limited/read.rs

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@ use pin_project_lite::pin_project;
1616
use tokio::io::{AsyncRead, AsyncWrite, ReadBuf};
1717
use tokio::time::{Instant, Sleep};
1818

19+
use crate::NilLimitedStats;
1920
use crate::limit::{GlobalLimitGroup, GlobalStreamLimit, StreamLimitAction, StreamLimiter};
2021
use crate::stream::AsyncStream;
2122

@@ -24,10 +25,7 @@ pub trait LimitedReaderStats {
2425
}
2526
pub type ArcLimitedReaderStats = Arc<dyn LimitedReaderStats + Send + Sync>;
2627

27-
#[derive(Default)]
28-
pub struct NilLimitedReaderStats(());
29-
30-
impl LimitedReaderStats for NilLimitedReaderStats {
28+
impl LimitedReaderStats for NilLimitedStats {
3129
fn add_read_bytes(&self, _size: usize) {}
3230
}
3331

lib/vey-io-ext/src/stream/limited/write.rs

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -16,17 +16,15 @@ use pin_project_lite::pin_project;
1616
use tokio::io::AsyncWrite;
1717
use tokio::time::{Instant, Sleep};
1818

19+
use crate::NilLimitedStats;
1920
use crate::limit::{GlobalLimitGroup, GlobalStreamLimit, StreamLimitAction, StreamLimiter};
2021

2122
pub trait LimitedWriterStats {
2223
fn add_write_bytes(&self, size: usize);
2324
}
2425
pub type ArcLimitedWriterStats = Arc<dyn LimitedWriterStats + Send + Sync>;
2526

26-
#[derive(Default)]
27-
pub struct NilLimitedWriterStats(());
28-
29-
impl LimitedWriterStats for NilLimitedWriterStats {
27+
impl LimitedWriterStats for NilLimitedStats {
3028
fn add_write_bytes(&self, _size: usize) {}
3129
}
3230

lib/vey-io-ext/src/udp/stats.rs

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,8 @@
55

66
use std::sync::Arc;
77

8+
use crate::NilLimitedStats;
9+
810
pub trait LimitedRecvStats {
911
fn add_recv_bytes(&self, size: usize);
1012
fn add_recv_packet(&self) {
@@ -14,6 +16,14 @@ pub trait LimitedRecvStats {
1416
}
1517
pub type ArcLimitedRecvStats = Arc<dyn LimitedRecvStats + Send + Sync>;
1618

19+
impl LimitedRecvStats for NilLimitedStats {
20+
fn add_recv_bytes(&self, _size: usize) {}
21+
22+
fn add_recv_packet(&self) {}
23+
24+
fn add_recv_packets(&self, _n: usize) {}
25+
}
26+
1727
pub trait LimitedSendStats {
1828
fn add_send_bytes(&self, size: usize);
1929
fn add_send_packet(&self) {
@@ -23,6 +33,14 @@ pub trait LimitedSendStats {
2333
}
2434
pub type ArcLimitedSendStats = Arc<dyn LimitedSendStats + Send + Sync>;
2535

36+
impl LimitedSendStats for NilLimitedStats {
37+
fn add_send_bytes(&self, _size: usize) {}
38+
39+
fn add_send_packet(&self) {}
40+
41+
fn add_send_packets(&self, _n: usize) {}
42+
}
43+
2644
#[cfg(test)]
2745
mod tests {
2846
use std::sync::atomic::{AtomicUsize, Ordering};

vey-proxy/src/escape/comply_audit/mod.rs

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -215,6 +215,34 @@ impl EscaperInternal for ComplyAuditEscaper {
215215
audit_ctx.set_handle(self.audit_handle.clone());
216216
}
217217

218+
async fn _nested_tcp_connect(
219+
&self,
220+
task_conf: &TcpConnectTaskConf<'_>,
221+
egress_notes: &mut EgressNotes,
222+
task_notes: &ServerTaskNotes,
223+
audit_ctx: &mut AuditContext,
224+
) -> TcpConnectResult {
225+
egress_notes.escaper.clone_from(&self.config.name);
226+
self.stats.add_request_passed();
227+
self._update_audit_context(audit_ctx);
228+
self.next
229+
._nested_tcp_connect(task_conf, egress_notes, task_notes, audit_ctx)
230+
.await
231+
}
232+
233+
async fn _nested_udp_connect(
234+
&self,
235+
task_conf: &UdpConnectTaskConf<'_>,
236+
egress_notes: &mut EgressNotes,
237+
task_notes: &ServerTaskNotes,
238+
) -> UdpConnectResult {
239+
egress_notes.escaper.clone_from(&self.config.name);
240+
self.stats.add_request_passed();
241+
self.next
242+
._nested_udp_connect(task_conf, egress_notes, task_notes)
243+
.await
244+
}
245+
218246
async fn _new_http_forward_connection(
219247
&self,
220248
_task_conf: &TcpConnectTaskConf<'_>,

vey-proxy/src/escape/comply_context/mod.rs

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -240,6 +240,35 @@ impl EscaperInternal for ComplyContextEscaper {
240240
}
241241
}
242242

243+
async fn _nested_tcp_connect(
244+
&self,
245+
task_conf: &TcpConnectTaskConf<'_>,
246+
egress_notes: &mut EgressNotes,
247+
task_notes: &ServerTaskNotes,
248+
audit_ctx: &mut AuditContext,
249+
) -> TcpConnectResult {
250+
egress_notes.escaper.clone_from(&self.config.name);
251+
self._update_egress_path(task_notes);
252+
self.stats.add_request_passed();
253+
self.next
254+
._nested_tcp_connect(task_conf, egress_notes, task_notes, audit_ctx)
255+
.await
256+
}
257+
258+
async fn _nested_udp_connect(
259+
&self,
260+
task_conf: &UdpConnectTaskConf<'_>,
261+
egress_notes: &mut EgressNotes,
262+
task_notes: &ServerTaskNotes,
263+
) -> UdpConnectResult {
264+
egress_notes.escaper.clone_from(&self.config.name);
265+
self._update_egress_path(task_notes);
266+
self.stats.add_request_passed();
267+
self.next
268+
._nested_udp_connect(task_conf, egress_notes, task_notes)
269+
.await
270+
}
271+
243272
async fn _new_http_forward_connection(
244273
&self,
245274
_task_conf: &TcpConnectTaskConf<'_>,

vey-proxy/src/escape/direct_fixed/http_forward/mod.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@
66

77
use std::sync::Arc;
88

9-
use vey_io_ext::{AsyncStream, LimitedBufReader, LimitedWriter, NilLimitedReaderStats};
9+
use vey_io_ext::{AsyncStream, LimitedBufReader, LimitedWriter, NilLimitedStats};
1010

1111
use super::{DirectFixedEscaper, DirectFixedEscaperStats};
1212
use crate::escape::EgressNotes;
@@ -94,7 +94,7 @@ impl DirectFixedEscaper {
9494

9595
let ups_r = LimitedBufReader::new_unlimited(
9696
ups_r,
97-
Arc::new(NilLimitedReaderStats::default()),
97+
Arc::new(NilLimitedStats::default()),
9898
wrapper_stats.clone(),
9999
);
100100
let ups_w = LimitedWriter::new(ups_w, wrapper_stats);

vey-proxy/src/escape/direct_fixed/mod.rs

Lines changed: 26 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -363,7 +363,7 @@ impl Escaper for DirectFixedEscaper {
363363
) -> UdpConnectResult {
364364
self.stats.interface.add_udp_connect_attempted();
365365
egress_notes.escaper.clone_from(&self.config.name);
366-
self.udp_connect_to(task_conf, egress_notes, task_notes, task_stats)
366+
self.new_udp_connection(task_conf, egress_notes, task_notes, task_stats)
367367
.await
368368
}
369369

@@ -428,6 +428,31 @@ impl EscaperInternal for DirectFixedEscaper {
428428
capability
429429
}
430430

431+
async fn _nested_tcp_connect(
432+
&self,
433+
task_conf: &TcpConnectTaskConf<'_>,
434+
egress_notes: &mut EgressNotes,
435+
task_notes: &ServerTaskNotes,
436+
_audit_ctx: &mut AuditContext,
437+
) -> TcpConnectResult {
438+
self.stats.interface.add_tcp_connect_attempted();
439+
egress_notes.escaper.clone_from(&self.config.name);
440+
self.nested_tcp_connect(task_conf, egress_notes, task_notes)
441+
.await
442+
}
443+
444+
async fn _nested_udp_connect(
445+
&self,
446+
task_conf: &UdpConnectTaskConf<'_>,
447+
egress_notes: &mut EgressNotes,
448+
task_notes: &ServerTaskNotes,
449+
) -> UdpConnectResult {
450+
self.stats.interface.add_udp_connect_attempted();
451+
egress_notes.escaper.clone_from(&self.config.name);
452+
self.nested_udp_connect(task_conf, egress_notes, task_notes)
453+
.await
454+
}
455+
431456
async fn _new_http_forward_connection(
432457
&self,
433458
task_conf: &TcpConnectTaskConf<'_>,

vey-proxy/src/escape/direct_fixed/stats.rs

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@ use std::sync::Arc;
99
use arc_swap::ArcSwapOption;
1010

1111
use vey_daemon::stat::remote::{TcpConnectionTaskRemoteStats, UdpConnectTaskRemoteStats};
12-
use vey_io_ext::{LimitedReaderStats, LimitedWriterStats};
12+
use vey_io_ext::{LimitedReaderStats, LimitedRecvStats, LimitedSendStats, LimitedWriterStats};
1313
use vey_types::metrics::{MetricTagMap, NodeName};
1414
use vey_types::stats::{StatId, TcpIoSnapshot, UdpIoSnapshot};
1515

@@ -124,6 +124,26 @@ impl LimitedWriterStats for DirectFixedEscaperStats {
124124
}
125125
}
126126

127+
impl LimitedRecvStats for DirectFixedEscaperStats {
128+
fn add_recv_bytes(&self, size: usize) {
129+
self.udp.io.add_in_bytes(size as u64);
130+
}
131+
132+
fn add_recv_packets(&self, n: usize) {
133+
self.udp.io.add_in_packets(n);
134+
}
135+
}
136+
137+
impl LimitedSendStats for DirectFixedEscaperStats {
138+
fn add_send_bytes(&self, size: usize) {
139+
self.udp.io.add_out_bytes(size as u64);
140+
}
141+
142+
fn add_send_packets(&self, n: usize) {
143+
self.udp.io.add_out_packets(n);
144+
}
145+
}
146+
127147
impl TcpConnectionTaskRemoteStats for DirectFixedEscaperStats {
128148
fn add_read_bytes(&self, size: u64) {
129149
self.tcp.io.add_in_bytes(size);

0 commit comments

Comments
 (0)