From 61499c6dc7ac327d9bf46367a79c781c598a0025 Mon Sep 17 00:00:00 2001 From: Nathan Simony <98903773+nathansimony@users.noreply.github.com> Date: Thu, 17 Sep 2026 11:13:28 +0200 Subject: [PATCH] Fix stalled connection attempts during Tuya handshakes --- index.js | 112 ++++++++++++++-------- test/connection.js | 229 +++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 302 insertions(+), 39 deletions(-) create mode 100644 test/connection.js diff --git a/index.js b/index.js index 5618d73..f80dc46 100644 --- a/index.js +++ b/index.js @@ -94,6 +94,7 @@ class TuyaDevice extends EventEmitter { this._responseTimeout = 2; // Seconds this._connectTimeout = 5; // Seconds + this._connectTimeoutTimer = null; this._pingPongPeriod = 10; // Seconds this._pingPongTimeout = null; this._lastPingAt = new Date(); @@ -548,10 +549,23 @@ class TuyaDevice extends EventEmitter { this.connectPromise.reject = rej; } + /** + * Reject and release a pending connection before notifying event listeners. + */ + _rejectConnect(error) { + if (this.connectPromise) { + const promise = this.connectPromise; + delete this.connectPromise; + promise.reject(error); + } + } + /** * Finish connecting and resolve */ _finishConnect() { + clearTimeout(this._connectTimeoutTimer); + this._connectTimeoutTimer = null; this._connected = true; /** @@ -615,13 +629,15 @@ class TuyaDevice extends EventEmitter { } this.createDeferredConnectPromise(); + const {connectPromise} = this; - this.client = new net.Socket(); + const client = new net.Socket(); + this.client = client; - // Default connect timeout is ~1 minute, - // 5 seconds is a more reasonable default - // since `retry` is used. - this.client.setTimeout(this._connectTimeout * 1000, () => { + // Bound the entire connection attempt, including session-key negotiation. + // A socket inactivity timeout can be disabled by TCP connect or postponed + // indefinitely by incoming data before the handshake completes. + this._connectTimeoutTimer = setTimeout(() => { /** * Emitted on socket error, usually a * result of a connection timeout. @@ -629,19 +645,24 @@ class TuyaDevice extends EventEmitter { * @event TuyaDevice#error * @property {Error} error error event */ - // this.emit('error', new Error('connection timed out')); - this.client.destroy(); - this.emit('error', new Error('connection timed out')); - if (this.connectPromise) { - this.connectPromise.reject(new Error('connection timed out')); - delete this.connectPromise; + if (this.client !== client) { + return; } - }); + + const error = new Error('connection timed out'); + this._rejectConnect(error); + this.disconnect(); + this.emit('error', error); + }, this._connectTimeout * 1000); // Add event listeners to socket // Parse response data - this.client.on('data', data => { + client.on('data', data => { + if (this.client !== client) { + return; + } + debug(`Received data: ${data.toString('hex')}`); let packets; @@ -682,31 +703,36 @@ class TuyaDevice extends EventEmitter { }); // Handle errors - this.client.on('error', err => { + client.on('error', err => { + if (this.client !== client) { + return; + } + debug('Error event from socket.', this.device.ip, err); + this._rejectConnect(err); + this.disconnect(); this.emit('error', new Error('Error from socket: ' + err.message)); - - if (!this._connected && this.connectPromise) { - this.connectPromise.reject(err); - delete this.connectPromise; - } - - this.client.destroy(); }); // Handle socket closure - this.client.on('close', () => { + client.on('close', () => { + // An error/disconnected listener may already have started a new attempt. + if (this.client !== client) { + return; + } + debug(`Socket closed: ${this.device.ip}`); this.disconnect(); }); - this.client.on('connect', async () => { - debug('Socket connected.'); + client.on('connect', () => { + if (this.client !== client || client.destroyed) { + return; + } - // Remove connect timeout - this.client.setTimeout(0); + debug('Socket connected.'); if (this.device.version === '3.4' || this.device.version === '3.5') { // Negotiate session key then emit 'connected' @@ -721,9 +747,12 @@ class TuyaDevice extends EventEmitter { }); debug('Protocol 3.4, 3.5: Negotiate Session Key - Send Msg 0x03'); - this.client.write(buffer); + client.write(buffer); } catch (error) { debug('Error binding key for protocol 3.4, 3.5: ' + error); + this._rejectConnect(error); + this.disconnect(); + this.emit('error', error); } return; @@ -733,9 +762,14 @@ class TuyaDevice extends EventEmitter { }); debug(`Connecting to ${this.device.ip}...`); - this.client.connect(this.device.port, this.device.ip); + try { + client.connect(this.device.port, this.device.ip); + } catch (error) { + this._rejectConnect(error); + this.disconnect(); + } - return this.connectPromise; + return connectPromise; } _packetHandler(packet) { @@ -759,11 +793,8 @@ class TuyaDevice extends EventEmitter { const expLocalHmac = packet.payload.slice(16, 16 + 32).toString('hex'); if (expLocalHmac !== calcLocalHmac) { const err = new Error(`HMAC mismatch(keys): expected ${expLocalHmac}, was ${calcLocalHmac}. ${packet.payload.toString('hex')}`); - if (this.connectPromise) { - this.connectPromise.reject(err); - delete this.connectPromise; - } - + this._rejectConnect(err); + this.disconnect(); this.emit('error', err); return; } @@ -931,19 +962,20 @@ class TuyaDevice extends EventEmitter { * close the socket and exit gracefully. */ disconnect() { - if (!this._connected) { - return; - } - debug('Disconnect'); + const wasConnected = this._connected; this._connected = false; this.device.parser.cipher.setSessionKey(null); // Clear timeouts + clearTimeout(this._connectTimeoutTimer); + this._connectTimeoutTimer = null; clearInterval(this._pingPongInterval); clearTimeout(this._pingPongTimeout); + this._rejectConnect(new Error('Connection closed before it was established')); + if (this.client) { this.client.destroy(); } @@ -956,7 +988,9 @@ class TuyaDevice extends EventEmitter { * goes off the network. * @event TuyaDevice#disconnected */ - this.emit('disconnected'); + if (wasConnected) { + this.emit('disconnected'); + } } /** diff --git a/test/connection.js b/test/connection.js new file mode 100644 index 0000000..05f6900 --- /dev/null +++ b/test/connection.js @@ -0,0 +1,229 @@ +import net from 'net'; +import crypto from 'crypto'; +import test from 'ava'; +import delay from 'delay'; + +const TuyAPI = require('..'); +const key = '0123456789abcdef'; + +// A minimal wire-protocol peer: complete the real session-key exchange without +// relying on TuyAPI's connection state or calling its private packet handlers. +function handshakeResponse(request, version, valid = true) { + let localNonce; + if (version === '3.4') { + const decipher = crypto.createDecipheriv('aes-128-ecb', key, null); + localNonce = Buffer.concat([decipher.update(request.slice(16, -36)), decipher.final()]); + } else { + const decipher = crypto.createDecipheriv('aes-128-gcm', key, request.slice(18, 30)); + decipher.setAAD(request.slice(4, 18)); + decipher.setAuthTag(request.slice(-20, -4)); + localNonce = Buffer.concat([decipher.update(request.slice(30, -20)), decipher.final()]); + } + + const hmac = crypto.createHmac('sha256', key).update(localNonce).digest(); + if (!valid) { + hmac[0] ^= 0xFF; + } + + const payload = Buffer.concat([Buffer.alloc(16, 0x42), hmac]); + if (version === '3.4') { + const cipher = crypto.createCipheriv('aes-128-ecb', key, null); + const encrypted = Buffer.concat([cipher.update(payload), cipher.final()]); + const header = Buffer.alloc(20); + header.writeUInt32BE(0x55AA, 0); + header.writeUInt32BE(1, 4); + header.writeUInt32BE(4, 8); + header.writeUInt32BE(encrypted.length + 40, 12); + const body = Buffer.concat([header, encrypted]); + return Buffer.concat([body, crypto.createHmac('sha256', key).update(body).digest(), Buffer.from('0000aa55', 'hex')]); + } + + const header = Buffer.alloc(18); + header.writeUInt32BE(0x6699, 0); + header.writeUInt32BE(1, 6); + header.writeUInt32BE(4, 10); + header.writeUInt32BE(payload.length + 32, 14); + const iv = Buffer.alloc(12, 0x24); + const cipher = crypto.createCipheriv('aes-128-gcm', key, iv); + cipher.setAAD(header.slice(4)); + const encrypted = Buffer.concat([cipher.update(Buffer.concat([Buffer.alloc(4), payload])), cipher.final()]); + return Buffer.concat([header, iv, encrypted, cipher.getAuthTag(), Buffer.from('00009966', 'hex')]); +} + +async function setup(t, version, firstConnection) { + const sockets = new Set(); + let attempts = 0; + let received; + const firstRequest = new Promise(resolve => { + received = resolve; + }); + const server = net.createServer(socket => { + sockets.add(socket); + socket.on('error', () => {}); + socket.on('close', () => sockets.delete(socket)); + const attempt = ++attempts; + socket.once('data', request => { + if (attempt === 1) { + received(); + firstConnection(socket, request); + } else { + socket.write(handshakeResponse(request, version)); + } + }); + }); + await new Promise(resolve => server.listen(0, '127.0.0.1', resolve)); + const device = new TuyAPI({id: 'test-device', key, ip: '127.0.0.1', port: server.address().port, version, issueGetOnConnect: false}); + const errors = []; + device.on('error', error => errors.push(error)); + t.context = {device, server, sockets}; + return {device, errors, firstRequest, attempts: () => attempts}; +} + +// Bound failures even against the unpatched implementation, whose connect() +// never settles in the cases covered below. +async function bounded(promise) { + let timer; + try { + return await Promise.race([promise, new Promise((resolve, reject) => { + timer = setTimeout(() => reject(new Error('Test deadline: connect() is still pending')), 1500); + })]); + } finally { + clearTimeout(timer); + } +} + +test.afterEach.always(async t => { + const {device, server, sockets} = t.context; + device.disconnect(); + if (device.client) { + device.client.destroy(); + } + + for (const socket of sockets) { + socket.destroy(); + } + + await new Promise(resolve => server.close(resolve)); +}); + +for (const version of ['3.4', '3.5']) { + test.serial(`${version}: close during handshake rejects and permits reconnection`, async t => { + const {device, attempts} = await setup(t, version, socket => socket.end()); + const connecting = device.connect(); + t.is(device.connect(), connecting); + await t.throwsAsync(bounded(connecting), {message: 'Connection closed before it was established'}); + t.false(device.isConnected()); + t.true(device.client.destroyed); + t.is(device.connectPromise, undefined); + t.true(await bounded(device.connect())); + t.true(device.isConnected()); + t.is(attempts(), 2); + }); + + test.serial(`${version}: silent handshake times out and a successful retry stays connected`, async t => { + const {device, errors} = await setup(t, version, () => {}); + device._connectTimeout = 0.15; + await t.throwsAsync(bounded(device.connect()), {message: 'connection timed out'}); + t.true(device.client.destroyed); + t.is(errors.length, 1); + t.true(await bounded(device.connect())); + await delay(250); + t.true(device.isConnected()); + t.is(errors.length, 1); + }); + + test.serial(`${version}: incoming traffic cannot postpone the handshake deadline`, async t => { + const {device, errors} = await setup(t, version, socket => { + const timer = setInterval(() => socket.write(Buffer.alloc(1)), 20); + socket.once('close', () => clearInterval(timer)); + }); + device._connectTimeout = 0.15; + await t.throwsAsync(bounded(device.connect()), {message: 'connection timed out'}); + t.true(errors.length > 1); + t.is(errors[errors.length - 1].message, 'connection timed out'); + t.true(await bounded(device.connect())); + t.true(device.isConnected()); + }); + + test.serial(`${version}: cancelling a handshake permits an immediate retry`, async t => { + const {device, firstRequest} = await setup(t, version, () => {}); + const rejection = t.throwsAsync(bounded(device.connect()), {message: 'Connection closed before it was established'}); + await firstRequest; + device.disconnect(); + const retry = device.connect(); + await rejection; + t.true(await bounded(retry)); + await delay(20); + t.true(device.isConnected()); + }); + + test.serial(`${version}: socket errors reject and release the pending attempt`, async t => { + const {device, errors, firstRequest} = await setup(t, version, () => {}); + const error = new Error('Socket failed during handshake'); + const rejection = t.throwsAsync(bounded(device.connect())); + await firstRequest; + device.client.destroy(error); + t.is(await rejection, error); + t.is(errors.length, 1); + t.true(await bounded(device.connect())); + t.true(device.isConnected()); + }); + + test.serial(`${version}: authentication failure releases the socket before retrying`, async t => { + const {device, errors} = await setup(t, version, (socket, request) => socket.write(handshakeResponse(request, version, false))); + await t.throwsAsync(bounded(device.connect()), {message: /HMAC mismatch/}); + t.true(device.client.destroyed); + t.is(errors.length, 1); + t.true(await bounded(device.connect())); + t.true(device.isConnected()); + }); + + test.serial(`${version}: an error listener can retry before the old socket closes`, async t => { + const {device} = await setup(t, version, () => {}); + device._connectTimeout = 0.15; + let retry; + device.once('error', () => { + retry = device.connect(); + }); + await t.throwsAsync(bounded(device.connect()), {message: 'connection timed out'}); + t.truthy(retry); + t.true(await bounded(retry)); + await delay(250); + t.true(device.isConnected()); + }); + + test.serial(`${version}: a disconnected listener can immediately establish a new session`, async t => { + const {device, attempts} = await setup(t, version, (socket, request) => socket.write(handshakeResponse(request, version))); + t.true(await bounded(device.connect())); + let retry; + device.once('disconnected', () => { + retry = device.connect(); + }); + device.disconnect(); + t.true(await bounded(retry)); + await delay(20); + t.true(device.isConnected()); + t.is(attempts(), 2); + }); + + test.serial(`${version}: handshake preparation errors reject immediately`, async t => { + const {device, errors} = await setup(t, version, () => {}); + const error = new Error('Cannot prepare handshake'); + device.device.parser.cipher.random = () => { + throw error; + }; + + t.is(await t.throwsAsync(bounded(device.connect())), error); + t.true(device.client.destroyed); + t.deepEqual(errors, [error]); + }); +} + +test.serial('invalid socket options reject without retaining a pending attempt', async t => { + const {device} = await setup(t, '3.5', () => {}); + device.device.port = -1; + await t.throwsAsync(device.connect(), {instanceOf: RangeError}); + t.is(device.connectPromise, undefined); + t.true(device.client.destroyed); + t.is(device._connectTimeoutTimer, null); +});