Skip to content

Commit 7e361a6

Browse files
committed
fix(usbmux): handle decoder errors and preserve unconsumed bytes on connect
1 parent dfa37d4 commit 7e361a6

3 files changed

Lines changed: 88 additions & 30 deletions

File tree

src/lib/usbmux/index.ts

Lines changed: 12 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ import {BaseSocketService} from '../../base-socket-service.js';
66
import {getLogger} from '../logger.js';
77
import {type PairRecord, processPlistResponse} from '../pair-record/index.js';
88
import {type RawPairRecordResponse} from '../pair-record/pair-record.js';
9-
import {LengthBasedSplitter, parsePlist} from '../plist/index.js';
9+
import {parsePlist} from '../plist/index.js';
1010
import type {PlistDictionary} from '../types.js';
1111
import {type DecodedUsbmux, UsbmuxDecoder} from '../usbmux/usbmux-decoder.js';
1212
import {UsbmuxEncoder} from '../usbmux/usbmux-encoder.js';
@@ -39,7 +39,6 @@ const log = getLogger('Usbmux');
3939
export const USBMUXD_PORT = 27015;
4040
export const DEFAULT_USBMUXD_SOCKET = '/var/run/usbmuxd';
4141
export const DEFAULT_USBMUXD_HOST = '127.0.0.1';
42-
export const MAX_FRAME_SIZE = 100 * 1024 * 1024; // 1MB
4342

4443
// Result codes from usbmuxd
4544
export const USBMUX_RESULT = {
@@ -68,7 +67,6 @@ export interface SocketOptions {
6867
*/
6968
export class Usbmux extends BaseSocketService {
7069
private readonly _decoder: UsbmuxDecoder;
71-
private readonly _splitter: LengthBasedSplitter;
7270
private readonly _encoder: UsbmuxEncoder;
7371
private _tag: number;
7472
private readonly _responseCallbacks: Record<number, (data: DecodedUsbmux) => void>;
@@ -81,16 +79,10 @@ export class Usbmux extends BaseSocketService {
8179
super(socketClient);
8280

8381
this._decoder = new UsbmuxDecoder();
84-
this._splitter = new LengthBasedSplitter({
85-
readableStream: socketClient,
86-
littleEndian: true,
87-
maxFrameLength: MAX_FRAME_SIZE,
88-
lengthFieldOffset: 0,
89-
lengthFieldLength: 4,
90-
lengthAdjustment: 0,
91-
});
92-
9382
this._socketClient.pipe(this._decoder);
83+
this._decoder.on('error', (err: Error) => {
84+
log.error(`Usbmux decoder error: ${err.message}`);
85+
});
9486

9587
this._encoder = new UsbmuxEncoder();
9688
this._encoder.pipe(this._socketClient);
@@ -219,20 +211,17 @@ export class Usbmux extends BaseSocketService {
219211

220212
if (data.payload.Number === USBMUX_RESULT.OK) {
221213
// Detach constructor-owned consumers from the raw socket so the caller
222-
// gets full byte-stream ownership. Leaving the splitter/decoder attached
223-
// lets their unread buffers grow, eventually applying backpressure that
224-
// stalls long-lived forwards.
225-
//
226-
// Order matters: unpipe() must happen before removeAllListeners() /
227-
// shutdown(). Removing listeners first breaks Node's pipe cleanup, leaving
228-
// orphaned source listeners that continue pushing bytes into detached
229-
// streams and eventually stall the socket.
230-
this._socketClient.unpipe(this._splitter);
231-
this._splitter.unpipe(this._decoder);
214+
// gets full byte-stream ownership.
232215
this._socketClient.unpipe(this._decoder);
233216
this._encoder.unpipe(this._socketClient);
234-
this._splitter.shutdown();
235217
this._decoder.removeAllListeners('data');
218+
219+
// Hand any unconsumed bytes back to the raw socket stream so the caller
220+
// receives them immediately when attaching listeners.
221+
const pending = this._decoder.buffer;
222+
if (pending.length > 0) {
223+
this._socketClient.unshift(pending);
224+
}
236225
return this._socketClient;
237226
} else if (data.payload.Number === USBMUX_RESULT.CONNREFUSED) {
238227
throw new Error(`Connection was refused to port ${port}`);

src/lib/usbmux/usbmux-decoder.ts

Lines changed: 14 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,10 @@ export class UsbmuxDecoder extends Transform {
2424
super({objectMode: true});
2525
}
2626

27+
get buffer(): Buffer {
28+
return this._buffer;
29+
}
30+
2731
_transform(chunk: Buffer, encoding: BufferEncoding, callback: TransformCallback): void {
2832
// Append the new chunk to the internal buffer
2933
this._buffer = Buffer.concat([this._buffer, chunk]);
@@ -38,12 +42,16 @@ export class UsbmuxDecoder extends Transform {
3842
break; // Wait for more data
3943
}
4044

41-
// Extract the full message
42-
const message = this._buffer.slice(0, totalLength);
43-
this._decode(message);
45+
// Extract the full message and remove it from the buffer before decoding,
46+
// so any synchronous data listeners see only the unconsumed remainder.
47+
const message = this._buffer.subarray(0, totalLength);
48+
this._buffer = this._buffer.subarray(totalLength);
4449

45-
// Remove the processed message from the buffer
46-
this._buffer = this._buffer.slice(totalLength);
50+
try {
51+
this._decode(message);
52+
} catch (err) {
53+
return callback(err instanceof Error ? err : new Error(String(err)));
54+
}
4755
}
4856
callback();
4957
}
@@ -56,7 +64,7 @@ export class UsbmuxDecoder extends Transform {
5664
tag: data.readUInt32LE(12),
5765
};
5866

59-
const payload = data.slice(HEADER_LENGTH);
67+
const payload = data.subarray(HEADER_LENGTH);
6068
this.push({header, payload: parsePlist(payload)} as DecodedUsbmux);
6169
}
6270
}

test/unit/usbmux/usbmux.spec.ts

Lines changed: 62 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,11 +1,18 @@
11
import assert from 'node:assert/strict';
2-
import {type Server, type Socket} from 'node:net';
2+
import {type AddressInfo, type Server, type Socket, createConnection, createServer} from 'node:net';
3+
import {resolve} from 'node:path';
34
import {afterEach, beforeEach, describe, it} from 'node:test';
5+
import {fileURLToPath} from 'node:url';
6+
7+
import {fs, node} from '@appium/support';
48

59
import {type Device, Usbmux} from '../../../src/lib/usbmux/index.js';
10+
import {UsbmuxDecoder} from '../../../src/lib/usbmux/usbmux-decoder.js';
611
import {prioritizeUsbOverNetworkForDuplicateUdids} from '../../../src/lib/usbmux/utils.js';
712
import {UDID, fixtures, getServerWithFixtures} from '../fixtures/index.js';
813

14+
const PKG_ROOT = node.getModuleRootSync('appium-ios-remotexpc', fileURLToPath(import.meta.url));
15+
916
const DUP_UDID = 'aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa';
1017

1118
function mockUsbmuxDevice(
@@ -95,6 +102,60 @@ describe('usbmux', function () {
95102
}
96103
});
97104

105+
it('should preserve remainder bytes in decoder buffer when partial chunk arrives', function () {
106+
const decoder = new UsbmuxDecoder();
107+
const chunk = Buffer.from([0x05, 0x00, 0x00, 0x00]);
108+
decoder.write(chunk);
109+
assert.deepStrictEqual(decoder.buffer, chunk);
110+
});
111+
112+
it('should emit error on malformed payload instead of throwing synchronously in decoder', async function () {
113+
const decoder = new UsbmuxDecoder();
114+
const errorPromise = new Promise<Error>((resolve) => {
115+
decoder.once('error', resolve);
116+
});
117+
118+
const header = Buffer.alloc(16);
119+
header.writeUInt32LE(20, 0); // length: 20 bytes total (16 header + 4 payload)
120+
header.writeUInt32LE(1, 4); // version
121+
header.writeUInt32LE(8, 8); // type
122+
header.writeUInt32LE(1, 12); // tag
123+
const invalidPayload = Buffer.from('bad!');
124+
125+
decoder.write(Buffer.concat([header, invalidPayload]));
126+
const err = await errorPromise;
127+
assert.ok(err instanceof Error);
128+
});
129+
130+
it('should unshift unconsumed trailing bytes received alongside connect result', async function () {
131+
const connectFixture = Buffer.from(
132+
await fs.readFile(resolve(PKG_ROOT, 'test', 'unit', 'fixtures', 'usbmuxconnectmessage.bin')),
133+
);
134+
const trailingGreeting = Buffer.from('REMOTE_SERVICE_GREETING');
135+
136+
server = createServer((s: Socket) => {
137+
s.once('data', (data: Buffer) => {
138+
const clientTag = data.readUInt32LE(12);
139+
connectFixture.writeUInt32LE(clientTag, 12);
140+
s.write(Buffer.concat([connectFixture, trailingGreeting]));
141+
});
142+
});
143+
server.listen();
144+
const addr = server.address() as AddressInfo;
145+
socket = createConnection(addr.port);
146+
147+
usbmux = new Usbmux(socket);
148+
const connectedSocket = await usbmux.connect(1, 62078);
149+
150+
const received = await new Promise<Buffer>((resolve) => {
151+
connectedSocket.once('data', resolve);
152+
connectedSocket.resume();
153+
});
154+
155+
assert.deepStrictEqual(received, trailingGreeting);
156+
connectedSocket.destroy();
157+
});
158+
98159
it('should order duplicate UDIDs with USB before Network', function () {
99160
const net = mockUsbmuxDevice(2, DUP_UDID, 'Network');
100161
const usb = mockUsbmuxDevice(1, DUP_UDID, 'USB');

0 commit comments

Comments
 (0)