From 0cffd60173a1f347862260062a8e107192419b7b Mon Sep 17 00:00:00 2001 From: Anso Date: Mon, 18 May 2026 22:08:06 -0400 Subject: [PATCH] fix(mesh): clean up activeStreams on src socket close/error (F-10) (#1105) The cross-node MeshService.openCrossNode handler attached src.on('close') and src.on('error') listeners that only called tcpStream.destroy() and never deleted the activeStreams Map entry. MeshTcpStreamLike.destroy() sends a tcp_close frame to the remote but does not synchronously emit the local close event; that fires only when the remote sends back a tcp_close ack via _dispatchClose (or when the entire tunnel tears down and the bridge force-emits close on every stream). Failed dials whose remote ack was lost (peer gone, network drop, or the F-9-era timeouts) therefore leaked records into activeStreams until the tunnel itself idle-closed, producing the monotonically-climbing activeStreamCount documented in the mesh E2E audit (F-10). Add a cleanupRecord() closure that clears the open timer and deletes the record idempotently, then route all four lifecycle handlers (tcpStream.on('error'), tcpStream.on('close'), src.on('close'), src.on('error')) through it. Map.delete is naturally idempotent so the double-delete that fires when both sides close normally is a free no-op. The same-node openSameNode path already had this shape via its existing teardown() closure; this brings the cross-node path to parity. Three new Vitest cases in mesh-service.test.ts use an EventEmitter- backed fake src so emit('close') actually delivers to the listener (the existing makeFakeSocket uses vi.fn for .on which records calls but never fires them). They cover src.close, src.error, and the idempotency case where both src and tcpStream emit close in sequence. --- backend/src/__tests__/mesh-service.test.ts | 102 +++++++++++++++++++++ backend/src/services/MeshService.ts | 20 ++-- 2 files changed, 116 insertions(+), 6 deletions(-) diff --git a/backend/src/__tests__/mesh-service.test.ts b/backend/src/__tests__/mesh-service.test.ts index 5fe56193..eab9896f 100644 --- a/backend/src/__tests__/mesh-service.test.ts +++ b/backend/src/__tests__/mesh-service.test.ts @@ -1089,6 +1089,108 @@ describe('MeshService.openCrossNode (BUG-4)', () => { vi.useRealTimers(); } }); + + // F-10 regression: src.on('close')/src.on('error') must delete the + // activeStreams entry, not just call tcpStream.destroy(). MeshTcpStream- + // Like.destroy() only sends a tcp_close frame; the handle's own 'close' + // event waits for the remote ack. If that ack never lands (peer gone, + // network drop), the record sits in activeStreams until tunnel idle- + // close. Use an EventEmitter-backed fake src so emit('close') actually + // delivers to the listener (the plain makeFakeSocket uses vi.fn for .on + // which records calls but never fires them). + function makeEmittingSocket(): EventEmitter & { + destroy: ReturnType; + end: ReturnType; + write: ReturnType; + } { + const ee = new EventEmitter() as EventEmitter & { + destroy: ReturnType; + end: ReturnType; + write: ReturnType; + }; + ee.destroy = vi.fn(); + ee.end = vi.fn(); + ee.write = vi.fn(); + return ee; + } + + it('src.on(close) deletes activeStreams entry even when tcpStream never emits close', async () => { + const svc = MeshService.getInstance(); + const target: MeshTarget = { + nodeId: 14, stack: 'audit-mesh-pilot', service: 'echo', + port: 9001, alias: 'echo.audit-mesh-pilot.sencho-pilot-test.sencho', + }; + const fakeStream = makeFakeStream(50); + vi.spyOn( + svc as unknown as { dialMeshTcpStream: (t: MeshTarget) => MeshTcpStreamLike | null }, + 'dialMeshTcpStream', + ).mockReturnValue(fakeStream); + + const fakeSrc = makeEmittingSocket(); + await (svc as unknown as { openCrossNode: (t: MeshTarget, s: unknown) => Promise }) + .openCrossNode(target, fakeSrc); + + const activeStreams = (svc as unknown as { activeStreams: Map }).activeStreams; + expect(activeStreams.size).toBe(1); + expect(activeStreams.has(50)).toBe(true); + + fakeSrc.emit('close'); + + expect(activeStreams.size).toBe(0); + expect(fakeStream.destroy).toHaveBeenCalled(); + }); + + it('src.on(error) deletes activeStreams entry symmetrically', async () => { + const svc = MeshService.getInstance(); + const target: MeshTarget = { + nodeId: 14, stack: 'audit-mesh-pilot', service: 'echo', + port: 9001, alias: 'echo.audit-mesh-pilot.sencho-pilot-test.sencho', + }; + const fakeStream = makeFakeStream(51); + vi.spyOn( + svc as unknown as { dialMeshTcpStream: (t: MeshTarget) => MeshTcpStreamLike | null }, + 'dialMeshTcpStream', + ).mockReturnValue(fakeStream); + + const fakeSrc = makeEmittingSocket(); + await (svc as unknown as { openCrossNode: (t: MeshTarget, s: unknown) => Promise }) + .openCrossNode(target, fakeSrc); + + const activeStreams = (svc as unknown as { activeStreams: Map }).activeStreams; + expect(activeStreams.size).toBe(1); + + fakeSrc.emit('error', new Error('connection reset by peer')); + + expect(activeStreams.size).toBe(0); + expect(fakeStream.destroy).toHaveBeenCalled(); + }); + + it('cleanup is idempotent when both src and tcpStream emit close', async () => { + const svc = MeshService.getInstance(); + const target: MeshTarget = { + nodeId: 14, stack: 'audit-mesh-pilot', service: 'echo', + port: 9001, alias: 'echo.audit-mesh-pilot.sencho-pilot-test.sencho', + }; + const fakeStream = makeFakeStream(52); + vi.spyOn( + svc as unknown as { dialMeshTcpStream: (t: MeshTarget) => MeshTcpStreamLike | null }, + 'dialMeshTcpStream', + ).mockReturnValue(fakeStream); + + const fakeSrc = makeEmittingSocket(); + await (svc as unknown as { openCrossNode: (t: MeshTarget, s: unknown) => Promise }) + .openCrossNode(target, fakeSrc); + + const activeStreams = (svc as unknown as { activeStreams: Map }).activeStreams; + expect(activeStreams.size).toBe(1); + + expect(() => { + fakeSrc.emit('close'); + fakeStream.emit('close'); + }).not.toThrow(); + + expect(activeStreams.size).toBe(0); + }); }); describe('MeshService.openCrossNode without reverseDialer', () => { diff --git a/backend/src/services/MeshService.ts b/backend/src/services/MeshService.ts index a47abdbe..d030199a 100644 --- a/backend/src/services/MeshService.ts +++ b/backend/src/services/MeshService.ts @@ -1888,6 +1888,16 @@ export class MeshService extends EventEmitter implements MeshForwarderHost { const clearOpenTimer = () => { if (openTimer) { clearTimeout(openTimer); openTimer = null; } }; + // Idempotent stream cleanup. F-10: src.on('close'/'error') used to + // only destroy tcpStream and never delete the activeStreams entry, + // so failed dials whose remote tcp_close ack was lost (or whose peer + // was already gone) leaked records until the whole tunnel idle- + // closed. Map.delete is naturally idempotent, so it's safe for both + // the src side and the tcpStream side to call this. + const cleanupRecord = () => { + clearOpenTimer(); + this.activeStreams.delete(record.streamId); + }; tcpStream.on('open', () => { clearOpenTimer(); @@ -1903,18 +1913,16 @@ export class MeshService extends EventEmitter implements MeshForwarderHost { try { src.write(chunk); } catch { /* ignore */ } }); tcpStream.on('error', (err: Error) => { - clearOpenTimer(); + cleanupRecord(); this.logActivity({ source: 'pilot', level: 'error', type: 'tunnel.fail', nodeId: target.nodeId, alias: target.alias, streamId: record.streamId, message: sanitizeForLog(err.message), }); - this.activeStreams.delete(record.streamId); try { src.destroy(); } catch { /* ignore */ } }); tcpStream.on('close', () => { - clearOpenTimer(); - this.activeStreams.delete(record.streamId); + cleanupRecord(); try { src.end(); } catch { /* ignore */ } }); src.on('data', (chunk: Buffer) => { @@ -1922,8 +1930,8 @@ export class MeshService extends EventEmitter implements MeshForwarderHost { tcpStream.write(chunk); }); src.on('end', () => tcpStream.end()); - src.on('close', () => { clearOpenTimer(); tcpStream.destroy(); }); - src.on('error', () => { clearOpenTimer(); tcpStream.destroy(); }); + src.on('close', () => { cleanupRecord(); try { tcpStream.destroy(); } catch { /* ignore */ } }); + src.on('error', () => { cleanupRecord(); try { tcpStream.destroy(); } catch { /* ignore */ } }); } private registerActiveStream(alias: string, streamId?: number): ActiveStreamRecord {