diff --git a/js/core/src/__tests__/yas.test.ts b/js/core/src/__tests__/yas.test.ts index 1735fa6b..a8b24058 100644 --- a/js/core/src/__tests__/yas.test.ts +++ b/js/core/src/__tests__/yas.test.ts @@ -514,6 +514,7 @@ import { transfersFor, type YasConnectionOptions, type YasFontFace, + type YasRequestFailure, type YasRelayRoute, type YasTransport, type YasTransferDescriptor, @@ -4584,6 +4585,63 @@ describe("YAS v1", () => { expect(abortObserved).toBe(true); }); + it("reports every rejected Request to failure listeners", async () => { + const { transport, connection } = await connected(); + const failures: YasRequestFailure[] = []; + connection.onRequestFailure((failure) => failures.push(failure)); + const rejection = (promise: Promise) => + promise.then( + () => undefined, + (error: unknown) => error, + ); + + const refused = rejection( + connection.request(YAS_FAMILY_CORE, YAS_CORE_PING, new Uint8Array()), + ); + const request = lastRequest(transport); + transport.push( + encodeYasFrame({ + family: request.family, + kind: request.kind, + class: YAS_CLASS_RESULT, + requestId: request.requestId, + payload: encodeResultPayload(YAS_STATUS_INVALID, new Uint8Array()), + }), + ); + const unadvertised = rejection( + connection.request(YAS_FAMILY_RELAY, 0x7fff), + ); + const lost = rejection( + connection.request(YAS_FAMILY_CORE, YAS_CORE_PING, new Uint8Array()), + ); + transport.setStatus("disconnected"); + + expect(failures).toEqual([ + { family: YAS_FAMILY_CORE, kind: YAS_CORE_PING, error: await refused }, + { family: YAS_FAMILY_RELAY, kind: 0x7fff, error: await unadvertised }, + { family: YAS_FAMILY_CORE, kind: YAS_CORE_PING, error: await lost }, + ]); + expect(failures[0]?.error).toMatchObject({ status: YAS_STATUS_INVALID }); + }); + + it("reports a send failure that fails the session once", async () => { + const { transport, connection } = await connected(); + const failures: YasRequestFailure[] = []; + connection.onRequestFailure((failure) => failures.push(failure)); + transport.send = () => { + transport.setStatus("error"); + throw new Error("edge write failed"); + }; + + const error = await connection + .request(YAS_FAMILY_CORE, YAS_CORE_PING, new Uint8Array()) + .catch((rejection: unknown) => rejection); + + expect(failures).toEqual([ + { family: YAS_FAMILY_CORE, kind: YAS_CORE_PING, error }, + ]); + }); + it("preserves a synchronous send failure through transport close status", async () => { const { transport, connection } = await connected(); const sendError = new Error("edge write failed"); diff --git a/js/core/src/yas/session.ts b/js/core/src/yas/session.ts index fa34dc74..d0ee29eb 100644 --- a/js/core/src/yas/session.ts +++ b/js/core/src/yas/session.ts @@ -149,6 +149,15 @@ export interface YasCatalogChange { type CatalogChangeListener = (change: YasCatalogChange) => void; +/** A Request whose promise rejected, reported before the caller sees it. */ +export interface YasRequestFailure { + family: number; + kind: number; + error: unknown; +} + +type RequestFailureListener = (failure: YasRequestFailure) => void; + interface PendingRequest { family: number; kind: number; @@ -326,6 +335,7 @@ export class YasConnection { private readonly invalidationListeners = new Set(); private readonly readyListeners = new Set(); private readonly catalogChangeListeners = new Set(); + private readonly requestFailureListeners = new Set(); private readonly familyLimitValidators = new Map< number, (limits: readonly YasExtension[]) => void @@ -520,12 +530,16 @@ export class YasConnection { sensitive?: boolean, ): Promise { if (!this.ready) - return Promise.reject( + return this.rejectRequest( + family, + kind, new YasDisconnectedError("YAS session is not ready"), ); this.family(family); if (!this.operationAdvertised(family, YAS_CLASS_REQUEST, kind)) - return Promise.reject( + return this.rejectRequest( + family, + kind, new YasResultError( YAS_STATUS_UNSUPPORTED, new Uint8Array(0), @@ -547,12 +561,16 @@ export class YasConnection { sensitive?: boolean, ): Promise { if (!this.ready) - return Promise.reject( + return this.rejectRequest( + family, + kind, new YasDisconnectedError("YAS session is not ready"), ); this.family(family); if (!this.operationAdvertised(family, YAS_CLASS_REQUEST, kind)) - return Promise.reject( + return this.rejectRequest( + family, + kind, new YasResultError( YAS_STATUS_UNSUPPORTED, new Uint8Array(0), @@ -571,12 +589,16 @@ export class YasConnection { sensitive?: boolean, ): Promise { if (!this.ready) - return Promise.reject( + return this.rejectRequest( + family, + kind, new YasDisconnectedError("YAS session is not ready"), ); this.family(family); if (!this.operationAdvertised(family, YAS_CLASS_REQUEST, kind)) - return Promise.reject( + return this.rejectRequest( + family, + kind, new YasResultError( YAS_STATUS_UNSUPPORTED, new Uint8Array(0), @@ -683,6 +705,19 @@ export class YasConnection { return () => this.readyListeners.delete(listener); } + /** + * Observe every Request whose promise rejects: a non-OK Result, a Result + * body that failed to decode, a Request the server does not advertise, a + * frame that could not be written, or a session lost with the Request in + * flight. Runs synchronously before the rejection reaches the caller. A + * non-OK status returned through {@link requestResult} resolves and is not + * reported. + */ + onRequestFailure(listener: RequestFailureListener): () => void { + this.requestFailureListeners.add(listener); + return () => this.requestFailureListeners.delete(listener); + } + /** Observe applied FAMILY_UPDATE and SESSION_INFO catalogue replacements. */ onCatalogChange(listener: CatalogChangeListener): () => void { this.catalogChangeListeners.add(listener); @@ -822,13 +857,13 @@ export class YasConnection { const requestId = this.allocateRequestId(); let rejectPromise!: (error: unknown) => void; const promise = new Promise((resolve, reject) => { - rejectPromise = reject; + rejectPromise = this.reportingRejection(family, kind, reject); this.pending.set(requestId, { family, kind, decode, resolve: resolve as (value: unknown) => void, - reject, + reject: rejectPromise, }); }); try { @@ -859,13 +894,13 @@ export class YasConnection { const requestId = this.allocateRequestId(); let rejectPromise!: (error: unknown) => void; const promise = new Promise((resolve, reject) => { - rejectPromise = reject; + rejectPromise = this.reportingRejection(family, kind, reject); this.pending.set(requestId, { family, kind, preserveResult: true, resolve: resolve as (value: unknown) => void, - reject, + reject: rejectPromise, }); }); try { @@ -1456,6 +1491,41 @@ export class YasConnection { } } + private rejectRequest( + family: number, + kind: number, + error: unknown, + ): Promise { + this.emitRequestFailure({ family, kind, error }); + return Promise.reject(error); + } + + private reportingRejection( + family: number, + kind: number, + reject: (error: unknown) => void, + ): (error: unknown) => void { + let rejected = false; + return (error) => { + // A send that fails the session rejects the Request there first; the + // send error that follows never reaches the caller. + if (rejected) return; + rejected = true; + this.emitRequestFailure({ family, kind, error }); + reject(error); + }; + } + + private emitRequestFailure(failure: YasRequestFailure): void { + for (const listener of [...this.requestFailureListeners]) { + try { + listener(failure); + } catch (error) { + this.reportListenerError("request failure", error); + } + } + } + private emitInvalidation(invalidation: YasInvalidation): void { for (const listener of [...this.invalidationListeners]) { try {