diff --git a/README.en.md b/README.en.md index f017116..85c2fe8 100644 --- a/README.en.md +++ b/README.en.md @@ -189,6 +189,8 @@ Such tools take over all traffic and often cut cloudflared's tunnel-edge connect | `client/` | "Phone access" settings tab + mobile adaptation (dsh-web-mobile port) | | `bin/dsh-pocket.mjs` | CLI: LAN/public modes, prints URL + QR | +When the browser cancels an HTTP request or disconnects, the proxy closes the corresponding upstream request, including while response headers are pending. Completing a request body normally keeps its pending response alive. + ## 🛠 Development ```sh diff --git a/README.md b/README.md index 6664f83..b1194af 100644 --- a/README.md +++ b/README.md @@ -191,6 +191,8 @@ npx @deepseek-ai/dsh web | `client/` | 设置页「手机访问」+ 移动端适配(dsh-web-mobile 移植) | | `bin/dsh-pocket.mjs` | CLI:局域网/公网模式,打印 URL + 二维码 | +浏览器取消 HTTP 请求或断开连接时,代理会关闭对应的上游请求,即使上游尚未返回响应头。请求体正常发送完成不会取消仍在等待的响应。 + ## 🛠 开发 ```sh diff --git a/lib/proxy.mjs b/lib/proxy.mjs index 740e13e..f8bab03 100644 --- a/lib/proxy.mjs +++ b/lib/proxy.mjs @@ -999,8 +999,16 @@ export function createPocketProxy({ port = 3081, host = '0.0.0.0', upstream = DE proxyRes.on('close', () => { if (!res.writableEnded) res.destroy(); }); }, ); + // 浏览器可在上游返回响应头之前取消搜索;请求一创建就接管断连清理。 + res.once('close', () => proxyReq.destroy()); + req.once('aborted', () => proxyReq.destroy()); proxyReq.on('error', (err) => { - if (!res.headersSent) res.writeHead(502, { 'content-type': 'text/plain; charset=utf-8' }); + if (res.destroyed || res.writableEnded) return; + if (res.headersSent) { + res.destroy(); + return; + } + res.writeHead(502, { 'content-type': 'text/plain; charset=utf-8' }); res.end(`dsh-pocket: 无法连接上游 dsh web(${upstream.host}:${upstream.port})——先启动 dsh web | ${err.message}`); }); req.pipe(proxyReq); diff --git a/test/proxy-cancellation.test.js b/test/proxy-cancellation.test.js new file mode 100644 index 0000000..ac44fda --- /dev/null +++ b/test/proxy-cancellation.test.js @@ -0,0 +1,126 @@ +import { once } from 'node:events'; +import { createServer, request } from 'node:http'; +import assert from 'node:assert/strict'; +import { test } from 'node:test'; + +import { createPocketProxy } from '../lib/proxy.mjs'; + +async function fixture(t, handler, options = {}) { + const upstream = createServer(handler); + t.after(async () => { + upstream.closeAllConnections(); + await new Promise((resolve) => upstream.close(resolve)); + }); + upstream.listen(0, '127.0.0.1'); + await once(upstream, 'listening'); + const proxy = await createPocketProxy({ + port: 0, + host: '127.0.0.1', + upstream: { host: '127.0.0.1', port: upstream.address().port }, + ...options, + }); + t.after(() => proxy.close()); + return proxy; +} + +function clientRequest(t, proxy, options = {}) { + const client = request({ host: '127.0.0.1', port: proxy.port, path: '/api/session/search', ...options }); + const closed = new Promise((resolve) => client.once('close', resolve)); + client.on('error', (error) => { + assert.equal(error.code, 'ECONNRESET'); + }); + t.after(async () => { + client.destroy(); + await closed; + }); + return { client, closed }; +} + +for (const method of ['GET', 'POST']) { + test(`cancelling ${method} before response headers closes the upstream request`, async (t) => { + const received = Promise.withResolvers(); + const proxy = await fixture(t, (req, res) => { + req.resume(); + req.once('end', () => received.resolve({ req, res })); + }); + const { client, closed } = clientRequest(t, proxy, { method }); + client.end(method === 'POST' ? JSON.stringify({ query: 'fixture' }) : undefined); + const upstream = await received.promise; + assert.equal(upstream.req.complete, true); + assert.equal(upstream.res.headersSent, false); + const cancelled = once(upstream.res, 'close', { signal: AbortSignal.timeout(2_000) }); + + client.destroy(); + await closed; + await cancelled; + + assert.equal(upstream.res.destroyed, true); + assert.equal(upstream.req.socket.destroyed, true); + }); +} + +test('cancelling an unfinished upload closes the upstream without an unhandled error', async (t) => { + const received = Promise.withResolvers(); + const proxy = await fixture(t, (req, res) => { + req.on('error', (error) => assert.equal(error.code, 'ECONNRESET')); + req.once('data', () => received.resolve({ req, res })); + }); + const { client, closed } = clientRequest(t, proxy, { method: 'POST', headers: { 'content-length': '1000' } }); + client.write('partial'); + const upstream = await received.promise; + assert.equal(upstream.req.complete, false); + const cancelled = once(upstream.res, 'close', { signal: AbortSignal.timeout(2_000) }); + + client.destroy(); + await closed; + await cancelled; + + assert.equal(upstream.req.socket.destroyed, true); +}); + +for (const scenario of [ + { name: 'buffered HTML', status: 200, type: 'text/html', body: 'pending' }, + { name: 'buffered 403', status: 403, type: 'text/plain', body: 'forbid' }, + { name: 'compressed JSON', status: 200, type: 'application/json', body: '{"pending":"' + 'fixture '.repeat(2000) }, +]) { + test(`cancelling a ${scenario.name} response closes the upstream`, async (t) => { + const headers = Promise.withResolvers(); + let upstream; + const proxy = await fixture(t, (req, res) => { + upstream = res; + res.writeHead(scenario.status, { 'content-type': scenario.type }); + res.write(scenario.body); + }, { log: () => headers.resolve() }); + const { client, closed } = clientRequest(t, proxy, { + headers: { accept: 'text/html', 'accept-encoding': 'gzip' }, + }); + client.end(); + await headers.promise; + const cancelled = once(upstream, 'close', { signal: AbortSignal.timeout(2_000) }); + + client.destroy(); + await closed; + await cancelled; + + assert.equal(upstream.destroyed, true); + }); +} + +test('a completed request body keeps waiting for its delayed response', async (t) => { + const received = Promise.withResolvers(); + const proxy = await fixture(t, (req, res) => { + req.resume(); + req.once('end', () => received.resolve(res)); + }); + const { client } = clientRequest(t, proxy, { method: 'POST' }); + const response = once(client, 'response'); + client.end('complete'); + const upstream = await received.promise; + upstream.writeHead(200, { 'content-type': 'application/json' }); + upstream.end('{"ok":true}'); + const [res] = await response; + const chunks = []; + for await (const chunk of res) chunks.push(chunk); + assert.equal(res.statusCode, 200); + assert.equal(Buffer.concat(chunks).toString('utf8'), '{"ok":true}'); +});