Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions README.en.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -191,6 +191,8 @@ npx @deepseek-ai/dsh web
| `client/` | 设置页「手机访问」+ 移动端适配(dsh-web-mobile 移植) |
| `bin/dsh-pocket.mjs` | CLI:局域网/公网模式,打印 URL + 二维码 |

浏览器取消 HTTP 请求或断开连接时,代理会关闭对应的上游请求,即使上游尚未返回响应头。请求体正常发送完成不会取消仍在等待的响应。

## 🛠 开发

```sh
Expand Down
10 changes: 9 additions & 1 deletion lib/proxy.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
126 changes: 126 additions & 0 deletions test/proxy-cancellation.test.js
Original file line number Diff line number Diff line change
@@ -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: '<html><head></head><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}');
});