Hello World

吞风吻雨葬落日 欺山赶海踏雪径

0%

Spring Boot 2.x 内嵌 VLESS / Trojan over WebSocket 代理隧道实现

1. 技术栈与约束

维度 版本 / 约束 影响
Spring Boot 2.5.6 WebSocket API 为 javax.websocket.* / Spring 5.3.x,非 jakarta;不可用 Boot 3.x 特性
JDK 1.8 无虚拟线程、无 java.net.http.HttpClient(Java 11+)、无 List.of / var / String.repeat
内嵌容器 Tomcat(spring-boot-starter-web 默认) tomcat-embed-websocket 随 starter 引入;WS 消息回调跑在 Tomcat WS worker 线程池(默认上限 200)
代码组织 单 Controller 文件 一个类同时承担 HTTP 路由、WS 端点注册、WS 隧道逻辑三种角色
协议范围 VLESS + Trojan 首包嗅探分流,两种协议共享同一 WS 端点

2. 依赖引入

pom.xml 新增 spring-boot-starter-websocket,显式版本与 spring-boot-starter-web 对齐:

1
2
3
4
5
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-websocket</artifactId>
<version>2.5.6</version>
</dependency>

3. 架构设计(单 Controller 三合一)

一个类同时承担三种角色:

1
2
3
4
5
ProxyTunnelController
├── @Controller + @RequestMapping("/api/v1/proxy") ← HTTP 路由(订阅接口)
├── @EnableWebSocket + implements WebSocketConfigurer ← WS 端点注册(常量路径)
├── extends BinaryWebSocketHandler ← WS 隧道逻辑(协议嗅探 / 解析 / 转发)
└── @Bean ServletServerContainerFactoryBean ← 二进制缓冲区 1MB(lite-mode Bean)

要点:

  1. @EnableWebSocket 放在 @Controller 类上即生效(lite mode,Spring 会处理其 @Import),无需独立 @Configuration 类。
  2. registerWebSocketHandlersregistry.addHandler(this, WSPATH).setAllowedOrigins("*")WSPATH 为完整路径常量。Tomcat 对路径按最长匹配处理,与 HTTP 路由无冲突。
  3. @Bean ServletServerContainerFactoryBean:Boot 2.5 的 WebSocketServletAutoConfiguration@ConditionalOnMissingBean,用户 Bean 自动优先生效。

4. 协议格式

4.1 VLESS 首包字节布局

1
[version:1B][uuid:16B][addonsLen:1B][addons:nB][cmd:1B][port:2B BE][atyp:1B][addr:变长]
  • version:固定 0x00
  • uuid:16 字节认证凭据(UUID 去 - 后 hex 解码)
  • addonsLen + addons:附加数据长度与内容
  • cmd:必须为 0x01(TCP CONNECT)
  • port:2 字节大端序
  • atyp1=IPv4(4 字节点分十进制)/ 2=域名(1 字节长度 + UTF-8 文本)/ 3=IPv6(16 字节按 2 字节分组)
  • 隧道建立后服务端需先回发 2 字节握手响应 {0x00, 0x00}

4.2 Trojan 首包字节布局

1
[sha224(uuid-hex):56B][\r\n?][cmd:1B=0x01][atyp:1B][addr:变长][port:2B BE][\r\n?][payload]
  • sha224(uuid-hex):UUID 字符串的 SHA-224 摘要转 hex 字符串的 ASCII 字节(56 字节)
  • 可选 CRLF(\r\n
  • cmd:必须为 0x01(CONNECT)
  • atyp1=IPv4 / 3=域名 / 4=IPv6
  • 端口后可选尾部 CRLF,随后为业务数据

4.3 嗅探顺序

先按 VLESS 判定(首字节 0 且 UUID 匹配),再按 Trojan 判定(SHA-224 前缀匹配),均不匹配返回 null → 关闭连接。凭据比较采用恒定时间比较,防止基于时序的侧信道攻击。

5. 线程模型

Java 8 无虚拟线程,模型设计为「两条通道 + 一把锁」:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
                         ┌──────────────────────────────┐
WS 客户端 ◄── Binary ──► │ Tomcat WS worker 线程池 │
│ (默认 200,消息回调串行) │
└──────┬───────────────────────┘
回调线程直写 │ sendLock 同步
(阻塞式,天然背压)

Socket.getOutputStream() ──► TCP 目标

泵线程读 │
┌──────┴───────────────────────┐
│ 自建有界 ThreadPoolExecutor │
│ core=8 max=200 daemon │
└──────────────────────────────┘
  1. WS → TCP:在 Tomcat WS 回调线程中直接 socket.getOutputStream().write(...)(阻塞式)。回调阻塞期间容器不再向该 session 分发下一帧 → TCP 接收窗口收缩 → 客户端自动减速,天然背压。代价是占用容器 WS 线程,并发上限受 Tomcat WS 线程池约束(单实例 200,可通过 -Dorg.apache.tomcat.websocket.threadPoolSize=N 调整)。

  2. TCP → WS:每隧道一个泵任务,提交到自建线程池。拒绝策略为「关闭隧道」而非 CallerRuns(CallerRuns 会反压到泵提交线程,语义混乱)。

  3. sendMessage 线程安全:泵线程与 WS 回调 / 关闭并发调用 session.sendMessage 会抛 IllegalStateException,全部发送路径套 synchronized (sendLock)

6. 二进制缓冲区

不设置时 Tomcat 默认 8KB,大帧(如刷页面、下载)会被容器直接关闭连接。通过 ServletServerContainerFactoryBean(内部走 javax ServerContainer)将 maxBinaryMessageBufferSize 设为 1MB,并设置 maxSessionIdleTimeout 10 分钟,避免 SLB / Ingress 空闲超时前服务端先杀会话。

7. 隧道生命周期

7.1 建立流程(openTunnel)

  1. 解析首包得到协议类型与目标地址,凭据或格式非法则关闭连接
  2. 黑名单校验目标域名,命中则关闭连接
  3. 建立 TCP 连接(setTcpNoDelay(true) + 限时建连)并注册隧道
  4. VLESS 协议先发送 2 字节握手响应,再写入首包剩余业务数据
  5. 提交 TCP 到 WS 的泵送任务;任一步骤失败即关闭连接,保证无残留

注册时机:握手响应与首包直写之前tunnels.put(session, tunnel),保证后续帧可路由。

7.2 WebSocket 生命周期回调

回调 行为
handleBinaryMessage 隧道未建立时调用 openTunnel;已建立时 tunnel.write(payload) 直写 TCP,写失败即 closeTunnel
afterConnectionClosed 客户端断开 → closeTunnel
handleTransportError 传输异常 → 记录日志并 closeTunnel
@PreDestroy destroy 遍历关闭全部活跃隧道、清空映射、pumpPool.shutdownNow()

7.3 关闭与清理

  • closeTunnelmarkClosed() 保证关闭只执行一次,随后移除映射、关闭 Socket 与 Session(均幂等 try/catch)
  • Tunnel.markClosedsynchronized (sendLock) 内置位,防止关闭与发送竞争
  • ConcurrentHashMap<WebSocketSession, Tunnel> 记录活跃隧道

8. 订阅接口

GET /api/v1/proxy/subscribe,返回 base64 编码的 VLESS / Trojan 节点订阅内容。流程:

  1. 缓存窗口过期时占用 ISP_LOCK 窗口并触发后台 ISP 探测,探测失败保留旧值
  2. 以常量接入域名 / 端口与当前 ISP 拼装两条节点 URL
  3. 拼接后 base64 编码,以 text/plain 返回

节点 URL 模板(WSPATH 中的 / 编码为 %2F):

1
2
vless://{UUID}@{DOMAIN}:{PORT}?encryption=none&security=tls&sni={DOMAIN}&fp=chrome&type=ws&host={DOMAIN}&path={WSPATH_ENCODED}#{NAME}-{ISP}
trojan://{UUID}@{DOMAIN}:{PORT}?security=tls&sni={DOMAIN}&fp=chrome&type=ws&host={DOMAIN}&path={WSPATH_ENCODED}#{NAME}-{ISP}

ISP 探测细节:

  • 后台守护线程执行,避免占用容器请求线程
  • 双源顺序探测,任一源返回有效数据即返回,全部失败返回 null(保留旧值)
  • 字段优先级:国家码取 country_code / countryCode,运营商取 isp / org / organization
  • 结果清洗:空格替换为下划线,去除首尾 -,拼为「国家码-运营商」
  • 专用 RestTemplate 连接与读取超时均为 1500ms

9. 可调参数

常量 说明
UUID 自定义 UUID VLESS 首包 UUID 认证与 Trojan 首包哈希种子,上线前必须更换
WSPATH /api/v1/proxy/websocket WebSocket 接入路径
CONNECT_TIMEOUT_MS 10_000 TCP 连接建立超时(毫秒)
PUMP_BUFFER_SIZE 16 * 1024 TCP 向 WS 泵送的读缓冲区大小
MAX_BINARY_BUFFER 1024 * 1024 二进制帧缓冲上限(字节)
IDLE_TIMEOUT_MS 10min 会话空闲超时
PUMP_CORE / PUMP_MAX / PUMP_KEEP_ALIVE_SEC 8 / 200 / 60 泵线程池参数
SUB_DOMAIN / SUB_PORT / SUB_SECURITY / SUB_NAME your-domain.com / 443 / tls / ProxyNode 订阅节点接入参数
SUB_TTL_MS 10min 订阅 ISP 信息刷新间隔
ISP_TIMEOUT_MS 1500 ISP 探测单源超时(毫秒)
BLOCKED_DOMAINS 10 个测速域名 目标域名黑名单(含子域名)

10. 完整代码

以下为完整的单文件实现,包名与路径均为示例值,按实际项目调整即可:

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
package com.example.proxy;

import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.context.annotation.Bean;
import org.springframework.http.HttpEntity;
import org.springframework.http.HttpHeaders;
import org.springframework.http.HttpMethod;
import org.springframework.http.MediaType;
import org.springframework.http.ResponseEntity;
import org.springframework.http.client.SimpleClientHttpRequestFactory;
import org.springframework.stereotype.Controller;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.client.RestTemplate;
import org.springframework.web.socket.BinaryMessage;
import org.springframework.web.socket.CloseStatus;
import org.springframework.web.socket.WebSocketSession;
import org.springframework.web.socket.config.annotation.EnableWebSocket;
import org.springframework.web.socket.config.annotation.WebSocketConfigurer;
import org.springframework.web.socket.config.annotation.WebSocketHandlerRegistry;
import org.springframework.web.socket.handler.BinaryWebSocketHandler;
import org.springframework.web.socket.server.standard.ServletServerContainerFactoryBean;

import javax.annotation.PreDestroy;
import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
import java.net.InetSocketAddress;
import java.net.Socket;
import java.nio.ByteBuffer;
import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
import java.security.NoSuchAlgorithmException;
import java.util.Arrays;
import java.util.Base64;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.SynchronousQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;

/**
* 代理隧道控制器。
*
* <p>通过 WebSocket 二进制帧承载两条协议隧道(VLESS / Trojan),首包完成协议识别与目标地址解析后,
* 建立 TCP 连接并在 WebSocket 与 TCP 之间双向转发;同时提供订阅接口返回两条协议的节点 URL。
* 全部逻辑集中在单个类内,不依赖额外业务模块。
*
* <p>参数(UUID / 路径前缀 / 池大小 / 超时等)以类常量形式集中定义,可按环境调整。
*/
@Slf4j
@Controller
@EnableWebSocket
@RequestMapping("/api/v1/proxy")
public class ProxyTunnelController extends BinaryWebSocketHandler implements WebSocketConfigurer {

// ==================== 可调参数 ====================
/** 认证凭据,用于 VLESS 首包 UUID 校验与 Trojan 首包哈希校验 */
private static final String UUID = "your-uuid-here";
/** WebSocket 接入路径,与 HTTP 路由前缀对齐 */
private static final String WSPATH = "/api/v1/proxy/websocket";
/** TCP 连接建立超时(毫秒) */
private static final int CONNECT_TIMEOUT_MS = 10_000;
/** TCP 向 WS 泵送的读缓冲区大小 */
private static final int PUMP_BUFFER_SIZE = 16 * 1024;
/** 二进制帧缓冲上限(字节),避免默认 8KB 导致大帧被容器断开 */
private static final int MAX_BINARY_BUFFER = 1024 * 1024;
/** 会话空闲超时(毫秒) */
private static final long IDLE_TIMEOUT_MS = TimeUnit.MINUTES.toMillis(10);
/** 泵线程池核心线程数 */
private static final int PUMP_CORE = 8;
/** 泵线程池最大线程数 */
private static final int PUMP_MAX = 200;
/** 泵线程池空闲存活时间 */
private static final long PUMP_KEEP_ALIVE_SEC = 60L;
/** 订阅节点接入域名,上线前按环境替换 */
private static final String SUB_DOMAIN = "your-domain.com";
/** 订阅节点接入端口 */
private static final int SUB_PORT = 443;
/** 订阅节点传输安全参数(tls / none) */
private static final String SUB_SECURITY = "tls";
/** 订阅节点备注名前缀 */
private static final String SUB_NAME = "ProxyNode";
/** 订阅元信息刷新间隔(毫秒) */
private static final long SUB_TTL_MS = TimeUnit.MINUTES.toMillis(10);
/** ISP 探测接口的单源超时(毫秒) */
private static final long ISP_TIMEOUT_MS = 1500L;
/** ISP 双源探测端点 */
private static final String[] ISP_ENDPOINTS = {"https://api.ip.sb/geoip", "http://ip-api.com/json"};

/** 目标地址黑名单(及其子域名) */
private static final List<String> BLOCKED_DOMAINS = Collections.unmodifiableList(Arrays.asList(
"speedtest.net", "fast.com", "speedtest.cn", "speed.cloudflare.com",
"speedof.me", "testmy.net", "bandwidth.place", "speed.io",
"librespeed.org", "speedcheck.org"));

// ==================== 协议常量 ====================
/** VLESS 协议 UUID 的二进制形式(hex 解码,16 字节) */
private static final byte[] UUID_BYTES = hexToBytes(UUID.replace("-", ""));
/** Trojan 协议首包哈希:SHA-224(UUID) 的 hex 字符串字节(56 字节) */
private static final byte[] TROJAN_HASH = buildTrojanHash();

/** 当前节点 ISP 描述(国家码-运营商),探测失败保留旧值 */
private static volatile String nodeIsp = "Unknown";
/** 上次 ISP 探测时间戳,配合 ISP_LOCK 避免并发重复探测 */
private static volatile long ispRefreshedAt = 0L;
/** ISP 探测窗口锁 */
private static final Object ISP_LOCK = new Object();
/** ISP 探测响应解析器 */
private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
/** ISP 探测 HTTP 客户端 */
private static final RestTemplate ISP_REST = buildIspRestTemplate();

/** 活跃隧道集合:session -> tunnel */
private final ConcurrentHashMap<WebSocketSession, Tunnel> tunnels = new ConcurrentHashMap<>();
/** TCP 向 WS 方向泵送线程池,有界且拒绝即关闭隧道 */
private final ThreadPoolExecutor pumpPool = new ThreadPoolExecutor(
PUMP_CORE, PUMP_MAX, PUMP_KEEP_ALIVE_SEC, TimeUnit.SECONDS,
new SynchronousQueue<Runnable>(),
r -> {
Thread t = new Thread(r, "proxy-tunnel-pump");
t.setDaemon(true);
return t;
},
(r, executor) -> {
log.warn("Proxy pump pool saturated, close tunnel");
if (r instanceof PumpTask) {
((PumpTask) r).closeNow();
}
});

// ==================== WebSocket 端点注册与容器配置 ====================

@Override
public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) {
registry.addHandler(this, WSPATH).setAllowedOrigins("*");
}

@Bean
public ServletServerContainerFactoryBean proxyWebSocketContainer() {
ServletServerContainerFactoryBean container = new ServletServerContainerFactoryBean();
container.setMaxBinaryMessageBufferSize(MAX_BINARY_BUFFER);
container.setMaxSessionIdleTimeout(IDLE_TIMEOUT_MS);
return container;
}

// ==================== 订阅接口 ====================

/**
* 返回 base64 编码的 VLESS / Trojan 节点订阅内容。
* 步骤:
* 1. 缓存窗口过期时占用窗口并触发后台 ISP 探测,探测失败保留旧值
* 2. 以常量接入域名/端口与当前 ISP 拼装两条节点 URL
* 3. 拼接后 base64 编码,以 text/plain 返回
*/
@GetMapping("/subscribe")
public ResponseEntity<String> subscribe() {
synchronized (ISP_LOCK) {
if (System.currentTimeMillis() - ispRefreshedAt > SUB_TTL_MS) {
ispRefreshedAt = System.currentTimeMillis();
refreshIspAsync();
}
}
String namePart = SUB_NAME + "-" + nodeIsp;
String wsPathEncoded = WSPATH.replace("/", "%2F");
String vlessUrl = "vless://" + UUID + "@" + SUB_DOMAIN + ":" + SUB_PORT
+ "?encryption=none&security=" + SUB_SECURITY + "&sni=" + SUB_DOMAIN
+ "&fp=chrome&type=ws&host=" + SUB_DOMAIN + "&path=" + wsPathEncoded + "#" + namePart;
String trojanUrl = "trojan://" + UUID + "@" + SUB_DOMAIN + ":" + SUB_PORT
+ "?security=" + SUB_SECURITY + "&sni=" + SUB_DOMAIN
+ "&fp=chrome&type=ws&host=" + SUB_DOMAIN + "&path=" + wsPathEncoded + "#" + namePart;
String base64 = Base64.getEncoder()
.encodeToString((vlessUrl + "\n" + trojanUrl).getBytes(StandardCharsets.UTF_8));
return ResponseEntity.ok().contentType(MediaType.TEXT_PLAIN).body(base64 + "\n");
}

/**
* 后台线程执行 ISP 探测,避免占用容器请求线程。
*/
private static void refreshIspAsync() {
Thread t = new Thread(() -> {
String result = fetchIsp();
if (result != null) {
nodeIsp = result;
}
}, "proxy-sub-isp");
t.setDaemon(true);
t.start();
}

/**
* 双源顺序探测出口 ISP 信息,任一源返回有效数据即拼为"国家码-运营商"返回,全部失败返回 null。
*/
private static String fetchIsp() {
for (String endpoint : ISP_ENDPOINTS) {
try {
HttpHeaders headers = new HttpHeaders();
headers.set("User-Agent", "Mozilla/5.0");
ResponseEntity<String> resp = ISP_REST.exchange(
endpoint, HttpMethod.GET, new HttpEntity<Void>(headers), String.class);
JsonNode node = OBJECT_MAPPER.readTree(resp.getBody());
String code = firstText(node, "country_code", "countryCode");
String provider = firstText(node, "isp", "org", "organization");
String cleaned = (code + "-" + provider).replace(" ", "_");
while (cleaned.startsWith("-")) {
cleaned = cleaned.substring(1);
}
while (cleaned.endsWith("-")) {
cleaned = cleaned.substring(0, cleaned.length() - 1);
}
if (!cleaned.isEmpty()) {
return cleaned;
}
} catch (Exception e) {
log.debug("ISP source {} failed: {}", endpoint, e.getMessage());
}
}
return null;
}

/**
* 按候选字段顺序取第一个非空文本值。
*/
private static String firstText(JsonNode node, String... fields) {
for (String field : fields) {
JsonNode value = node.get(field);
if (value != null && value.isTextual() && !value.asText().isEmpty()) {
return value.asText();
}
}
return "";
}

/**
* 构建 ISP 探测专用 HTTP 客户端,连接与读取超时均为 ISP_TIMEOUT_MS。
*/
private static RestTemplate buildIspRestTemplate() {
SimpleClientHttpRequestFactory factory = new SimpleClientHttpRequestFactory();
factory.setConnectTimeout((int) ISP_TIMEOUT_MS);
factory.setReadTimeout((int) ISP_TIMEOUT_MS);
return new RestTemplate(factory);
}

// ==================== WebSocket 生命周期 ====================

@Override
protected void handleBinaryMessage(WebSocketSession session, BinaryMessage message) {
Tunnel tunnel = tunnels.get(session);
if (tunnel == null) {
openTunnel(session, message.getPayload());
return;
}
try {
tunnel.write(message.getPayload());
} catch (IOException e) {
log.debug("Proxy write to tcp failed, close tunnel: {}", e.getMessage());
closeTunnel(tunnel);
}
}

@Override
public void afterConnectionClosed(WebSocketSession session, CloseStatus status) {
closeTunnel(tunnels.get(session));
}

@Override
public void handleTransportError(WebSocketSession session, Throwable exception) {
log.debug("Proxy transport error, close tunnel: {}", exception.getMessage());
closeTunnel(tunnels.get(session));
}

@PreDestroy
public void destroy() {
for (Tunnel tunnel : tunnels.values()) {
closeTunnel(tunnel);
}
tunnels.clear();
pumpPool.shutdownNow();
}

// ==================== 隧道建立 ====================

/**
* 建立隧道。
* 步骤:
* 1. 解析首包得到协议类型与目标地址,凭据或格式非法则关闭连接
* 2. 黑名单校验目标域名,命中则关闭连接
* 3. 建立 TCP 连接并注册隧道
* 4. VLESS 协议先发送 2 字节握手响应,再写入首包剩余业务数据
* 5. 提交 TCP 到 WS 的泵送任务,任一步骤失败即关闭连接,保证无残留
*/
private void openTunnel(WebSocketSession session, ByteBuffer payload) {
byte[] msg = toBytes(payload);
Target target = sniff(msg);
if (target == null) {
log.debug("Proxy handshake rejected: unsupported protocol or bad credential");
closeQuietly(session);
return;
}
if (isBlockedDomain(target.host)) {
log.debug("Proxy handshake rejected: host blocked: {}", target.host);
closeQuietly(session);
return;
}
Socket socket = new Socket();
try {
socket.setTcpNoDelay(true);
socket.connect(new InetSocketAddress(target.host, target.port), CONNECT_TIMEOUT_MS);
} catch (IOException e) {
log.debug("Proxy connect failed to {}:{}: {}", target.host, target.port, e.getMessage());
closeQuietly(socket);
closeQuietly(session);
return;
}
Tunnel tunnel = new Tunnel(session, socket);
tunnels.put(session, tunnel);
try {
if (target.vless) {
session.sendMessage(new BinaryMessage(new byte[]{0x00, 0x00}));
}
if (target.payloadOffset < msg.length) {
OutputStream out = socket.getOutputStream();
out.write(msg, target.payloadOffset, msg.length - target.payloadOffset);
out.flush();
}
} catch (IOException e) {
log.debug("Proxy write first payload failed: {}", e.getMessage());
closeTunnel(tunnel);
return;
}
pumpPool.execute(new PumpTask(tunnel));
}

/**
* 关闭隧道:仅首次调用生效,随后清理映射并关闭两端。
*/
private void closeTunnel(Tunnel tunnel) {
if (tunnel == null || !tunnel.markClosed()) {
return;
}
tunnels.remove(tunnel.session);
closeQuietly(tunnel.socket);
closeQuietly(tunnel.session);
}

private static void closeQuietly(Socket socket) {
if (socket != null) {
try {
socket.close();
} catch (IOException ignored) {
}
}
}

private static void closeQuietly(WebSocketSession session) {
if (session != null) {
try {
if (session.isOpen()) {
session.close(CloseStatus.NORMAL);
}
} catch (IOException ignored) {
}
}
}

// ==================== 协议嗅探与解析 ====================

/**
* 识别首包协议并解析目标地址。
* 先按 VLESS 判定(首字节 0 且 UUID 匹配),再按 Trojan 判定(SHA-224 哈希前缀匹配),
* 均不匹配返回 null。
*/
private static Target sniff(byte[] msg) {
if (msg.length > 17 && msg[0] == 0 && constantTimeEquals(UUID_BYTES, 0, msg, 1, 16)) {
return parseVless(msg);
}
if (msg.length >= 58 && constantTimeEquals(TROJAN_HASH, 0, msg, 0, 56)) {
return parseTrojan(msg);
}
return null;
}

/**
* 解析 VLESS 首包:version(1) + uuid(16) + addonsLen(1) + addons(n) + cmd(1) + port(2) + atyp(1) + addr。
* cmd 必须为 0x01(TCP CONNECT),任何越界或非法寻址返回 null。
*/
private static Target parseVless(byte[] msg) {
int i = (msg[17] & 0xFF) + 18;
if (i >= msg.length) {
return null;
}
int cmd = msg[i++] & 0xFF;
if (cmd != 0x01) {
return null;
}
if (i + 2 > msg.length) {
return null;
}
int port = ((msg[i] & 0xFF) << 8) | (msg[i + 1] & 0xFF);
i += 2;
if (i >= msg.length) {
return null;
}
int atyp = msg[i++] & 0xFF;
String host;
switch (atyp) {
case 1:
if (i + 4 > msg.length) {
return null;
}
host = joinIpv4(msg, i);
i += 4;
break;
case 2:
if (i >= msg.length) {
return null;
}
int len = msg[i] & 0xFF;
if (i + 1 + len > msg.length) {
return null;
}
host = new String(msg, i + 1, len, StandardCharsets.UTF_8);
i += 1 + len;
break;
case 3:
if (i + 16 > msg.length) {
return null;
}
host = joinIpv6(msg, i);
i += 16;
break;
default:
return null;
}
if (host.isEmpty()) {
return null;
}
return new Target(host, port, i, true);
}

/**
* 解析 Trojan 首包:sha224 哈希(56) + [CRLF] + cmd(1) + atyp(1) + addr + port(2) + [CRLF]。
* cmd 必须为 0x01(CONNECT),任何越界或非法寻址返回 null。
*/
private static Target parseTrojan(byte[] msg) {
int i = 56;
if (msg[i] == 0x0d && msg[i + 1] == 0x0a) {
i += 2;
}
if (i >= msg.length) {
return null;
}
int cmd = msg[i++] & 0xFF;
if (cmd != 0x01) {
return null;
}
if (i >= msg.length) {
return null;
}
int atyp = msg[i++] & 0xFF;
String host;
switch (atyp) {
case 1:
if (i + 4 > msg.length) {
return null;
}
host = joinIpv4(msg, i);
i += 4;
break;
case 3:
if (i >= msg.length) {
return null;
}
int len = msg[i] & 0xFF;
if (i + 1 + len > msg.length) {
return null;
}
host = new String(msg, i + 1, len, StandardCharsets.UTF_8);
i += 1 + len;
break;
case 4:
if (i + 16 > msg.length) {
return null;
}
host = joinIpv6(msg, i);
i += 16;
break;
default:
return null;
}
if (i + 2 > msg.length) {
return null;
}
int port = ((msg[i] & 0xFF) << 8) | (msg[i + 1] & 0xFF);
i += 2;
if (i < msg.length && msg[i] == 0x0d && i + 1 < msg.length && msg[i + 1] == 0x0a) {
i += 2;
}
if (host.isEmpty()) {
return null;
}
return new Target(host, port, i, false);
}

// ==================== 黑名单 ====================

private static boolean isBlockedDomain(String host) {
String h = host.toLowerCase();
for (String blocked : BLOCKED_DOMAINS) {
if (h.equals(blocked) || h.endsWith("." + blocked)) {
return true;
}
}
return false;
}

// ==================== 字节工具 ====================

private static String joinIpv4(byte[] msg, int i) {
return (msg[i] & 0xFF) + "." + (msg[i + 1] & 0xFF) + "."
+ (msg[i + 2] & 0xFF) + "." + (msg[i + 3] & 0xFF);
}

private static String joinIpv6(byte[] msg, int i) {
StringBuilder sb = new StringBuilder();
for (int j = 0; j < 16; j += 2) {
if (j > 0) {
sb.append(':');
}
sb.append(Integer.toHexString(((msg[i + j] & 0xFF) << 8) | (msg[i + j + 1] & 0xFF)));
}
return sb.toString();
}

private static byte[] toBytes(ByteBuffer buffer) {
ByteBuffer dup = buffer.duplicate();
byte[] out = new byte[dup.remaining()];
dup.get(out);
return out;
}

private static byte[] hexToBytes(String hex) {
int len = hex.length();
byte[] out = new byte[len / 2];
for (int i = 0; i < len; i += 2) {
out[i / 2] = (byte) ((Character.digit(hex.charAt(i), 16) << 4)
| Character.digit(hex.charAt(i + 1), 16));
}
return out;
}

private static String toHex(byte[] bytes) {
StringBuilder sb = new StringBuilder(bytes.length * 2);
for (byte b : bytes) {
sb.append(Character.forDigit((b >> 4) & 0xF, 16));
sb.append(Character.forDigit(b & 0xF, 16));
}
return sb.toString();
}

private static byte[] buildTrojanHash() {
try {
byte[] digest = MessageDigest.getInstance("SHA-224").digest(UUID.getBytes(StandardCharsets.UTF_8));
return toHex(digest).getBytes(StandardCharsets.ISO_8859_1);
} catch (NoSuchAlgorithmException e) {
throw new IllegalStateException("SHA-224 algorithm not available", e);
}
}

/**
* 恒定时间比较,等价 timingSafeEqual,防止基于时序的凭据侧信道。
*/
private static boolean constantTimeEquals(byte[] a, int offA, byte[] b, int offB, int len) {
int diff = 0;
for (int i = 0; i < len; i++) {
diff |= a[offA + i] ^ b[offB + i];
}
return diff == 0;
}

// ==================== 内部结构 ====================

/** 首包解析结果:目标地址 + 首包业务数据偏移 + 协议类型标记 */
private static final class Target {
final String host;
final int port;
final int payloadOffset;
/** VLESS 协议标记:隧道建立后需先发送 2 字节握手响应 */
final boolean vless;

Target(String host, int port, int payloadOffset, boolean vless) {
this.host = host;
this.port = port;
this.payloadOffset = payloadOffset;
this.vless = vless;
}
}

/**
* 单条隧道的双向转发状态:持有 session 与 socket,负责 WS 数据写 TCP 与 TCP 泵送回 WS。
* sendLock 保证 session.sendMessage 串行,closed 标记保证关闭只执行一次。
*/
private final class Tunnel {
final WebSocketSession session;
final Socket socket;
final Object sendLock = new Object();
private volatile boolean closed;

Tunnel(WebSocketSession session, Socket socket) {
this.session = session;
this.socket = socket;
}

/**
* WS 方向数据直写 TCP(在容器回调线程执行,阻塞即天然背压)。
*/
void write(ByteBuffer data) throws IOException {
OutputStream out = socket.getOutputStream();
if (data.hasArray()) {
out.write(data.array(), data.arrayOffset() + data.position(), data.remaining());
} else {
byte[] tmp = new byte[data.remaining()];
data.duplicate().get(tmp);
out.write(tmp);
}
out.flush();
}

/**
* TCP 数据泵送回 WS:持续读 socket 并发送二进制帧,EOF/异常后统一收口。
*/
void pump() {
byte[] buf = new byte[PUMP_BUFFER_SIZE];
try (InputStream in = socket.getInputStream()) {
int n;
while (!closed && (n = in.read(buf)) != -1) {
if (n == 0) {
continue;
}
byte[] chunk = Arrays.copyOf(buf, n);
synchronized (sendLock) {
if (!closed && session.isOpen()) {
session.sendMessage(new BinaryMessage(ByteBuffer.wrap(chunk)));
}
}
}
} catch (IOException e) {
log.debug("Proxy pump stream ended: {}", e.getMessage());
} finally {
closeTunnel(this);
}
}

boolean markClosed() {
synchronized (sendLock) {
if (closed) {
return false;
}
closed = true;
return true;
}
}

/**
* 触发隧道关闭,供线程池拒绝策略调用。
*/
void closeNow() {
closeTunnel(this);
}
}

/** 泵送任务包装,线程池拒绝时可直接关闭隧道 */
private static final class PumpTask implements Runnable {
private final Tunnel tunnel;

PumpTask(Tunnel tunnel) {
this.tunnel = tunnel;
}

@Override
public void run() {
tunnel.pump();
}

void closeNow() {
tunnel.closeNow();
}
}
}

11. 网关配置

前置 SLB / Nginx Ingress 需正确透传 WebSocket Upgrade:

Nginx 示例:

1
2
3
4
5
6
7
8
9
10
location /api/v1/proxy/websocket {
proxy_pass http://backend;
proxy_http_version 1.1;
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection "upgrade";
proxy_set_header Host $host;
proxy_read_timeout 600s;
proxy_send_timeout 600s;
proxy_buffering off;
}

K8s ingress-nginx 注解:

1
2
3
4
5
6
7
nginx.ingress.kubernetes.io/proxy-read-timeout: "600"
nginx.ingress.kubernetes.io/proxy-send-timeout: "600"
nginx.ingress.kubernetes.io/configuration-snippet: |
proxy_http_version 1.1;
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection "upgrade";
proxy_buffering off;

要点:

  • proxy_read_timeout / proxy_send_timeout 须大于服务端 maxSessionIdleTimeout(10min),建议 600s
  • proxy_buffering off 避免大帧被截断
  • 验证:握手响应 101 Switching Protocols、10 分钟空闲连接存活、512KB 帧传输无截断

12. 验证方案

  1. 编译mvn clean install 通过
  2. 回环 echo 自测nc -l 8080 作为回环目标;Node 测试脚本(ws 库)构造 VLESS 首包(目标 127.0.0.1:8080)连接 ws://127.0.0.1:8088/api/v1/proxy/websocket,发送后应收到 echo 回传二进制帧
  3. Trojan 同理:脚本计算 sha224(UUID) hex 前 56 字节构造首包,验证解析与回传
  4. 错误凭据:错误 UUID / 错误 Trojan 哈希 → 连接被立即关闭
  5. 黑名单:目标 speedtest.net:443 → 连接关闭
  6. 大帧:发送 512KB 二进制帧验证缓冲区配置生效(无 8KB 截断)
  7. 线程安全:并发 50 隧道压测,观察无 IllegalStateException、无泄漏线程

13. 风险与注意事项

风险 等级 说明与对策
安全与合规 凭据一旦泄露该端口即成为开放代理跳板。对策:UUID 上线前必须更换;高危端口(25/465/587 等)拒绝逻辑可选实施
Tomcat WS 线程池上限(200) 单实例隧道并发上限 = min(WS 线程池, 泵线程池 200)。对策:-Dorg.apache.tomcat.websocket.threadPoolSize=N 调整
网关空闲超时 proxy-read-timeout(常见 60s)会切断空闲隧道。对策:客户端开启心跳;服务端 maxSessionIdleTimeout 设 10min 并配合网关超时调优
APM 字节码增强 APM SDK 对 Socket 读写增强可能带来长连接转发开销。对策:上线后观测泵线程 CPU;必要时对代理类豁免
sendMessage 并发异常 全部发送路径 synchronized(sendLock),关闭与泵送并发安全
泵线程池耗尽 有界 200 + 拒绝即关隧道,防止线程爆炸拖垮整机
内存拷贝开销 每帧 BinaryMessage 一次拷贝,百级并发可接受

14. 设计要点总结

要点 说明
单 Controller 三合一 @Controller + @EnableWebSocket + extends BinaryWebSocketHandler + implements WebSocketConfigurer,lite mode 下 @Bean 直接生效
协议嗅探分流 首包先匹配 VLESS(首字节 0 + UUID),再匹配 Trojan(SHA-224 前缀),恒定时间比较防侧信道
天然背压 WS 回调线程阻塞式写 TCP,容器不再分发下一帧 → TCP 接收窗口收缩 → 客户端减速
有界线程池 + 拒绝即关 防止线程爆炸拖垮整机,拒绝策略直接关闭隧道而非 CallerRuns
sendLock 串行化 泵线程与回调 / 关闭并发调用 sendMessage 全部 synchronized(sendLock)
markClosed 幂等 保证关闭只执行一次,防止重复关闭异常
二进制缓冲区 1MB 避免默认 8KB 导致大帧被容器断开
空闲超时 10min 配合网关空闲超时调优,避免服务端先杀会话
订阅 ISP 缓存 双源探测 + 10min TTL + 后台线程 + 失败保留旧值