Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
58 changes: 58 additions & 0 deletions js/core/src/__tests__/yas.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -514,6 +514,7 @@ import {
transfersFor,
type YasConnectionOptions,
type YasFontFace,
type YasRequestFailure,
type YasRelayRoute,
type YasTransport,
type YasTransferDescriptor,
Expand Down Expand Up @@ -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<unknown>) =>
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");
Expand Down
90 changes: 80 additions & 10 deletions js/core/src/yas/session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -326,6 +335,7 @@ export class YasConnection {
private readonly invalidationListeners = new Set<InvalidationListener>();
private readonly readyListeners = new Set<ReadyListener>();
private readonly catalogChangeListeners = new Set<CatalogChangeListener>();
private readonly requestFailureListeners = new Set<RequestFailureListener>();
private readonly familyLimitValidators = new Map<
number,
(limits: readonly YasExtension[]) => void
Expand Down Expand Up @@ -520,12 +530,16 @@ export class YasConnection {
sensitive?: boolean,
): Promise<Uint8Array> {
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),
Expand All @@ -547,12 +561,16 @@ export class YasConnection {
sensitive?: boolean,
): Promise<YasResultEnvelope> {
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),
Expand All @@ -571,12 +589,16 @@ export class YasConnection {
sensitive?: boolean,
): Promise<T> {
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),
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -822,13 +857,13 @@ export class YasConnection {
const requestId = this.allocateRequestId();
let rejectPromise!: (error: unknown) => void;
const promise = new Promise<T>((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 {
Expand Down Expand Up @@ -859,13 +894,13 @@ export class YasConnection {
const requestId = this.allocateRequestId();
let rejectPromise!: (error: unknown) => void;
const promise = new Promise<YasResultEnvelope>((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 {
Expand Down Expand Up @@ -1456,6 +1491,41 @@ export class YasConnection {
}
}

private rejectRequest<T>(
family: number,
kind: number,
error: unknown,
): Promise<T> {
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 {
Expand Down
Loading