diff --git a/documentation/docs/call/from_outside_your_app.mdx b/documentation/docs/call/from_outside_your_app.mdx
index c186bf75e..601038f9a 100644
--- a/documentation/docs/call/from_outside_your_app.mdx
+++ b/documentation/docs/call/from_outside_your_app.mdx
@@ -95,39 +95,72 @@ const context = new ExternalContext({
## Calling readers reactively
-Using an `ExternalContext`, you can also call `reader` methods
-_reactively_: you receive a new answer every time the state changes.
-Each answer is a pair of a `response` and an `aborted`: the reader's
-response, or the error it answered with (a declared error, a denied
-`authorizer()`, a state that is not constructed yet). An error does
-not end the read; Reboot keeps the subscription open and yields again
-once the state changes, so you decide whether to keep waiting, to
-`break`, or to `raise aborted`. Failed connections are retried for
-you. For example:
+You can also call `reader` methods _reactively_: you receive a new
+result every time the state changes. Each result is either the reader's
+`response` or the error it raised (`aborted`): a declared error, a
+denied `authorizer()`, a state that does not exist yet. An error does
+not end the read: the subscription stays open and yields again once the
+state changes, so you decide whether to keep reading, stop, or raise the
+error yourself. Failed connections are retried for you.
+
+:::note
+
+The Node.js `ExternalContext` cannot read reactively yet. From Node.js,
+you can use a `WebContext` from `@reboot-dev/reboot-web` instead.
+Running it in Node.js is not supported beyond the case shown below. It
+needs the `fetch` and `WebSocket` globals (Node.js 22 has both by
+default; Node.js 20 needs `--experimental-websocket`).
+
+:::
+
+For example:
-
+(CODE:src=../../../tests/reboot/documentation/chat_room.py&lines=39-49) -->
+
```py
-fig = Fig.ref(fig_id)
-async for response, aborted in fig.reactively().get_position(context):
+async for response, aborted in chat_room.reactively().messages(context):
if aborted is not None:
- print(f"{fig_id}: {aborted}")
+ # The reader raised an error, e.g., the chat room does not
+ # exist yet. The read continues and yields again once the
+ # state changes.
+ print(f"Could not read messages: {aborted}")
continue
- print(f"{fig_id}: {response}")
+ assert response is not None
+ print(response.messages)
+ if "Hello, World!" in response.messages:
+ break
```
-
+
+
+
```ts
-// COMING SOON!
+const [responses] = await chatRoom.reactively().messages(context);
+for await (const { response, aborted } of responses) {
+ if (aborted !== undefined) {
+ // The reader raised an error, e.g., the chat room does not
+ // exist yet. The read continues and yields again once the
+ // state changes.
+ console.log(`Could not read messages: ${aborted.message}`);
+ continue;
+ }
+ console.log(response.messages);
+ if (response.messages.includes("Hello, World!")) {
+ break;
+ }
+}
```
+
+
diff --git a/documentation/docs/call/from_react.mdx b/documentation/docs/call/from_react.mdx
index 4a375e7b1..b608acb21 100644
--- a/documentation/docs/call/from_react.mdx
+++ b/documentation/docs/call/from_react.mdx
@@ -601,4 +601,13 @@ outage. These retries are _idempotent_, so this is perfectly safe to do.
An error is considered unretryable if it originates from the application. All
other errors are retried with an exponential backoff.
+A reader hook never stops retrying. If the connection drops, the hook
+reconnects with backoff; while it reconnects, `isLoading` is `true` and
+`response` and `aborted` keep their last values. If the reader raises an
+error (a declared error, a denied `authorizer()`, `StateNotConstructed`),
+the hook sets `aborted` to that error, clears `response`, and reconnects
+with backoff. When the reader later returns a response, the hook sets
+`response` and clears `aborted`. If the error means the read should not
+continue, stop rendering the hook.
+
diff --git a/reboot/examples/chat-room/frontend/reboot-non-react-web/src/main.ts b/reboot/examples/chat-room/frontend/reboot-non-react-web/src/main.ts
index 97e7cffde..b01d9ee87 100644
--- a/reboot/examples/chat-room/frontend/reboot-non-react-web/src/main.ts
+++ b/reboot/examples/chat-room/frontend/reboot-non-react-web/src/main.ts
@@ -1,5 +1,9 @@
+import type { ResponseOrAborted } from "@reboot-dev/reboot-web";
import { WebContext } from "@reboot-dev/reboot-web";
-import { ChatRoom } from "../../api/chat_room/v1/chat_room_rbt_web";
+import {
+ ChatRoom,
+ ChatRoomMessagesAborted,
+} from "../../api/chat_room/v1/chat_room_rbt_web";
const root = document.getElementById("messages");
const button = document.getElementById("button");
@@ -29,9 +33,18 @@ async function handleClick(element: HTMLInputElement) {
async function bindToElement(
element: HTMLElement,
- generator: AsyncGenerator
+ generator: AsyncGenerator<
+ ResponseOrAborted
+ >
) {
- for await (const response of generator) {
+ for await (const { response, aborted } of generator) {
+ if (aborted !== undefined) {
+ // The reader raised an error, e.g., the chat room has not been
+ // created yet. The read keeps going and yields again once the
+ // state changes.
+ element.innerHTML = `${aborted.message}
`;
+ continue;
+ }
element.innerHTML = `${response.messages
.map((msg: string) => `${msg}
`)
.join("")}`;
diff --git a/reboot/plugin/skills/upgrade/migrations/next/reactively-yields-response-or-aborted.md b/reboot/plugin/skills/upgrade/migrations/next/reactively-yields-response-or-aborted.md
index 38e2dd97d..754774b96 100644
--- a/reboot/plugin/skills/upgrade/migrations/next/reactively-yields-response-or-aborted.md
+++ b/reboot/plugin/skills/upgrade/migrations/next/reactively-yields-response-or-aborted.md
@@ -43,3 +43,43 @@ Also:
`response, aborted = await anext(subscription)`.
- A reader whose response type is empty yields `(None, None)` on
success, so check `aborted`, not `response`, for those.
+
+## Web `reactively()` reads yield `{ response, aborted }` objects
+
+The non-React web client's `Type.ref(id).reactively().(context)`
+generator (`@reboot-dev/reboot-web`) used to yield bare responses and
+to swallow every error the reader raised, silently reconnecting
+forever, so a `for await` over it simply went quiet on a declared
+error, a denied `authorizer()`, or `StateNotConstructed`. It now
+yields a `ResponseOrAborted` object for every result: `{ response }`
+for a response, `{ aborted }` (the method's `Aborted`)
+for an error. An error does not end the read: the generator keeps going
+and yields again once the state changes. It never throws.
+
+Find every use of such a read, matching only the opening parenthesis
+since a formatter may have split the call across lines:
+
+ grep -rn "\.reactively(" --include=*.ts --include=*.tsx --include=*.js
+
+Destructure the item in every `for await` over one. Before:
+
+```ts
+for await (const response of responses) {
+ render(response);
+}
+```
+
+After:
+
+```ts
+for await (const { response, aborted } of responses) {
+ if (aborted !== undefined) {
+ // Decide: `continue` to wait for the state to change, `break` to
+ // stop reading, or `throw aborted`.
+ continue;
+ }
+ render(response);
+}
+```
+
+Code that pulls items with `responses.next()` gets the object too.
diff --git a/reboot/templates/reboot_web.ts.j2 b/reboot/templates/reboot_web.ts.j2
index fa40f8213..2c822accb 100644
--- a/reboot/templates/reboot_web.ts.j2
+++ b/reboot/templates/reboot_web.ts.j2
@@ -445,7 +445,7 @@ class _Reactively {
partialRequest?: {{ client.proto.state_name }}.Partial{{ method.proto.name }}Request,
options?: { signal?: AbortSignal },
): Promise<[
- AsyncGenerator<{{ client.proto.state_name }}.{{ method.proto.name }}Response, void, unknown>,
+ AsyncGenerator, void, unknown>,
(newRequest: {{ client.proto.state_name }}.Partial{{ method.proto.name }}Request) => void
]> {
const request = {{ client.proto.state_name }}{{ method.proto.name }}RequestToProtobuf(partialRequest);
@@ -458,9 +458,14 @@ class _Reactively {
id: this.#id,
requestType: {{ method.input_type }},
responseType: {{ method.output_type }},
+ abortedType: {{ client.proto.state_name | to_camel }}{{ method.proto.name | to_camel }}Aborted,
request: request,
signal: options?.signal,
- bearerToken: context.bearerToken,
+ // Read the token from `context` on every attempt, so that a
+ // token set with `context.setBearerToken()` after the read
+ // started is sent on the next attempt.
+ bearerToken: async () => await context.bearerToken?.(),
+ onUnauthenticated: context.onUnauthenticated,
websockets: context.websockets,
}
);
@@ -470,10 +475,13 @@ class _Reactively {
setRequest(typedRequest);
};
- async function* typedGenerator(): AsyncGenerator<{{ client.proto.state_name }}.{{ method.proto.name }}Response, void, unknown> {
- for await (const response of generator) {
- const typedResponse = {{ client.proto.state_name }}{{ method.proto.name }}ResponseFromProtobufShape(response);
- yield typedResponse;
+ async function* typedGenerator(): AsyncGenerator, void, unknown> {
+ for await (const { response, aborted } of generator) {
+ if (aborted !== undefined) {
+ yield { aborted };
+ continue;
+ }
+ yield { response: {{ client.proto.state_name }}{{ method.proto.name }}ResponseFromProtobufShape(response) };
}
};
diff --git a/reboot/web/index.ts b/reboot/web/index.ts
index 512460398..d9b8bfe85 100644
--- a/reboot/web/index.ts
+++ b/reboot/web/index.ts
@@ -4,6 +4,7 @@ import {
Backoff,
Event,
Status,
+ StatusCode,
TRANSACTION_SHOULD_RETRY_REASONS_WITHOUT_BACKOFF,
assert,
check_bufbuild_protobuf_library,
@@ -86,9 +87,13 @@ export class Deferred {
}
}
-interface AbortedType {
- new (...args: any[]): A;
- fromStatus(status: Status): A;
+// The type of an `Aborted` class, e.g., `typeof GreeterGreetAborted`:
+// a constructor with a static `fromStatus`. Its type parameter,
+// `AbortedType`, is the type of the instances it makes, e.g.,
+// `GreeterGreetAborted`. Mirrors bufbuild's `MessageType`.
+interface AbortedMessageType {
+ new (...args: any[]): AbortedType;
+ fromStatus(status: Status): AbortedType;
}
// A hook the SPA's provider plugs into `httpCall` so a 401 from a
@@ -101,8 +106,7 @@ export type OnUnauthenticated = () => Promise;
export async function httpCall<
RequestType extends Message,
ResponseType extends Message,
- A extends Aborted,
- AT extends AbortedType
+ AbortedType extends Aborted
>({
url,
method,
@@ -121,7 +125,7 @@ export async function httpCall<
stateRef: string;
requestType: MessageType;
responseType: MessageType;
- abortedType: AT;
+ abortedType: AbortedMessageType;
request: RequestType;
idempotencyKey: string;
options?: {
@@ -153,7 +157,7 @@ export async function httpCall<
// `/__/oauth/refresh`. The retry happens at most once.
let didRefresh = false;
let response: Response | undefined;
- let aborted: A | undefined;
+ let aborted: AbortedType | undefined;
// The age of the transaction this call started, once an error has
// told us: the root transaction id of its first attempt. A retry
@@ -168,7 +172,7 @@ export async function httpCall<
const headers = await buildHeaders();
const result = await (async (): Promise<{
response?: Response;
- aborted?: A;
+ aborted?: AbortedType;
}> => {
const backoff = new Backoff();
// A `TransactionShouldRetry` may ask us to retry immediately,
@@ -324,9 +328,28 @@ export async function httpCall<
}
}
+// Reads `method` reactively. The returned generator yields an item for
+// every result the server sends, until the caller stops iterating or
+// `signal` aborts:
+//
+// - `{ response }` when the reader returned a response.
+// - `{ aborted }` when the reader raised an error, e.g., a declared
+// error, a denied authorization, or `StateNotConstructed`.
+//
+// An error does not end the generator. After yielding `{ aborted }` it
+// waits with backoff, reconnects, and yields the next result. The
+// caller decides whether to keep iterating, `break`, or `throw
+// aborted`. A dropped connection is retried with backoff and yields
+// nothing. When the error is `Unauthenticated` and `onUnauthenticated`
+// is set, the hook is called to renew the session before `{ aborted }`
+// is yielded; if it returns `true` the generator reconnects immediately
+// instead. If the hook does not renew the session, `{ aborted }` is
+// yielded and the hook is called again on the next attempt, after the
+// backoff.
export function reactively<
RequestType extends Message,
- ResponseType extends Message
+ ResponseType extends Message,
+ AbortedType extends Aborted
>({
url,
state,
@@ -334,9 +357,11 @@ export function reactively<
id,
requestType,
responseType,
+ abortedType,
request,
signal,
bearerToken,
+ onUnauthenticated,
websockets = false,
}: {
url: string;
@@ -345,12 +370,17 @@ export function reactively<
id: string;
requestType: MessageType;
responseType: MessageType;
+ // Typed by the instance it makes (rather than by a second type
+ // parameter for the class) so that `AbortedType` is inferred as the
+ // method's `Aborted` and not as the base class.
+ abortedType: AbortedMessageType;
request?: RequestType;
signal?: AbortSignal;
bearerToken?: () => Promise;
+ onUnauthenticated?: OnUnauthenticated;
websockets: boolean;
}): [
- AsyncGenerator,
+ AsyncGenerator, void, unknown>,
(newRequest: PartialMessage) => void
] {
if (request !== undefined) {
@@ -396,7 +426,11 @@ export function reactively<
}
};
- async function* responses(): AsyncGenerator {
+ async function* responses(): AsyncGenerator<
+ ResponseOrAborted,
+ void,
+ unknown
+ > {
// Wait for either the first request or an abort.
await Promise.race([
firstRequest.wait(),
@@ -411,9 +445,24 @@ export function reactively<
const backoff = new Backoff();
+ // True after `onUnauthenticated` returned `true`, until the attempt
+ // that follows delivers a response or raises an error. If that
+ // attempt is `Unauthenticated` too, the hook is not called again
+ // right away: the error is yielded, and the hook is called on the
+ // attempt after the backoff. This keeps a hook that returns `true`
+ // without renewing the session from reconnecting in a tight loop.
+ let didRefresh = false;
+
assert(request !== undefined);
while (signal === undefined || !signal.aborted) {
+ // `setRequest()` aborts `responsesAbortController` and then
+ // replaces it with a new controller. Keep this attempt's signal:
+ // after the stream ends, `attemptSignal.aborted` is true only if
+ // `setRequest()` ended it, in which case the loop reads the new
+ // request immediately instead of waiting on the backoff.
+ const attemptSignal = responsesAbortController.signal;
+
try {
// The reactive read path multiplexes many RPCs over one
// WebSocket and each call may carry a different bearer (a
@@ -436,25 +485,88 @@ export function reactively<
const queryResponses = reactiveReader({
endpoint: `${url}/__/reboot/rpc/${stateRef}`,
request: queryRequest,
- signal: responsesAbortController.signal,
+ signal: attemptSignal,
websockets,
});
for await (const queryResponse of queryResponses) {
backoff.reset({ log: `[Reboot] Call to \`${method}\` succeeded` });
+ didRefresh = false;
if (queryResponse.responseOrStatus.case === "response") {
const response = responseType.fromBinary(
queryResponse.responseOrStatus.value
);
- yield response;
+ yield { response };
}
}
+
+ if (attemptSignal.aborted) {
+ // `setRequest()` closed the stream; read the new request.
+ continue;
+ }
+
+ // The server closed the stream, e.g., because it is shutting
+ // down. Wait with backoff, then reconnect.
+ await backoff.wait({
+ log: `[Reboot] Reactive call to \`${method}\` ended; retrying with backoff ...`,
+ });
} catch (e) {
- if (signal === undefined || !signal.aborted) {
+ if (signal !== undefined && signal.aborted) {
+ return;
+ }
+
+ if (attemptSignal.aborted) {
+ // `setRequest()` closed the stream; read the new request.
+ continue;
+ }
+
+ if (e instanceof Status) {
+ // The reader raised an error. If it is `Unauthenticated`, the
+ // session may be stale: ask `onUnauthenticated` to renew it
+ // before yielding the error.
+ if (
+ e.code === StatusCode.UNAUTHENTICATED &&
+ onUnauthenticated !== undefined &&
+ !didRefresh
+ ) {
+ let refreshed = false;
+ try {
+ refreshed = await onUnauthenticated();
+ } catch {
+ // The hook failed, e.g., because the backend is
+ // restarting; yield the original error.
+ }
+ if (refreshed) {
+ // Reconnect right away; the next attempt reads the
+ // renewed bearer token.
+ didRefresh = true;
+ continue;
+ }
+ }
+
+ // The hook is called again on the next `Unauthenticated`
+ // attempt, which comes after the backoff below.
+ didRefresh = false;
+
+ // Yield the error. It does not end the read: the state may
+ // change so that the next read returns a response.
+ yield { aborted: abortedType.fromStatus(e) };
+
+ if (signal !== undefined && signal.aborted) {
+ return;
+ }
+
await backoff.wait({
- log: `[Reboot] Retrying call to \`${method}\` with backoff ...`,
+ log: `[Reboot] Reactive call to \`${method}\` failed with ${e.message}; retrying with backoff ...`,
});
+ continue;
}
+
+ // A transport failure, e.g., a disconnect or a server that is
+ // restarting: reconnect with backoff.
+ await backoff.wait({
+ log: `[Reboot] Retrying call to \`${method}\` with backoff ...`,
+ });
}
}
}
@@ -1220,7 +1332,9 @@ export class WebContext {
// retries if it resolves `true`. Wire it to whatever renews the
// session — e.g. a `refreshBearer(...)`-based function in a
// browser app, or a custom refresh flow elsewhere. Without it,
- // unary 401s surface immediately to the caller.
+ // unary 401s surface immediately to the caller. A reactive read
+ // calls it on each `Unauthenticated` attempt and reconnects right
+ // away if it resolves `true`.
onUnauthenticated?: OnUnauthenticated;
constructor(
diff --git a/tests/reboot/documentation/BUILD.bazel b/tests/reboot/documentation/BUILD.bazel
index 3362c6f8f..e468033b2 100644
--- a/tests/reboot/documentation/BUILD.bazel
+++ b/tests/reboot/documentation/BUILD.bazel
@@ -12,6 +12,7 @@ load(
"//reboot:rules.bzl",
"js_proto_library",
"js_reboot_library",
+ "js_reboot_web_library",
"py_reboot_library",
)
load("//reboot/nodejs:rules.bzl", "js_reboot_test")
@@ -256,6 +257,17 @@ js_reboot_library_from_zod(
zod = "chat_room_zod.ts",
)
+# The browser client for the same Zod `ChatRoom`, from the proto the
+# macro above generates. `chat_room.ts` reads the chat room reactively
+# through it, which the Node.js client cannot do.
+js_reboot_web_library(
+ name = "chat_room_zod_web_ts_reboot",
+ srcs = [":chat_room_zod_js_reboot_proto"],
+ proto = ":chat_room_zod_js_reboot_proto_file",
+ proto_deps = [":chat_room_zod_js_reboot_proto"],
+ deps = [":chat_room_zod_js_reboot_js_proto"],
+)
+
ts_project(
name = "test_chat_room_zod_ts",
srcs = [
@@ -600,3 +612,61 @@ ts_project(
"//:node_modules/zod",
],
)
+
+py_reboot_library_from_pydantic(
+ name = "chat_room_py_reboot",
+ py_deps = [":chat_room_pydantic_py"],
+ pydantic = ":chat_room_pydantic.py",
+)
+
+py_test(
+ name = "test_chat_room_py",
+ srcs = ["chat_room.py"],
+ main = "chat_room.py",
+ deps = [
+ ":chat_room_py_reboot",
+ ":chat_room_pydantic_py",
+ "//reboot/aio:applications_py",
+ "//reboot/aio:tests_py",
+ ],
+)
+
+ts_project(
+ name = "chat_room_ts",
+ srcs = [
+ "chat_room.ts",
+ ":package.json",
+ ],
+ declaration = True,
+ tsconfig = {
+ "compilerOptions": {
+ "declaration": True,
+ "module": "nodenext",
+ "moduleResolution": "nodenext",
+ "target": "es2020",
+ },
+ },
+ deps = [
+ ":chat_room_zod_js_reboot",
+ ":chat_room_zod_ts",
+ ":chat_room_zod_web_ts_reboot",
+ "//:node_modules/@reboot-dev/reboot",
+ "//:node_modules/@reboot-dev/reboot-api",
+ "//:node_modules/@reboot-dev/reboot-web",
+ "//:node_modules/@types/node",
+ "//:node_modules/zod",
+ ],
+)
+
+js_reboot_test(
+ name = "test_chat_room_ts",
+ data = [
+ ":chat_room_ts",
+ ],
+ entry_point = "chat_room.js",
+ # A reactive read over plain HTTP goes over a websocket, and the
+ # global `WebSocket` a browser provides is only behind this flag in
+ # Node.js 20.
+ env = {"NODE_OPTIONS": "--experimental-websocket"},
+ visibility = ["//visibility:public"],
+)
diff --git a/tests/reboot/documentation/chat_room.py b/tests/reboot/documentation/chat_room.py
new file mode 100644
index 000000000..ed5dc3696
--- /dev/null
+++ b/tests/reboot/documentation/chat_room.py
@@ -0,0 +1,85 @@
+import unittest
+from rbt.v1alpha1.errors_pb2 import StateNotConstructed
+from reboot.aio.applications import Application
+from reboot.aio.auth.authorizers import allow
+from reboot.aio.contexts import ReaderContext, WriterContext
+from reboot.aio.external import ExternalContext
+from reboot.aio.tests import Reboot
+from tests.reboot.documentation.chat_room_pydantic import (
+ MessagesResponse,
+ SendRequest,
+)
+from tests.reboot.documentation.chat_room_pydantic_rbt import ChatRoom
+
+
+class ChatRoomServicer(ChatRoom.Servicer):
+
+ def authorizer(self):
+ return allow()
+
+ async def messages(
+ self,
+ context: ReaderContext,
+ ) -> MessagesResponse:
+ return MessagesResponse(messages=self.state.messages or [])
+
+ async def send(
+ self,
+ context: WriterContext,
+ request: SendRequest,
+ ) -> None:
+ self.state.messages = [*(self.state.messages or []), request.message]
+
+
+async def print_messages(
+ chat_room: ChatRoom.WeakReference,
+ context: ExternalContext,
+) -> None:
+ # The docs snippet: from `async for` through `break`.
+ async for response, aborted in chat_room.reactively().messages(context):
+ if aborted is not None:
+ # The reader raised an error, e.g., the chat room does not
+ # exist yet. The read continues and yields again once the
+ # state changes.
+ print(f"Could not read messages: {aborted}")
+ continue
+ assert response is not None
+ print(response.messages)
+ if "Hello, World!" in response.messages:
+ break
+
+
+class ChatRoomTest(unittest.IsolatedAsyncioTestCase):
+ """Runs the reactive read that the "Calling readers reactively"
+ section of `documentation/docs/call/from_outside_your_app.mdx`
+ shows, so the snippet is real and tested. The `ChatRoom` is the
+ Pydantic API definition in `chat_room_pydantic.py`, which the
+ React docs show."""
+
+ async def asyncSetUp(self) -> None:
+ self.rbt = Reboot()
+ await self.rbt.start()
+ await self.rbt.up(Application(servicers=[ChatRoomServicer]))
+
+ async def asyncTearDown(self) -> None:
+ await self.rbt.stop()
+
+ async def test_read_messages_reactively(self) -> None:
+ context = self.rbt.create_external_context(name=self.id())
+ chat_room = ChatRoom.ref("reboot-chat-room")
+ await chat_room.send(context, message="Hello, World!")
+ await print_messages(chat_room, context)
+
+ async def test_read_before_the_chat_room_exists(self) -> None:
+ context = self.rbt.create_external_context(name=self.id())
+ chat_room = ChatRoom.ref("another-chat-room")
+ responses = chat_room.reactively().messages(context)
+ response, aborted = await anext(responses)
+ await responses.aclose()
+ self.assertIsNone(response)
+ assert aborted is not None
+ self.assertIsInstance(aborted.error, StateNotConstructed)
+
+
+if __name__ == '__main__':
+ unittest.main()
diff --git a/tests/reboot/documentation/chat_room.ts b/tests/reboot/documentation/chat_room.ts
new file mode 100644
index 000000000..e6c01068d
--- /dev/null
+++ b/tests/reboot/documentation/chat_room.ts
@@ -0,0 +1,94 @@
+import {
+ Application,
+ ReaderContext,
+ Reboot,
+ WriterContext,
+ allow,
+} from "@reboot-dev/reboot";
+import { WebContext } from "@reboot-dev/reboot-web";
+import { strict as assert } from "node:assert";
+import test from "node:test";
+import { ChatRoom as ChatRoomBackend } from "./chat_room_zod_rbt.js";
+import { ChatRoom } from "./chat_room_zod_rbt_web.js";
+
+// The servicer for the Zod `ChatRoom` in `chat_room_zod.ts`, which the
+// React docs show as the TypeScript API definition.
+class ChatRoomServicer extends ChatRoomBackend.Servicer {
+ authorizer() {
+ return allow();
+ }
+
+ async messages(
+ context: ReaderContext,
+ request: ChatRoomBackend.MessagesRequest
+ ): Promise {
+ return { messages: this.state.messages };
+ }
+
+ async send(
+ context: WriterContext,
+ request: ChatRoomBackend.SendRequest
+ ): Promise {
+ this.state.messages.push(request.message);
+ }
+}
+
+// Runs the reactive read that the "Calling readers reactively" section
+// of `documentation/docs/call/from_outside_your_app.mdx` shows, so the
+// snippet is real and tested. Reactive reads from TypeScript use the
+// web client (`@reboot-dev/reboot-web`), which works from Node.js too.
+test("chat room", async (t) => {
+ let rbt: Reboot;
+ t.before(async () => {
+ rbt = new Reboot();
+ await rbt.start();
+ // The web client connects to a URL, which the local Envoy provides.
+ await rbt.up(new Application({ servicers: [ChatRoomServicer] }), {
+ localEnvoy: true,
+ });
+ });
+ t.after(async () => {
+ await rbt.stop();
+ });
+
+ await t.test("read messages reactively", async () => {
+ const context = new WebContext({ url: rbt.url() });
+ const chatRoom = ChatRoom.ref("reboot-chat-room");
+ await chatRoom.send(context, { message: "Hello, World!" });
+
+ // The docs snippet: from `const [responses]` through the loop.
+ const [responses] = await chatRoom.reactively().messages(context);
+ for await (const { response, aborted } of responses) {
+ if (aborted !== undefined) {
+ // The reader raised an error, e.g., the chat room does not
+ // exist yet. The read continues and yields again once the
+ // state changes.
+ console.log(`Could not read messages: ${aborted.message}`);
+ continue;
+ }
+ console.log(response.messages);
+ if (response.messages.includes("Hello, World!")) {
+ break;
+ }
+ }
+ });
+
+ await t.test("read before the chat room exists", async () => {
+ const context = new WebContext({ url: rbt.url() });
+ const abortController = new AbortController();
+ const [responses] = await ChatRoom.ref("another-chat-room")
+ .reactively()
+ .messages(context, {}, { signal: abortController.signal });
+ const result = await responses.next();
+ if (result.done === true) {
+ throw new Error("Expected an item");
+ }
+ const { response, aborted } = result.value;
+ assert(response === undefined);
+ assert(aborted !== undefined);
+ // A Zod API's errors are plain objects discriminated by `type`.
+ assert(aborted.error.type === "StateNotConstructed");
+ abortController.abort();
+ assert((await responses.next()).done);
+ });
+});
diff --git a/tests/reboot/greeter_rbt_web.golden.js b/tests/reboot/greeter_rbt_web.golden.js
index 554dc3755..c02647549 100755
--- a/tests/reboot/greeter_rbt_web.golden.js
+++ b/tests/reboot/greeter_rbt_web.golden.js
@@ -1664,9 +1664,14 @@ class _Reactively {
id: __classPrivateFieldGet(this, __Reactively_id, "f"),
requestType: greeter_pb.GreetRequest,
responseType: greeter_pb.GreetResponse,
+ abortedType: GreeterGreetAborted,
request: request,
signal: options === null || options === void 0 ? void 0 : options.signal,
- bearerToken: context.bearerToken,
+ // Read the token from `context` on every attempt, so that a
+ // token set with `context.setBearerToken()` after the read
+ // started is sent on the next attempt.
+ bearerToken: async () => { var _a; return await ((_a = context.bearerToken) === null || _a === void 0 ? void 0 : _a.call(context)); },
+ onUnauthenticated: context.onUnauthenticated,
websockets: context.websockets,
});
const setTypedRequest = (newRequest) => {
@@ -1674,9 +1679,12 @@ class _Reactively {
setRequest(typedRequest);
};
async function* typedGenerator() {
- for await (const response of generator) {
- const typedResponse = GreeterGreetResponseFromProtobufShape(response);
- yield typedResponse;
+ for await (const { response, aborted } of generator) {
+ if (aborted !== undefined) {
+ yield { aborted };
+ continue;
+ }
+ yield { response: GreeterGreetResponseFromProtobufShape(response) };
}
}
;
@@ -1691,9 +1699,14 @@ class _Reactively {
id: __classPrivateFieldGet(this, __Reactively_id, "f"),
requestType: Empty,
responseType: Empty,
+ abortedType: GreeterTryToConstructContextAborted,
request: request,
signal: options === null || options === void 0 ? void 0 : options.signal,
- bearerToken: context.bearerToken,
+ // Read the token from `context` on every attempt, so that a
+ // token set with `context.setBearerToken()` after the read
+ // started is sent on the next attempt.
+ bearerToken: async () => { var _a; return await ((_a = context.bearerToken) === null || _a === void 0 ? void 0 : _a.call(context)); },
+ onUnauthenticated: context.onUnauthenticated,
websockets: context.websockets,
});
const setTypedRequest = (newRequest) => {
@@ -1701,9 +1714,12 @@ class _Reactively {
setRequest(typedRequest);
};
async function* typedGenerator() {
- for await (const response of generator) {
- const typedResponse = GreeterTryToConstructContextResponseFromProtobufShape(response);
- yield typedResponse;
+ for await (const { response, aborted } of generator) {
+ if (aborted !== undefined) {
+ yield { aborted };
+ continue;
+ }
+ yield { response: GreeterTryToConstructContextResponseFromProtobufShape(response) };
}
}
;
@@ -1718,9 +1734,14 @@ class _Reactively {
id: __classPrivateFieldGet(this, __Reactively_id, "f"),
requestType: Empty,
responseType: Empty,
+ abortedType: GreeterTryToConstructExternalContextAborted,
request: request,
signal: options === null || options === void 0 ? void 0 : options.signal,
- bearerToken: context.bearerToken,
+ // Read the token from `context` on every attempt, so that a
+ // token set with `context.setBearerToken()` after the read
+ // started is sent on the next attempt.
+ bearerToken: async () => { var _a; return await ((_a = context.bearerToken) === null || _a === void 0 ? void 0 : _a.call(context)); },
+ onUnauthenticated: context.onUnauthenticated,
websockets: context.websockets,
});
const setTypedRequest = (newRequest) => {
@@ -1728,9 +1749,12 @@ class _Reactively {
setRequest(typedRequest);
};
async function* typedGenerator() {
- for await (const response of generator) {
- const typedResponse = GreeterTryToConstructExternalContextResponseFromProtobufShape(response);
- yield typedResponse;
+ for await (const { response, aborted } of generator) {
+ if (aborted !== undefined) {
+ yield { aborted };
+ continue;
+ }
+ yield { response: GreeterTryToConstructExternalContextResponseFromProtobufShape(response) };
}
}
;
@@ -1745,9 +1769,14 @@ class _Reactively {
id: __classPrivateFieldGet(this, __Reactively_id, "f"),
requestType: greeter_pb.TestLongRunningFetchRequest,
responseType: Empty,
+ abortedType: GreeterTestLongRunningFetchAborted,
request: request,
signal: options === null || options === void 0 ? void 0 : options.signal,
- bearerToken: context.bearerToken,
+ // Read the token from `context` on every attempt, so that a
+ // token set with `context.setBearerToken()` after the read
+ // started is sent on the next attempt.
+ bearerToken: async () => { var _a; return await ((_a = context.bearerToken) === null || _a === void 0 ? void 0 : _a.call(context)); },
+ onUnauthenticated: context.onUnauthenticated,
websockets: context.websockets,
});
const setTypedRequest = (newRequest) => {
@@ -1755,9 +1784,12 @@ class _Reactively {
setRequest(typedRequest);
};
async function* typedGenerator() {
- for await (const response of generator) {
- const typedResponse = GreeterTestLongRunningFetchResponseFromProtobufShape(response);
- yield typedResponse;
+ for await (const { response, aborted } of generator) {
+ if (aborted !== undefined) {
+ yield { aborted };
+ continue;
+ }
+ yield { response: GreeterTestLongRunningFetchResponseFromProtobufShape(response) };
}
}
;
@@ -1772,9 +1804,14 @@ class _Reactively {
id: __classPrivateFieldGet(this, __Reactively_id, "f"),
requestType: greeter_pb.GetWholeStateRequest,
responseType: GreeterProto,
+ abortedType: GreeterGetWholeStateAborted,
request: request,
signal: options === null || options === void 0 ? void 0 : options.signal,
- bearerToken: context.bearerToken,
+ // Read the token from `context` on every attempt, so that a
+ // token set with `context.setBearerToken()` after the read
+ // started is sent on the next attempt.
+ bearerToken: async () => { var _a; return await ((_a = context.bearerToken) === null || _a === void 0 ? void 0 : _a.call(context)); },
+ onUnauthenticated: context.onUnauthenticated,
websockets: context.websockets,
});
const setTypedRequest = (newRequest) => {
@@ -1782,9 +1819,12 @@ class _Reactively {
setRequest(typedRequest);
};
async function* typedGenerator() {
- for await (const response of generator) {
- const typedResponse = GreeterGetWholeStateResponseFromProtobufShape(response);
- yield typedResponse;
+ for await (const { response, aborted } of generator) {
+ if (aborted !== undefined) {
+ yield { aborted };
+ continue;
+ }
+ yield { response: GreeterGetWholeStateResponseFromProtobufShape(response) };
}
}
;
@@ -1799,9 +1839,14 @@ class _Reactively {
id: __classPrivateFieldGet(this, __Reactively_id, "f"),
requestType: Empty,
responseType: Empty,
+ abortedType: GreeterFailWithExceptionAborted,
request: request,
signal: options === null || options === void 0 ? void 0 : options.signal,
- bearerToken: context.bearerToken,
+ // Read the token from `context` on every attempt, so that a
+ // token set with `context.setBearerToken()` after the read
+ // started is sent on the next attempt.
+ bearerToken: async () => { var _a; return await ((_a = context.bearerToken) === null || _a === void 0 ? void 0 : _a.call(context)); },
+ onUnauthenticated: context.onUnauthenticated,
websockets: context.websockets,
});
const setTypedRequest = (newRequest) => {
@@ -1809,9 +1854,12 @@ class _Reactively {
setRequest(typedRequest);
};
async function* typedGenerator() {
- for await (const response of generator) {
- const typedResponse = GreeterFailWithExceptionResponseFromProtobufShape(response);
- yield typedResponse;
+ for await (const { response, aborted } of generator) {
+ if (aborted !== undefined) {
+ yield { aborted };
+ continue;
+ }
+ yield { response: GreeterFailWithExceptionResponseFromProtobufShape(response) };
}
}
;
@@ -1826,9 +1874,14 @@ class _Reactively {
id: __classPrivateFieldGet(this, __Reactively_id, "f"),
requestType: Empty,
responseType: Empty,
+ abortedType: GreeterFailWithAbortedAborted,
request: request,
signal: options === null || options === void 0 ? void 0 : options.signal,
- bearerToken: context.bearerToken,
+ // Read the token from `context` on every attempt, so that a
+ // token set with `context.setBearerToken()` after the read
+ // started is sent on the next attempt.
+ bearerToken: async () => { var _a; return await ((_a = context.bearerToken) === null || _a === void 0 ? void 0 : _a.call(context)); },
+ onUnauthenticated: context.onUnauthenticated,
websockets: context.websockets,
});
const setTypedRequest = (newRequest) => {
@@ -1836,9 +1889,12 @@ class _Reactively {
setRequest(typedRequest);
};
async function* typedGenerator() {
- for await (const response of generator) {
- const typedResponse = GreeterFailWithAbortedResponseFromProtobufShape(response);
- yield typedResponse;
+ for await (const { response, aborted } of generator) {
+ if (aborted !== undefined) {
+ yield { aborted };
+ continue;
+ }
+ yield { response: GreeterFailWithAbortedResponseFromProtobufShape(response) };
}
}
;
@@ -1853,9 +1909,14 @@ class _Reactively {
id: __classPrivateFieldGet(this, __Reactively_id, "f"),
requestType: greeter_pb.ReadRecursiveMessageRequest,
responseType: greeter_pb.ReadRecursiveMessageResponse,
+ abortedType: GreeterReadRecursiveMessageAborted,
request: request,
signal: options === null || options === void 0 ? void 0 : options.signal,
- bearerToken: context.bearerToken,
+ // Read the token from `context` on every attempt, so that a
+ // token set with `context.setBearerToken()` after the read
+ // started is sent on the next attempt.
+ bearerToken: async () => { var _a; return await ((_a = context.bearerToken) === null || _a === void 0 ? void 0 : _a.call(context)); },
+ onUnauthenticated: context.onUnauthenticated,
websockets: context.websockets,
});
const setTypedRequest = (newRequest) => {
@@ -1863,9 +1924,12 @@ class _Reactively {
setRequest(typedRequest);
};
async function* typedGenerator() {
- for await (const response of generator) {
- const typedResponse = GreeterReadRecursiveMessageResponseFromProtobufShape(response);
- yield typedResponse;
+ for await (const { response, aborted } of generator) {
+ if (aborted !== undefined) {
+ yield { aborted };
+ continue;
+ }
+ yield { response: GreeterReadRecursiveMessageResponseFromProtobufShape(response) };
}
}
;
diff --git a/tests/reboot/nodejs/reboot_web_test/BUILD.bazel b/tests/reboot/nodejs/reboot_web_test/BUILD.bazel
index a2d28e6d7..3498481d2 100644
--- a/tests/reboot/nodejs/reboot_web_test/BUILD.bazel
+++ b/tests/reboot/nodejs/reboot_web_test/BUILD.bazel
@@ -34,5 +34,9 @@ js_reboot_test(
":test_ts",
],
entry_point = "test.js",
+ # A reactive read over plain HTTP goes over a websocket, and the
+ # global `WebSocket` a browser provides is only behind this flag in
+ # Node.js 20.
+ env = {"NODE_OPTIONS": "--experimental-websocket"},
visibility = ["//visibility:public"],
)
diff --git a/tests/reboot/nodejs/reboot_web_test/test.ts b/tests/reboot/nodejs/reboot_web_test/test.ts
index 1f31e0502..0b6fa384e 100644
--- a/tests/reboot/nodejs/reboot_web_test/test.ts
+++ b/tests/reboot/nodejs/reboot_web_test/test.ts
@@ -4,13 +4,20 @@ import {
ReaderContext,
Reboot,
TokenVerifier,
+ allowIf,
+ hasVerifiedToken,
} from "@reboot-dev/reboot";
+import { errors_pb } from "@reboot-dev/reboot-api";
import { WebContext } from "@reboot-dev/reboot-web";
import { fork } from "child_process";
import { strict as assert } from "node:assert";
import test from "node:test";
import { v4 as uuidv4 } from "uuid";
-import { Greeter } from "../../greeter_rbt_web.js";
+import {
+ ErrorWithValue,
+ Greeter,
+ GreeterFailWithAbortedAborted,
+} from "../../greeter_rbt_web.js";
import { GreeterServicer } from "../greeter.js";
const TOKEN_FOR_TEST = "S3CR3T!";
@@ -21,6 +28,17 @@ const TOKEN_FOR_TEST = "S3CR3T!";
// retry, which we can't do yet to the best of our knowledge with the
// tests in 'tests/reboot/react'.
+// The next item of a reactive read, which must not be done.
+async function nextItem- (
+ items: AsyncGenerator
-
+): Promise
- {
+ const result = await items.next();
+ if (result.done === true) {
+ assert.fail("Expected another item");
+ }
+ return result.value;
+}
+
class StaticTokenVerifier extends TokenVerifier {
async verifyToken(
context: ReaderContext,
@@ -31,6 +49,23 @@ class StaticTokenVerifier extends TokenVerifier {
}
}
+// Verifies `TOKEN_FOR_TEST` and no other token.
+class OnlyTokenForTestVerifier extends TokenVerifier {
+ async verifyToken(
+ context: ReaderContext,
+ token?: string
+ ): Promise {
+ return token === TOKEN_FOR_TEST ? new Auth({ userId: "test" }) : null;
+ }
+}
+
+// A `GreeterServicer` whose methods require a verified token.
+class AuthenticatedGreeterServicer extends GreeterServicer {
+ authorizer() {
+ return allowIf({ all: [hasVerifiedToken] });
+ }
+}
+
test("Reboot", async (t) => {
await t.test("Non reactive calls", async (t) => {
const application = new Application({
@@ -197,6 +232,298 @@ test("Reboot", async (t) => {
assert(response.message == "Hi , I am Dr Jonathan the Friendly");
});
+ await t.test("Reactive reader", async (t) => {
+ const application = new Application({
+ servicers: [GreeterServicer],
+ });
+
+ const rbt = new Reboot();
+ await rbt.start();
+
+ t.after(async () => {
+ await rbt.stop();
+ });
+
+ await rbt.up(application, { localEnvoy: true });
+
+ const context = new WebContext({
+ url: rbt.url(),
+ });
+
+ const [greeter] = await Greeter.create(context, {
+ title: "Dr",
+ name: "Jonathan",
+ adjective: "Best",
+ });
+
+ const abortController = new AbortController();
+
+ const [items] = await greeter
+ .reactively()
+ .greet(context, {}, { signal: abortController.signal });
+
+ const first = await nextItem(items);
+ assert(first.response?.message == "Hi , I am Dr Jonathan the Best");
+
+ await greeter.setAdjective(context, {
+ adjective: "Friendly",
+ });
+
+ // The generator yields a response for each change to the state.
+ const second = await nextItem(items);
+ assert(second.response?.message == "Hi , I am Dr Jonathan the Friendly");
+
+ abortController.abort();
+ assert((await items.next()).done);
+ });
+
+ await t.test(
+ "Reactive reader yields a declared error and keeps reading",
+ async (t) => {
+ const application = new Application({
+ servicers: [GreeterServicer],
+ });
+
+ const rbt = new Reboot();
+ await rbt.start();
+
+ t.after(async () => {
+ await rbt.stop();
+ });
+
+ await rbt.up(application, { localEnvoy: true });
+
+ const context = new WebContext({
+ url: rbt.url(),
+ });
+
+ const [greeter] = await Greeter.create(context, {
+ title: "Dr",
+ name: "Jonathan",
+ adjective: "Best",
+ });
+
+ const abortController = new AbortController();
+
+ const [items] = await greeter
+ .reactively()
+ .failWithAborted(context, {}, { signal: abortController.signal });
+
+ // A declared error is yielded rather than thrown, and it does not
+ // end the read: the next attempt, after a backoff, yields it
+ // again.
+ for (let i = 0; i < 2; i++) {
+ const { response, aborted } = await nextItem(items);
+ assert(response === undefined);
+ assert(aborted instanceof GreeterFailWithAbortedAborted);
+ assert(aborted.error instanceof ErrorWithValue);
+ assert(aborted.error.value == "Hi!");
+ }
+
+ abortController.abort();
+ assert((await items.next()).done);
+ }
+ );
+
+ await t.test(
+ "Reactive reader calls `onUnauthenticated` until the session is renewed",
+ async (t) => {
+ const application = new Application({
+ servicers: [AuthenticatedGreeterServicer],
+ tokenVerifier: new OnlyTokenForTestVerifier(),
+ });
+
+ const rbt = new Reboot();
+ await rbt.start();
+
+ t.after(async () => {
+ await rbt.stop();
+ });
+
+ await rbt.up(application, { localEnvoy: true });
+
+ const [greeter] = await Greeter.create(
+ new WebContext({
+ url: rbt.url(),
+ bearerToken: TOKEN_FOR_TEST,
+ }),
+ {
+ title: "Dr",
+ name: "Jonathan",
+ adjective: "Best",
+ }
+ );
+
+ // The read starts with a token that the verifier rejects.
+ let token = "expired";
+ let calls = 0;
+
+ const context = new WebContext({
+ url: rbt.url(),
+ bearerToken: async () => token,
+ onUnauthenticated: async () => {
+ calls += 1;
+ if (calls === 1) {
+ // The renewal fails, e.g., because the backend is
+ // restarting.
+ return false;
+ }
+ if (calls === 2) {
+ // The renewal reports success, but the token is still the
+ // rejected one.
+ return true;
+ }
+ token = TOKEN_FOR_TEST;
+ return true;
+ },
+ });
+
+ const abortController = new AbortController();
+
+ const [items] = await greeter
+ .reactively()
+ .greet(context, {}, { signal: abortController.signal });
+
+ // The failed renewal yields the error.
+ const first = await nextItem(items);
+ assert(first.aborted?.error instanceof errors_pb.Unauthenticated);
+ assert.equal(calls, 1);
+
+ // The next attempt calls the hook again. It returns `true`, so
+ // the read reconnects right away, and because the token is still
+ // rejected the error is yielded without a third call.
+ const second = await nextItem(items);
+ assert(second.aborted?.error instanceof errors_pb.Unauthenticated);
+ assert.equal(calls, 2);
+
+ // The attempt after the backoff calls the hook a third time,
+ // which renews the token.
+ const third = await nextItem(items);
+ assert(third.response?.message == "Hi , I am Dr Jonathan the Best");
+ assert.equal(calls, 3);
+
+ abortController.abort();
+ assert((await items.next()).done);
+ }
+ );
+
+ await t.test(
+ "Reactive reader sends the token that `setBearerToken` set after it started",
+ async (t) => {
+ const application = new Application({
+ servicers: [AuthenticatedGreeterServicer],
+ tokenVerifier: new OnlyTokenForTestVerifier(),
+ });
+
+ const rbt = new Reboot();
+ await rbt.start();
+
+ t.after(async () => {
+ await rbt.stop();
+ });
+
+ await rbt.up(application, { localEnvoy: true });
+
+ const [greeter] = await Greeter.create(
+ new WebContext({
+ url: rbt.url(),
+ bearerToken: TOKEN_FOR_TEST,
+ }),
+ {
+ title: "Dr",
+ name: "Jonathan",
+ adjective: "Best",
+ }
+ );
+
+ // The read starts with a token that the verifier rejects, and
+ // `onUnauthenticated` renews it with `setBearerToken`.
+ let calls = 0;
+
+ const context: WebContext = new WebContext({
+ url: rbt.url(),
+ bearerToken: "expired",
+ onUnauthenticated: async () => {
+ calls += 1;
+ context.setBearerToken(TOKEN_FOR_TEST);
+ return true;
+ },
+ });
+
+ const abortController = new AbortController();
+
+ const [items] = await greeter
+ .reactively()
+ .greet(context, {}, { signal: abortController.signal });
+
+ // The attempt after the renewal sends the new token, so the
+ // first item is a response and not the `Unauthenticated` error.
+ const first = await nextItem(items);
+ assert(first.aborted === undefined);
+ assert(first.response?.message == "Hi , I am Dr Jonathan the Best");
+ assert.equal(calls, 1);
+
+ abortController.abort();
+ assert((await items.next()).done);
+ }
+ );
+
+ await t.test("Reactive reader retries a restarting server", async (t) => {
+ const application = new Application({
+ servicers: [GreeterServicer],
+ });
+
+ const rbt = new Reboot();
+ await rbt.start();
+
+ t.after(async () => {
+ await rbt.stop();
+ });
+
+ await rbt.up(application, { localEnvoy: true });
+
+ const context = new WebContext({
+ url: rbt.url(),
+ });
+
+ const [greeter] = await Greeter.create(context, {
+ title: "Dr",
+ name: "Jonathan",
+ adjective: "Best",
+ });
+
+ const abortController = new AbortController();
+
+ const [responses] = await greeter
+ .reactively()
+ .greet(context, {}, { signal: abortController.signal });
+
+ const first = await nextItem(responses);
+ assert(first.response?.message == "Hi , I am Dr Jonathan the Best");
+
+ // Restarting the server disconnects the reactive read, which
+ // reconnects rather than surface the disconnect, and then
+ // observes the mutation made after the restart.
+ await rbt.down();
+ await rbt.up(application, { localEnvoy: true });
+
+ await greeter.setAdjective(context, {
+ adjective: "Friendly",
+ });
+
+ while (true) {
+ const { response, aborted } = await nextItem(responses);
+ assert(aborted === undefined);
+ if (response?.message == "Hi , I am Dr Jonathan the Friendly") {
+ break;
+ }
+ }
+
+ // Aborting the read ends the generator.
+ abortController.abort();
+ assert((await responses.next()).done);
+ });
+
await t.test("Transaction", async (t) => {
const application = new Application({
servicers: [GreeterServicer],