Skip to content

Commit 741ef1b

Browse files
authored
Merge pull request cmliu#1436 from cmliu/alpha2.1
fix: 更新版本号至2026-08-09并优化XHTTP传输链路和连接管理
2 parents 072a925 + 4833ca6 commit 741ef1b

2 files changed

Lines changed: 86 additions & 78 deletions

File tree

CHANGELOG

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,11 @@
1+
## [2.1.20260809201158] - 2026-08-09 20:11:58
2+
3+
### Change
4+
5+
- 优化 **XHTTP** 传输链路:TCP 数据改由 `request.body.pipeTo(socket.writable)` 与 `socket.readable.pipeTo(IdentityTransformStream)` 双向直通,不再经过 ReadableStream + 上行写入队列 + BYOB/Grain 逐块搬运,显著降低 CPU 占用。
6+
- 优化 **XHTTP** 连接管理:`forwardataTCP` 新增"仅建立连接"模式,建连并写入首包后直接返回 socket 交由 pipe 直通,同时保留直连失败后自动切换反代的兜底逻辑(gRPC/WS 调用方不受影响)。
7+
- 优化 **XHTTP UDP** 处理:UDP 分支独立拆出为专用处理函数,保留 Trojan UDP 反代与 DNS 转发逻辑。
8+
19
## [2.1.20260729235734] - 2026-07-29 23:57:34
210

311
### Change

_worker.js

Lines changed: 78 additions & 78 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
const Version = '2026-07-29 23:57:34';
1+
const Version = '2026-08-09 20:11:58';
22
let config_JSON, 缓存SOCKS5白名单 = null, 调试日志打印 = false;
33
let SOCKS5白名单 = ['*tapecontent.net', '*cloudatacdn.com', '*loadshare.org', '*cdn-centaurus.com', 'scholar.google.com'];
44
const Pages静态页面 = 'https://edt-pages.github.io';
@@ -552,36 +552,75 @@ async function 处理XHTTP请求(request, yourUUID, 反代上下文 = {}) {
552552
return new Response('UDP is not supported', { status: 400 });
553553
}
554554

555-
const remoteConnWrapper = { socket: null, connectingPromise: null, retryConnect: null, downlinkDrain: Promise.resolve() };
556-
let 当前写入Socket = null;
557-
let 远端写入器 = null;
558-
const 失效远端连接 = () => 失效TCP连接世代(remoteConnWrapper);
559555
const responseHeaders = new Headers({
560556
'Content-Type': 'application/octet-stream',
561557
'X-Accel-Buffering': 'no',
562558
'Cache-Control': 'no-store'
563559
});
564560

565-
const 释放远端写入器 = () => {
566-
if (远端写入器) {
567-
try { 远端写入器.releaseLock() } catch (e) { }
568-
远端写入器 = null;
569-
}
570-
当前写入Socket = null;
561+
// UDP 分支:拆到独立函数(保留原逻辑)
562+
if (首包.isUDP) return 处理XHTTPUDP请求(首包, reader, request, 反代上下文, responseHeaders);
563+
564+
// ================= TCP 分支:pipe 对接 =================
565+
try { reader.releaseLock() } catch (e) { }
566+
567+
const remoteConnWrapper = { socket: null, connectingPromise: null, retryConnect: null, downlinkDrain: Promise.resolve() };
568+
const abortController = new AbortController();
569+
let 已清理 = false;
570+
const 清理 = (reason) => {
571+
if (已清理) return;
572+
已清理 = true;
573+
try { abortController.abort(reason) } catch (e) { }
574+
失效TCP连接世代(remoteConnWrapper); // 关闭 socket + 世代 +1
571575
};
572576

573-
const 获取远端写入器 = () => {
574-
const socket = remoteConnWrapper.socket;
575-
if (!socket) return null;
576-
if (socket !== 当前写入Socket) {
577-
释放远端写入器();
578-
当前写入Socket = socket;
579-
远端写入器 = socket.writable.getWriter();
577+
// ⚠️ 关键:ws 参数必须传占位对象(新版 forwardataTCP 2370 行会访问 ws.readyState,传 null 会 TypeError 崩掉)
578+
// closeSocketQuietly 对占位对象安全(无 close 方法 → TypeError 被内部 catch 吞掉)
579+
const 占位WS = { readyState: WebSocket.OPEN };
580+
581+
let socket;
582+
try {
583+
socket = await forwardataTCP(首包.hostname, 首包.port, 首包.rawData, 占位WS, 首包.respHeader, remoteConnWrapper, yourUUID, request, 反代上下文, 首包.协议 === 'trojan', 首包.原始数据, true);
584+
} catch (err) {
585+
log(`[XHTTP-Pipe] 连接失败: ${err?.message || err}`);
586+
清理(err);
587+
return new Response('bad gateway', { status: 502 });
588+
}
589+
if (!socket) {
590+
清理(new Error('socket is null'));
591+
return new Response('bad gateway', { status: 502 });
592+
}
593+
594+
// 上行:请求体直接 pipe 进 socket(首包残余数据已由 forwardataTCP 写入)
595+
const 上行Promise = (async () => {
596+
await request.body.pipeTo(socket.writable, { signal: abortController.signal });
597+
})();
598+
599+
// 下行:优先使用 IdentityTransformStream(若运行时不支持则回退到 TransformStream)
600+
const 响应流 = typeof IdentityTransformStream !== 'undefined'
601+
? new IdentityTransformStream()
602+
: new TransformStream();
603+
const 下行Promise = (async () => {
604+
const writer = 响应流.writable.getWriter();
605+
try {
606+
if (有效数据长度(首包.respHeader) > 0) await writer.write(首包.respHeader);
607+
} catch (error) {
608+
try { await writer.abort(error) } catch (e) { }
609+
throw error;
610+
} finally {
611+
try { writer.releaseLock() } catch (e) { }
580612
}
581-
return 远端写入器;
582-
};
613+
await socket.readable.pipeTo(响应流.writable, { signal: abortController.signal });
614+
})();
615+
616+
void 上行Promise.catch(清理);
617+
void 下行Promise.then(() => 清理(), 清理);
618+
void Promise.allSettled([上行Promise, 下行Promise]);
583619

584-
let XHTTP上行写入队列 = null;
620+
return new Response(响应流.readable, { status: 200, headers: responseHeaders });
621+
}
622+
623+
function 处理XHTTPUDP请求(首包, reader, request, 反代上下文, responseHeaders) {
585624
const 木马UDP上下文 = { 缓存: new Uint8Array(0), 反代地址: 反代上下文.木马反代地址 };
586625
return new Response(new ReadableStream({
587626
async start(controller) {
@@ -612,81 +651,38 @@ async function 处理XHTTP请求(request, yourUUID, 反代上下文 = {}) {
612651
try { controller.close() } catch (e) { }
613652
}
614653
};
615-
616-
const 上行写入队列 = XHTTP上行写入队列 = 创建上行写入队列({
617-
获取写入器: 获取远端写入器,
618-
获取连接任务: () => remoteConnWrapper.connectingPromise,
619-
释放写入器: 释放远端写入器,
620-
重试连接: async () => {
621-
if (typeof remoteConnWrapper.retryConnect !== 'function') throw new Error('retry unavailable');
622-
await remoteConnWrapper.retryConnect();
623-
},
624-
关闭连接: () => {
625-
失效远端连接();
626-
closeSocketQuietly(xhttpBridge);
627-
},
628-
名称: 'XHTTP上行'
629-
});
630-
631-
const 写入远端 = async (payload, allowRetry = true) => {
632-
return 上行写入队列.写入并等待(payload, allowRetry);
633-
};
634-
635654
let 转发失败 = false;
636655
try {
637-
if (首包.isUDP) {
638-
if (首包.协议 === 'trojan') {
639-
木马UDP上下文.目标主机 = 首包.hostname;
640-
木马UDP上下文.目标端口 = 首包.port;
641-
if (木马UDP上下文.反代地址) await 转发木马UDP数据(首包.原始数据, xhttpBridge, 木马UDP上下文, request);
642-
}
643-
if (!(首包.协议 === 'trojan' && 木马UDP上下文.反代地址) && 首包.rawData?.byteLength) {
644-
if (首包.协议 === 'trojan') await 转发木马UDP数据(首包.rawData, xhttpBridge, 木马UDP上下文, request);
645-
else await forwardataudp(首包.rawData, xhttpBridge, udpRespHeader, request);
646-
udpRespHeader = null;
647-
}
648-
} else {
649-
await forwardataTCP(首包.hostname, 首包.port, 首包.rawData, xhttpBridge, 首包.respHeader, remoteConnWrapper, yourUUID, request, 反代上下文, 首包.协议 === 'trojan', 首包.原始数据);
656+
if (首包.协议 === 'trojan') {
657+
木马UDP上下文.目标主机 = 首包.hostname;
658+
木马UDP上下文.目标端口 = 首包.port;
659+
if (木马UDP上下文.反代地址) await 转发木马UDP数据(首包.原始数据, xhttpBridge, 木马UDP上下文, request);
660+
}
661+
if (!(首包.协议 === 'trojan' && 木马UDP上下文.反代地址) && 首包.rawData?.byteLength) {
662+
if (首包.协议 === 'trojan') await 转发木马UDP数据(首包.rawData, xhttpBridge, 木马UDP上下文, request);
663+
else await forwardataudp(首包.rawData, xhttpBridge, udpRespHeader, request);
664+
udpRespHeader = null;
650665
}
651-
652666
while (true) {
653667
const { done, value } = await reader.read();
654668
if (done) break;
655669
if (!value || value.byteLength === 0) continue;
656-
if (首包.isUDP) {
657-
if (首包.协议 === 'trojan') await 转发木马UDP数据(value, xhttpBridge, 木马UDP上下文, request);
658-
else await forwardataudp(value, xhttpBridge, udpRespHeader, request);
659-
udpRespHeader = null;
660-
} else {
661-
if (!(await 写入远端(value))) throw new Error('Remote socket is not ready');
662-
}
663-
}
664-
665-
if (!首包.isUDP) {
666-
await 上行写入队列.等待空();
667-
const writer = 获取远端写入器();
668-
if (writer) {
669-
try { await writer.close() } catch (e) { }
670-
}
670+
if (首包.协议 === 'trojan') await 转发木马UDP数据(value, xhttpBridge, 木马UDP上下文, request);
671+
else await forwardataudp(value, xhttpBridge, udpRespHeader, request);
672+
udpRespHeader = null;
671673
}
672674
} catch (err) {
673675
转发失败 = true;
674676
log(`[XHTTP转发] 处理失败: ${err?.message || err}`);
675677
closeSocketQuietly(xhttpBridge);
676678
} finally {
677-
const 保持木马UDP反代下行 = !转发失败 && 首包.isUDP && 首包.协议 === 'trojan' && 木马UDP上下文.反代地址 && 木马UDP上下文.反代Socket;
678-
上行写入队列.清空();
679-
if (转发失败) 失效远端连接();
680-
释放远端写入器();
679+
const 保持木马UDP反代下行 = !转发失败 && 首包.协议 === 'trojan' && 木马UDP上下文.反代地址 && 木马UDP上下文.反代Socket;
681680
if (!保持木马UDP反代下行) try { 木马UDP上下文.反代Socket?.close() } catch (e) { }
682681
try { reader.releaseLock() } catch (e) { }
683682
}
684683
},
685684
cancel() {
686-
XHTTP上行写入队列?.清空();
687-
失效远端连接();
688685
try { 木马UDP上下文.反代Socket?.close() } catch (e) { }
689-
释放远端写入器();
690686
try { reader.releaseLock() } catch (e) { }
691687
}
692688
}), { status: 200, headers: responseHeaders });
@@ -2079,7 +2075,7 @@ async function SSAEAD解密(cryptoKey, nonceCounter, ciphertext) {
20792075
return new Uint8Array(pt);
20802076
}
20812077

2082-
async function forwardataTCP(host, portNum, rawData, ws, respHeader, remoteConnWrapper, yourUUID, request = null, 反代上下文 = {}, 允许木马反代 = false, 木马反代首包数据 = null) {
2078+
async function forwardataTCP(host, portNum, rawData, ws, respHeader, remoteConnWrapper, yourUUID, request = null, 反代上下文 = {}, 允许木马反代 = false, 木马反代首包数据 = null, 仅建立连接 = false) {
20832079
const ctx反代IP = 反代上下文.反代IP || '';
20842080
const ctx代理类型 = 反代上下文.代理类型 !== undefined ? 反代上下文.代理类型 : null;
20852081
const ctx代理全局 = 反代上下文.代理全局 !== undefined ? 反代上下文.代理全局 : false;
@@ -2116,6 +2112,7 @@ async function forwardataTCP(host, portNum, rawData, ws, respHeader, remoteConnW
21162112
throw new Error('connection superseded or client closed');
21172113
}
21182114
remoteConnWrapper.socket = socket;
2115+
if (仅建立连接) return socket;
21192116
connectStreams(socket, ws, 取出响应头, retryFunc, 连接仍有效, remoteConnWrapper).catch(err => {
21202117
if (!连接仍有效()) return;
21212118
log(`[TCP下行] 处理失败: ${err?.message || err}`);
@@ -2345,6 +2342,7 @@ async function forwardataTCP(host, portNum, rawData, ws, respHeader, remoteConnW
23452342
log(`[TCP转发] 启用 SOCKS5/HTTP/HTTPS/TURN/SSTP 全局代理`);
23462343
try {
23472344
await connecttoPry();
2345+
if (仅建立连接) return remoteConnWrapper.socket;
23482346
} catch (err) {
23492347
log(`[TCP转发] SOCKS5/HTTP/HTTPS/TURN/SSTP 代理连接失败: ${err.message}`);
23502348
throw err;
@@ -2360,6 +2358,7 @@ async function forwardataTCP(host, portNum, rawData, ws, respHeader, remoteConnW
23602358
if (remoteConnWrapper.generation !== 直连世代 || remoteConnWrapper.socket !== initialSocket) return;
23612359
await connecttoPry();
23622360
});
2361+
if (仅建立连接) return initialSocket;
23632362
} catch (err) {
23642363
log(`[TCP转发] 直连 ${host}:${portNum} 失败: ${err.message}`);
23652364
if (remoteConnWrapper.generation !== 直连世代) throw err;
@@ -2369,6 +2368,7 @@ async function forwardataTCP(host, portNum, rawData, ws, respHeader, remoteConnW
23692368
}
23702369
if (ws.readyState !== WebSocket.OPEN) throw err;
23712370
await connecttoPry();
2371+
if (仅建立连接) return remoteConnWrapper.socket;
23722372
}
23732373
}
23742374
}

0 commit comments

Comments
 (0)