Add transport status events (#790)

* Transport status events

Add symbol docs
Emit transport status events
Transport test suite

* Review fixes

* Remove core dependency

* HTTP transport use AbortSignal, error handling in TransportNode

* Improve stream handling

* Update packages/transport-web-serial/src/transport.ts

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>

* Fix linting

---------

Co-authored-by: philon- <philon-@users.noreply.github.com>
Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
This commit is contained in:
Jeremy Gallant
2025-08-23 20:44:03 -04:00
committed by GitHub
co-authored by Copilot philon-
parent 449fb3ac36
commit d7e32e9b03
15 changed files with 1884 additions and 181 deletions
@@ -0,0 +1,229 @@
import { describe, vi, expect, it, beforeEach, afterEach, type MockInstance } from "vitest";
import { runTransportContract } from "../../../tests/utils/transportContract";
import { TransportHTTP } from "./transport";
let abortTimeoutSpy: MockInstance | undefined;
beforeEach(() => {
abortTimeoutSpy = vi.spyOn(
globalThis.AbortSignal as unknown as { timeout(ms: number): AbortSignal },
"timeout",
).mockImplementation((ms: number) => {
const ctrl = new AbortController();
const abort = () =>
ctrl.abort(new DOMException("Timeout reached", "TimeoutError"));
// Uses setTimeout so vi.useFakeTimers() can fast-forward it
setTimeout(abort, ms);
return ctrl.signal;
});
});
afterEach(() => {
abortTimeoutSpy?.mockRestore();
});
function stubFetch() {
const inbox: Uint8Array[] = [];
let lastWritten: ArrayBuffer | undefined;
let forceNextReadToHang = false;
let forceNextReadToReturn500 = false;
function makeAbortAwareHang(signal?: AbortSignal): Promise<Response> {
return new Promise((_, reject) => {
const abort = () => reject(new DOMException("Aborted", "AbortError"));
if (signal?.aborted) {
abort();
return;
}
if (signal) {
signal.addEventListener("abort", abort, { once: true });
}
});
}
const mockFetch = vi.fn(async (url: string, init?: RequestInit) => {
const method = (init?.method ?? "GET").toUpperCase();
if (url.includes("/api/v1/toradio") && method === "OPTIONS") {
return { ok: true, status: 204 } as Response;
}
if (url.includes("/api/v1/toradio") && method === "PUT") {
lastWritten = init?.body as ArrayBuffer;
return { ok: true, status: 200 } as Response;
}
if (url.includes("/api/v1/fromradio") && method === "GET") {
if (forceNextReadToHang) {
forceNextReadToHang = false;
return makeAbortAwareHang(init?.signal ?? undefined);
}
if (forceNextReadToReturn500) {
forceNextReadToReturn500 = false;
return {
ok: false,
status: 500,
arrayBuffer: async () => new ArrayBuffer(0),
} as Response;
}
const next = inbox.shift() ?? new Uint8Array();
return {
ok: true,
status: 200,
arrayBuffer: async () => next.buffer,
} as Response;
}
return { ok: true, status: 200 } as Response;
});
vi.stubGlobal("fetch", mockFetch);
return {
pushIncoming: (u8: Uint8Array) => inbox.push(u8),
assertLastWritten: (u8: Uint8Array) => {
const got = new Uint8Array(lastWritten || new ArrayBuffer(0));
expect(got).toEqual(u8);
},
forceReadErrorOnce: () => {
forceNextReadToReturn500 = true;
},
forceReadTimeoutOnce: () => {
forceNextReadToHang = true;
},
getMock: () => mockFetch,
cleanup: () => vi.unstubAllGlobals(),
};
}
async function tickNextTimer() {
try {
await vi.advanceTimersToNextTimerAsync();
} catch {
await new Promise((r) => setTimeout(r, 5));
}
}
describe("TransportHTTP (contract)", () => {
runTransportContract({
name: "TransportHTTP",
setup: () => {
vi.useFakeTimers();
},
teardown: () => {
vi.useRealTimers();
vi.restoreAllMocks();
vi.unstubAllGlobals();
},
create: async () => {
(globalThis as unknown as { __http: ReturnType<typeof stubFetch> }).__http = stubFetch();
const transport = await TransportHTTP.create("127.0.0.1:80", false);
await tickNextTimer();
return transport;
},
pushIncoming: async (bytes) => {
(globalThis as unknown as { __http: ReturnType<typeof stubFetch> }).__http.pushIncoming(bytes);
await tickNextTimer();
},
assertLastWritten: (bytes) => {
(globalThis as unknown as { __http: ReturnType<typeof stubFetch> }).__http.assertLastWritten(bytes);
},
triggerDisconnect: async () => {
(globalThis as unknown as { __http: ReturnType<typeof stubFetch> }).__http.forceReadErrorOnce();
await tickNextTimer();
},
});
});
describe("TransportHTTP (extras)", () => {
let httpStub: ReturnType<typeof stubFetch> | undefined;
beforeEach(() => {
vi.useFakeTimers();
});
afterEach(() => {
vi.useRealTimers();
vi.restoreAllMocks();
httpStub?.cleanup();
httpStub = undefined;
});
async function createTransport(): Promise<TransportHTTP> {
httpStub = stubFetch();
const transport = await TransportHTTP.create("127.0.0.1:80", false);
await tickNextTimer();
return transport;
}
async function advanceOnePoll() {
await tickNextTimer();
}
it("emits DeviceDisconnected with reason 'read-timeout' when GET /fromradio hangs", async () => {
const transport = await createTransport();
const reader = transport.fromDevice.getReader();
httpStub!.forceReadTimeoutOnce();
await tickNextTimer();
await vi.advanceTimersByTimeAsync(8000);
let sawReadTimeout = false;
for (let i = 0; i < 6; i++) {
const { value } = await reader.read();
if (value?.type === "status" && value.data.reason === "read-timeout") {
sawReadTimeout = true;
break;
}
}
expect(sawReadTimeout).toBe(true);
reader.releaseLock();
await transport.disconnect();
});
it("stops polling after disconnect()", async () => {
const transport = await createTransport();
const fetchMock = httpStub!.getMock();
const callsBeforeDisconnect = fetchMock.mock.calls.length;
await transport.disconnect();
await advanceOnePoll();
await vi.runOnlyPendingTimersAsync();
const callsAfterDisconnect = fetchMock.mock.calls.length;
expect(callsAfterDisconnect).toBe(callsBeforeDisconnect);
});
it("emits DeviceDisconnected with reason 'read-timeout' when GET /fromradio hangs", async () => {
const transport = await createTransport();
const reader = transport.fromDevice.getReader();
httpStub!.forceReadTimeoutOnce();
await vi.advanceTimersToNextTimerAsync();
await vi.advanceTimersByTimeAsync(8000);
await Promise.resolve();
await Promise.resolve();
let sawReadTimeout = false;
for (let i = 0; i < 6; i++) {
const { value } = await reader.read();
if (value?.type === "status" && value.data.reason === "read-timeout") {
sawReadTimeout = true;
break;
}
}
expect(sawReadTimeout).toBe(true);
reader.releaseLock();
await transport.disconnect();
});
});
+178 -48
View File
@@ -1,14 +1,49 @@
import type { Types } from "@meshtastic/core";
import { Types } from "@meshtastic/core";
const FETCH_INTERVAL_MS = 3000;
const READ_TIMEOUT_MS = 7000;
const WRITE_TIMEOUT_MS = 4000;
function toArrayBuffer(uint8array: Uint8Array): ArrayBuffer {
if (
uint8array.buffer instanceof ArrayBuffer &&
uint8array.byteOffset === 0 &&
uint8array.byteLength === uint8array.buffer.byteLength
) {
return uint8array.buffer;
}
return uint8array.slice().buffer;
}
/**
* Provides HTTP(S) transport for Meshtastic devices.
*
* Implements {@link Types.Transport} using the device's HTTP API.
* Polls `/api/v1/fromradio` for incoming packets and writes to `/api/v1/toradio`.
*/
export class TransportHTTP implements Types.Transport {
private _toDevice: WritableStream<Uint8Array>;
private _fromDevice: ReadableStream<Types.DeviceOutput>;
private fromDeviceController?: ReadableStreamDefaultController<Types.DeviceOutput>;
private url: string;
private receiveBatchRequests: boolean;
private fetchInterval: number;
private fetching: boolean;
private interval: ReturnType<typeof setInterval> | undefined;
private inflightReadController?: AbortController;
private lastStatus: Types.DeviceStatusEnum =
Types.DeviceStatusEnum.DeviceDisconnected;
private closingByUser = false;
/**
* Probe the device and return a connected HTTP transport.
*
* @param address Hostname or IP address (with optional port).
* @param tls Use HTTPS if true, HTTP otherwise.
*/
public static async create(
address: string,
tls?: boolean,
@@ -17,97 +52,192 @@ export class TransportHTTP implements Types.Transport {
await fetch(`${connectionUrl}/api/v1/toradio`, {
method: "OPTIONS",
});
await Promise.resolve();
return new TransportHTTP(connectionUrl);
}
/**
* Construct a new HTTP transport for the given device URL.
*
* @param url Base URL of the device (`http://host:port` or `https://host:port`).
*/
constructor(url: string) {
this.url = url;
this.receiveBatchRequests = false;
this.fetchInterval = 3000;
this.fetchInterval = FETCH_INTERVAL_MS;
this.fetching = false;
this._toDevice = new WritableStream<Uint8Array>({
write: async (chunk) => {
await this.writeToRadio(chunk);
try {
await this.writeToRadio(chunk);
} catch (error) {
if (!this.closingByUser) {
this.emitStatus(
Types.DeviceStatusEnum.DeviceDisconnected,
this.isTimeoutOrAbort(error) ? "write-timeout" : "write-error",
);
}
throw error;
}
},
});
let controller: ReadableStreamDefaultController<Types.DeviceOutput>;
this._fromDevice = new ReadableStream<Types.DeviceOutput>({
start: (ctrl) => {
controller = ctrl;
this.fromDeviceController = ctrl;
this.emitStatus(Types.DeviceStatusEnum.DeviceConnecting);
// Start polling immediately
void this.safePoll();
this.interval = setInterval(
() => void this.safePoll(),
this.fetchInterval,
);
},
cancel: () => {
if (this.interval) {
clearInterval(this.interval);
}
this.interval = undefined;
},
});
this.interval = setInterval(async () => {
if (this.fetching) {
// We still have the previous request open
return;
}
this.fetching = true;
try {
await this.readFromRadio(controller);
} catch {
// TODO: Emit disconnection events for certain types of errors
}
this.fetching = false;
}, this.fetchInterval);
}
private async readFromRadio(
controller: ReadableStreamDefaultController<Types.DeviceOutput>,
): Promise<void> {
/** Poll `/api/v1/fromradio` and enqueue incoming packets. */
private async readFromRadio(): Promise<void> {
let readBuffer = new ArrayBuffer(1);
while (readBuffer.byteLength > 0) {
const response = await fetch(
`${this.url}/api/v1/fromradio?all=${
this.receiveBatchRequests ? "true" : "false"
}`,
{
method: "GET",
headers: {
Accept: "application/x-protobuf",
const inflight = new AbortController();
this.inflightReadController = inflight;
const signal = AbortSignal.any([
inflight.signal,
AbortSignal.timeout(READ_TIMEOUT_MS),
]);
try {
const response = await fetch(
`${this.url}/api/v1/fromradio?all=${this.receiveBatchRequests ? "true" : "false"}`,
{
method: "GET",
headers: { Accept: "application/x-protobuf" },
signal,
},
},
);
);
if (!response.ok) {
throw new Error(
`fromradio ${response.status} ${response.statusText}`,
);
}
readBuffer = await response.arrayBuffer();
this.emitStatus(Types.DeviceStatusEnum.DeviceConnected);
if (readBuffer.byteLength > 0) {
controller.enqueue({
type: "packet",
data: new Uint8Array(readBuffer),
});
readBuffer = await response.arrayBuffer();
if (readBuffer.byteLength > 0) {
this.fromDeviceController?.enqueue({
type: "packet",
data: new Uint8Array(readBuffer),
});
}
} finally {
this.inflightReadController = undefined;
}
}
}
/** Write a protobuf-encoded request to `/api/v1/toradio`. */
private async writeToRadio(data: Uint8Array): Promise<void> {
await fetch(`${this.url}/api/v1/toradio`, {
method: "PUT",
headers: {
"Content-Type": "application/x-protobuf",
},
body: data,
});
try {
const response = await fetch(`${this.url}/api/v1/toradio`, {
method: "PUT",
headers: { "Content-Type": "application/x-protobuf" },
body: toArrayBuffer(data),
signal: AbortSignal.timeout(WRITE_TIMEOUT_MS),
});
if (!response.ok) {
throw new Error(`toradio ${response.status} ${response.statusText}`);
}
} catch (error) {
if (!this.closingByUser) {
this.emitStatus(
Types.DeviceStatusEnum.DeviceDisconnected,
this.isTimeoutOrAbort(error) ? "write-timeout" : "write-error",
);
}
throw error;
}
}
/** Writable stream of bytes to the device. */
get toDevice(): WritableStream<Uint8Array> {
return this._toDevice;
}
/** Readable stream of {@link Types.DeviceOutput} from the device. */
get fromDevice(): ReadableStream<Types.DeviceOutput> {
return this._fromDevice;
}
/**
* Stop polling and emit `DeviceDisconnected("user")`.
*/
disconnect(): Promise<void> {
this.fetching = false;
this.closingByUser = true;
if (this.interval) {
clearInterval(this.interval);
}
this.interval = undefined;
this.fetching = false;
try {
this.inflightReadController?.abort();
} catch {}
this.inflightReadController = undefined;
this.emitStatus(Types.DeviceStatusEnum.DeviceDisconnected, "user");
return Promise.resolve();
}
private emitStatus(next: Types.DeviceStatusEnum, reason?: string): void {
if (next === this.lastStatus) {
return;
}
this.lastStatus = next;
this.fromDeviceController?.enqueue({
type: "status",
data: { status: next, reason },
});
}
private isTimeoutOrAbort(err: unknown): boolean {
return (
(err instanceof DOMException &&
(err.name === "AbortError" || err.name === "TimeoutError")) ||
(err instanceof Error &&
(err.name === "AbortError" || err.name === "TimeoutError"))
);
}
private async safePoll(): Promise<void> {
if (this.fetching) {
return;
}
this.fetching = true;
try {
await this.readFromRadio();
} catch (error) {
if (!this.closingByUser) {
this.emitStatus(
Types.DeviceStatusEnum.DeviceDisconnected,
this.isTimeoutOrAbort(error) ? "read-timeout" : "read-error",
);
}
} finally {
this.fetching = false;
}
}
}