feat(adb): support delayed ack

This commit is contained in:
Simon Chan 2024-02-01 23:04:21 +08:00
parent 59d78dae20
commit ac6dc1e57c
No known key found for this signature in database
GPG key ID: A8B69F750B9BCEDD
9 changed files with 267 additions and 61 deletions

View file

@ -13,7 +13,7 @@ import {
ConsumableWritableStream,
WritableStream,
} from "@yume-chan/stream-extra";
import { EMPTY_UINT8_ARRAY } from "@yume-chan/struct";
import { EMPTY_UINT8_ARRAY, NumberFieldType } from "@yume-chan/struct";
import type { AdbIncomingSocketHandler, AdbSocket, Closeable } from "../adb.js";
import { decodeUtf8, encodeUtf8 } from "../utils/index.js";
@ -32,6 +32,13 @@ export interface AdbPacketDispatcherOptions {
*/
appendNullToServiceString: boolean;
maxPayloadSize: number;
/**
* The number of bytes the device can send before receiving an ack packet.
* Set to 0 or any negative value to disable delayed ack.
* Otherwise the value must be in the range of unsigned 32-bit integer.
*/
initialDelayedAckBytes: number;
/**
* Whether to preserve the connection open after the `AdbPacketDispatcher` is closed.
*/
@ -39,6 +46,11 @@ export interface AdbPacketDispatcherOptions {
debugSlowRead?: boolean | undefined;
}
interface SocketOpenResult {
remoteId: number;
availableWriteBytes: number;
}
/**
* The dispatcher is the "dumb" part of the connection handling logic.
*
@ -79,6 +91,10 @@ export class AdbPacketDispatcher implements Closeable {
options: AdbPacketDispatcherOptions,
) {
this.options = options;
// Don't allow negative values in dispatcher
if (this.options.initialDelayedAckBytes < 0) {
this.options.initialDelayedAckBytes = 0;
}
connection.readable
.pipeTo(
@ -169,15 +185,39 @@ export class AdbPacketDispatcher implements Closeable {
}
#handleOkay(packet: AdbPacketData) {
if (this.#initializers.resolve(packet.arg1, packet.arg0)) {
let ackBytes: number;
if (this.options.initialDelayedAckBytes !== 0) {
if (packet.payload.byteLength !== 4) {
throw new Error(
"Invalid OKAY packet. Payload size should be 4",
);
}
ackBytes = NumberFieldType.Uint32.deserialize(packet.payload, true);
} else {
if (packet.payload.byteLength !== 0) {
throw new Error(
"Invalid OKAY packet. Payload size should be 0",
);
}
ackBytes = Infinity;
}
if (
this.#initializers.resolve(packet.arg1, {
remoteId: packet.arg0,
availableWriteBytes: ackBytes,
} satisfies SocketOpenResult)
) {
// Device successfully created the socket
return;
}
const socket = this.#sockets.get(packet.arg1);
if (socket) {
// Device has received last `WRTE` to the socket
socket.ack();
// When delayed ack is enabled, device has received `ackBytes` from the socket.
// When delayed ack is disabled, device has received last `WRTE` packet from the socket,
// `ackBytes` is `Infinity` in this case.
socket.ack(ackBytes);
return;
}
@ -186,6 +226,18 @@ export class AdbPacketDispatcher implements Closeable {
void this.sendPacket(AdbCommand.Close, packet.arg1, packet.arg0);
}
#sendOkay(localId: number, remoteId: number, ackBytes: number) {
let payload: Uint8Array;
if (this.options.initialDelayedAckBytes !== 0) {
payload = new Uint8Array(4);
new DataView(payload.buffer).setUint32(0, ackBytes, true);
} else {
payload = EMPTY_UINT8_ARRAY;
}
return this.sendPacket(AdbCommand.Okay, localId, remoteId, payload);
}
async #handleOpen(packet: AdbPacketData) {
// `AsyncOperationManager` doesn't support skipping IDs
// Use `add` + `resolve` to simulate this behavior
@ -193,9 +245,20 @@ export class AdbPacketDispatcher implements Closeable {
this.#initializers.resolve(localId, undefined);
const remoteId = packet.arg0;
let service = decodeUtf8(packet.payload);
if (service.endsWith("\0")) {
service = service.substring(0, service.length - 1);
let initialDelayedAckBytes = packet.arg1;
const service = decodeUtf8(packet.payload);
if (this.options.initialDelayedAckBytes === 0) {
if (initialDelayedAckBytes !== 0) {
throw new Error("Invalid OPEN packet. arg1 should be 0");
}
initialDelayedAckBytes = Infinity;
} else {
if (initialDelayedAckBytes === 0) {
throw new Error(
"Invalid OPEN packet. arg1 should be greater than 0",
);
}
}
const handler = this.#incomingSocketHandlers.get(service);
@ -211,11 +274,16 @@ export class AdbPacketDispatcher implements Closeable {
localCreated: false,
service,
});
controller.ack(initialDelayedAckBytes);
try {
await handler(controller.socket);
this.#sockets.set(localId, controller);
await this.sendPacket(AdbCommand.Okay, localId, remoteId);
await this.#sendOkay(
localId,
remoteId,
this.options.initialDelayedAckBytes,
);
} catch (e) {
await this.sendPacket(AdbCommand.Close, 0, remoteId);
}
@ -238,10 +306,10 @@ export class AdbPacketDispatcher implements Closeable {
}),
(async () => {
await socket.enqueue(packet.payload);
await this.sendPacket(
AdbCommand.Okay,
await this.#sendOkay(
packet.arg1,
packet.arg0,
packet.payload.length,
);
handled = true;
})(),
@ -255,11 +323,17 @@ export class AdbPacketDispatcher implements Closeable {
service += "\0";
}
const [localId, initializer] = this.#initializers.add<number>();
await this.sendPacket(AdbCommand.Open, localId, 0, service);
const [localId, initializer] =
this.#initializers.add<SocketOpenResult>();
await this.sendPacket(
AdbCommand.Open,
localId,
this.options.initialDelayedAckBytes,
service,
);
// Fulfilled by `handleOk`
const remoteId = await initializer;
const { remoteId, availableWriteBytes } = await initializer;
const controller = new AdbDaemonSocketController({
dispatcher: this,
localId,
@ -267,6 +341,7 @@ export class AdbPacketDispatcher implements Closeable {
localCreated: true,
service,
});
controller.ack(availableWriteBytes);
this.#sockets.set(localId, controller);
return controller.socket;