Skip to content
Merged
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
17 changes: 13 additions & 4 deletions documentation/docs/call/from_outside_your_app.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -96,18 +96,27 @@ const context = new ExternalContext({
## Calling readers reactively

Using an `ExternalContext`, you can also call `reader` methods
_reactively_: you receive a new response every time the state
changes. For example:
_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:

<Tabs groupId="language">
<TabItem value="python" label="Python" default>
<!-- MARKDOWN-AUTO-DOCS:START
(CODE:src=../../../reboot/demos/fig/backend/src/many_readers.py&lines=20-22) -->
(CODE:src=../../../reboot/demos/fig/backend/src/many_readers.py&lines=20-25) -->
<!-- The below code snippet is automatically added from ../../../reboot/demos/fig/backend/src/many_readers.py -->

```py
fig = Fig.ref(fig_id)
async for response in fig.reactively().get_position(context):
async for response, aborted in fig.reactively().get_position(context):
if aborted is not None:
print(f"{fig_id}: {aborted}")
continue
print(f"{fig_id}: {response}")
```

Expand Down
12 changes: 11 additions & 1 deletion reboot/bdd/steps.py
Original file line number Diff line number Diff line change
Expand Up @@ -1146,7 +1146,7 @@ async def _eventually_has(
)
)
try:
response = await asyncio.wait_for(
response, aborted = await asyncio.wait_for(
anext(responses), timeout=remaining
)
except asyncio.TimeoutError:
Expand All @@ -1160,6 +1160,16 @@ async def _eventually_has(
)
) from None
else:
if aborted is not None:
# The read was answered with an error, e.g., the
# state is not constructed yet. The reactive read
# keeps going and is answered again once the state
# changes, so keep waiting for a response.
last_error = AssertionError(
f"the last read was answered with {aborted}"
)
continue
assert response is not None
try:
_assert_properties(response, assertions)
except AssertionError as error:
Expand Down
10 changes: 7 additions & 3 deletions reboot/cli/commands/cloud/up.py
Original file line number Diff line number Diff line change
Expand Up @@ -222,9 +222,13 @@ async def cloud_up(args: argparse.Namespace) -> int:
traceback.print_exc()
terminal.fail("Please report this bug to the maintainers")

async for status_response in application.reactively().RevisionStatus(
context, revision_number=up_response.revision_number
):
async for status_response, status_aborted in application.reactively(
).RevisionStatus(context, revision_number=up_response.revision_number):
if status_aborted is not None:
print(f"🛑 unexpected error: {status_aborted}")
terminal.fail("Please report this bug to the maintainers")
assert status_response is not None

revision = status_response.revision
if revision.status == Status.UPPING:
# Keep waiting.
Expand Down
11 changes: 9 additions & 2 deletions reboot/demos/fig/backend/src/many_readers.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,14 +18,21 @@ async def simulate_reactive_browser_reader(port: Optional[int]):

async def reactively_print_fig_position(fig_id: str):
fig = Fig.ref(fig_id)
async for response in fig.reactively().get_position(context):
async for response, aborted in fig.reactively().get_position(context):
if aborted is not None:
print(f"{fig_id}: {aborted}")
continue
Comment thread
rileysdev marked this conversation as resolved.
print(f"{fig_id}: {response}")

fig_board = FigBoard.ref(FIG_BOARD_ID)

fig_tasks = {}

async for response in fig_board.reactively().list(context):
async for response, aborted in fig_board.reactively().list(context):
if aborted is not None:
print(f"{FIG_BOARD_ID}: {aborted}")
continue
Comment thread
rileysdev marked this conversation as resolved.
assert response is not None

for fig_id in response.fig_ids:
if fig_id not in fig_tasks:
Expand Down
16 changes: 11 additions & 5 deletions reboot/plugin/skills/python/references/testing-external-context.md
Original file line number Diff line number Diff line change
Expand Up @@ -144,22 +144,28 @@ with self.assertRaises(Aborted):
A "another session sees the change without refreshing" story is
tested with `reactively()`, the same push-based subscription the
generated React hooks use. `Type.ref(id).reactively().<reader>(context)`
returns an **async iterator** that yields a fresh response on every
state change; `anext()` pulls the next one. Never poll in a loop
with `asyncio.sleep` — that tests the sleep, not the reactivity.
returns an **async iterator** that yields a `(response, aborted)`
pair on every state change: 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 iterator; it
yields again once the state changes. `anext()` pulls the next pair.
Never poll in a loop with `asyncio.sleep` — that tests the sleep, not
the reactivity.

```python
# The "other browser session" subscribes...
subscription = TaskList.ref(list_id).reactively().get(bob)
first = await asyncio.wait_for(anext(subscription), timeout=10)
first, aborted = await asyncio.wait_for(anext(subscription), timeout=10)
self.assertIsNone(aborted)
self.assertEqual(first.tasks, [])

# ...another session writes...
await TaskList.ref(list_id).add_task(alice, title="Milk")

# ...and the subscription is pushed the update.
while True:
update = await asyncio.wait_for(anext(subscription), timeout=10)
update, aborted = await asyncio.wait_for(anext(subscription), timeout=10)
self.assertIsNone(aborted)
if len(update.tasks) == 1:
break
self.assertEqual(update.tasks[0].title, "Milk")
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
## Python `reactively()` reads yield `(response, aborted)` pairs

`Type.ref(id).reactively().<reader>(context)` used to yield bare
responses and to raise the method's `<Type>.<Method>Aborted` when the
backend answered the read with an error (a declared error raised by
the reader, a denied `authorizer()`, `StateNotConstructed`, ...),
which ended the read. It now yields a `(response, aborted)` pair for
every answer and never raises: for a response, `aborted` is `None`;
for an error, `response` is `None` and `aborted` is the method's
`<Type>.<Method>Aborted`. An error does not end the read: the
subscription stays open and yields again once the state changes.
Failed connections are still retried with backoff without yielding
anything.

Find every use of such a read. Match only the opening parenthesis,
since a formatter may have split the call across lines:

grep -rn "\.reactively(" --include=*.py

Unpack the pair in every `async for` over one. Before:

```python
async for response in cart.reactively().get(context):
render(response)
```

After:

```python
async for response, aborted in cart.reactively().get(context):
if aborted is not None:
# Decide: `continue` to wait for the state to change, `break`
# to stop reading, or `raise aborted`.
continue
render(response)
```

Also:

- Code that caught the `Aborted` around the loop (`try: async for ... except Type.MethodAborted:`) must handle `aborted` inside the loop
instead; the `except` is now unreachable.
- Code that pulls items with `anext(...)` gets the pair too:
`response, aborted = await anext(subscription)`.
- A reader whose response type is empty yields `(None, None)` on
success, so check `aborted`, not `response`, for those.
45 changes: 31 additions & 14 deletions reboot/templates/reboot.py.j2
Original file line number Diff line number Diff line change
Expand Up @@ -5201,7 +5201,7 @@ class {{ client.proto.state_name }}:
__context__: IMPORT_reboot_aio_external.ExternalContext | IMPORT_reboot_aio_contexts.ReaderContext | IMPORT_reboot_aio_contexts.WorkflowContext,
__request_or_options__: {{ client.proto.state_name }}.{{ method.proto.name }}Request,
__options__: IMPORT_typing.Optional[IMPORT_reboot_aio_call.Options] = None,
) -> IMPORT_typing.AsyncIterator[{% if method.has_non_none_response %}{{ client.proto.state_name }}.{{ method.proto.name }}Response{% else %}None{% endif %}]:
) -> IMPORT_typing.AsyncIterator[tuple[{% if method.has_non_none_response %}{{ client.proto.state_name }}.{{ method.proto.name }}Response{% else %}None{% endif %}, None] | tuple[None, {{ client.proto.state_name }}.{{ method.proto.name }}Aborted]]:
...

@IMPORT_typing.overload
Expand All @@ -5219,7 +5219,7 @@ class {{ client.proto.state_name }}:
{{ name }}: IMPORT_typing.Optional[{{ type }}] | Unset = UNSET,
{% endif %}
{% endfor %}
) -> IMPORT_typing.AsyncIterator[{% if method.has_non_none_response %}{{ client.proto.state_name }}.{{ method.proto.name }}Response{% else %}None{% endif %}]:
) -> IMPORT_typing.AsyncIterator[tuple[{% if method.has_non_none_response %}{{ client.proto.state_name }}.{{ method.proto.name }}Response{% else %}None{% endif %}, None] | tuple[None, {{ client.proto.state_name }}.{{ method.proto.name }}Aborted]]:
...
{% endif %}

Expand Down Expand Up @@ -5250,7 +5250,7 @@ class {{ client.proto.state_name }}:
{{ name }}: IMPORT_typing.Optional[{{ type }}] | Unset = UNSET,
{% endif %}
{% endfor %}
) -> IMPORT_typing.AsyncIterator[{% if method.has_non_none_response %}{{ client.proto.state_name }}.{{ method.proto.name }}Response{% else %}None{% endif %}]:
) -> IMPORT_typing.AsyncIterator[tuple[{% if method.has_non_none_response %}{{ client.proto.state_name }}.{{ method.proto.name }}Response{% else %}None{% endif %}, None] | tuple[None, {{ client.proto.state_name }}.{{ method.proto.name }}Aborted]]:
IMPORT_reboot_aio_types.assert_type(__context__, [IMPORT_reboot_aio_external.ExternalContext, IMPORT_reboot_aio_contexts.ReaderContext, IMPORT_reboot_aio_contexts.WorkflowContext])
{% if method.has_non_none_request %}
# UX improvement: check that neither positional argument was accidentally
Expand Down Expand Up @@ -5370,7 +5370,7 @@ class {{ client.proto.state_name }}:

__response__ = {{ method.output_type }}()
__response__.ParseFromString(__query_response__.response)
yield {{ client.proto.state_name }}{{ method.proto.name }}ResponseFromProto(__response__)
yield ({{ client.proto.state_name }}{{ method.proto.name }}ResponseFromProto(__response__), None)

except IMPORT_grpc.aio.AioRpcError as error:
# We expect to get disconnected from the server
Expand All @@ -5385,18 +5385,35 @@ class {{ client.proto.state_name }}:
)
await __query_backoff__()
continue
if error.code() == IMPORT_grpc.StatusCode.ABORTED:
Comment thread
rileysdev marked this conversation as resolved.
# Reconstitute the error that the server threw, if it was a declared error.
status = await IMPORT_rpc_status_async.from_call(__call__)
if status is not None:
raise {{ client.proto.state_name }}.{{ method.proto.name }}Aborted.from_status(
status
) from None
raise {{ client.proto.state_name }}.{{ method.proto.name }}Aborted.from_grpc_aio_rpc_error(

# The server answered the read with an error,
# e.g., a declared error raised by the reader, a
# denied authorizer, or a state that has not been
# constructed (yet). That answer is a value to
# the caller, not the end of the read: the state
# may change so that the next read succeeds, so
# we keep reading, and it is up to the caller to
# stop iterating (or to raise) if the error is
# final for them.
#
# Reconstitute the error that the server threw,
# if it was a declared error.
status = (
await IMPORT_rpc_status_async.from_call(__call__)
if __call__ is not None else None
)
if status is not None:
__aborted__ = {{ client.proto.state_name }}.{{ method.proto.name }}Aborted.from_status(
status
)
else:
__aborted__ = {{ client.proto.state_name }}.{{ method.proto.name }}Aborted.from_grpc_aio_rpc_error(
error
) from None
)

raise
yield (None, __aborted__)

await __query_backoff__()

# Keep the original functions on the client, so old code will
# continue to work, but use the new 'snake_case' method in
Expand Down
10 changes: 10 additions & 0 deletions tests/reboot/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -384,6 +384,16 @@ py_test(
],
)

py_test(
name = "reactive_reader_tests_py",
srcs = [":reactive_reader_tests.py"],
main = "reactive_reader_tests.py",
deps = [
":greeter_servicers_py",
"//reboot/aio:tests_py",
],
)

py_test(
name = "external_context_tests_py",
srcs = [":external_context_tests.py"],
Expand Down
4 changes: 3 additions & 1 deletion tests/reboot/bank.py
Original file line number Diff line number Diff line change
Expand Up @@ -968,7 +968,9 @@ async def test_reactive_method_call(
) -> Empty:
# Call a method reactively, and see that we don't crash.
account = Account.ref(self.state.account_ids[0])
async for _ in account.reactively().balance(context):
async for _, aborted in account.reactively().balance(context):
if aborted is not None:
raise aborted
# Yay, that worked.
break

Expand Down
7 changes: 5 additions & 2 deletions tests/reboot/dashboard/api_watcher_tests.py
Original file line number Diff line number Diff line change
Expand Up @@ -122,8 +122,11 @@ async def _wait_for_api(self, satisfied):
whenever it changes."""
context = self.rbt.create_external_context(name=self.id())

async for response in Dashboard.ref(DASHBOARD_ID
).reactively().Get(context):
async for response, aborted in Dashboard.ref(DASHBOARD_ID).reactively(
).Get(context):
if aborted is not None:
raise aborted
assert response is not None
if satisfied(response):
return response

Expand Down
26 changes: 19 additions & 7 deletions tests/reboot/dashboard/code_watcher_tests.py
Original file line number Diff line number Diff line change
Expand Up @@ -415,8 +415,11 @@ async def _code_changes(self, *, satisfied) -> list[Change]:
"""The code's entries in the changelog once they satisfy,
newest first, reading again whenever the changelog changes."""
context = self.rbt.create_external_context(name=self.id())
async for response in OrderedMap.ref(CHANGELOG_ID).reactively(
async for response, aborted in OrderedMap.ref(CHANGELOG_ID).reactively(
).ReverseRange(context, limit=100):
if aborted is not None:
raise aborted
assert response is not None
changes = [
Change.FromString(entry.bytes) for entry in response.entries
]
Expand All @@ -438,8 +441,11 @@ async def _servicers(self, *, satisfied):
"""
context = self.rbt.create_external_context(name=self.id())

async for response in Dashboard.ref(DASHBOARD_ID
).reactively().Get(context):
async for response, aborted in Dashboard.ref(DASHBOARD_ID).reactively(
).Get(context):
if aborted is not None:
raise aborted
assert response is not None
found: dict[str, list[Servicer]] = {}

for servicer in response.servicers:
Expand All @@ -455,8 +461,11 @@ async def _implementation(self, *, satisfied):
reading again whenever it changes."""
context = self.rbt.create_external_context(name=self.id())

async for response in Dashboard.ref(DASHBOARD_ID
).reactively().Get(context):
async for response, aborted in Dashboard.ref(DASHBOARD_ID).reactively(
).Get(context):
if aborted is not None:
raise aborted
assert response is not None
if satisfied(response):
return response

Expand Down Expand Up @@ -569,8 +578,11 @@ async def test_what_is_declared_and_what_implements_it_are_separate(

context = self.rbt.create_external_context(name=self.id())

async for response in Dashboard.ref(DASHBOARD_ID
).reactively().Get(context):
async for response, aborted in Dashboard.ref(DASHBOARD_ID).reactively(
).Get(context):
if aborted is not None:
raise aborted
assert response is not None
if any(
f'{api.package}.{state_type.name}' == 'shop.v1.Shop'
for api in response.apis.values()
Expand Down
8 changes: 6 additions & 2 deletions tests/reboot/dashboard/dashboard_tests.py
Original file line number Diff line number Diff line change
Expand Up @@ -392,8 +392,12 @@ async def _wait_for_suppress_open_on_restart(self, expected: bool) -> None:
write is what the page sends after the click, so seeing it is
how the test knows the choice reached the application."""
context = self.rbt.create_external_context(name=self.id())
async for response in Preferences.ref(PREFERENCES_ID
).reactively().Get(context):
async for response, aborted in Preferences.ref(
PREFERENCES_ID
).reactively().Get(context):
if aborted is not None:
raise aborted
assert response is not None
if response.suppress_open_on_restart == expected:
return

Expand Down
7 changes: 5 additions & 2 deletions tests/reboot/dashboard/features_watcher_tests.py
Original file line number Diff line number Diff line change
Expand Up @@ -88,8 +88,11 @@ async def _wait_for_features(self, satisfied):
again whenever they change."""
context = self.rbt.create_external_context(name=self.id())

async for response in Dashboard.ref(DASHBOARD_ID
).reactively().Get(context):
async for response, aborted in Dashboard.ref(DASHBOARD_ID).reactively(
).Get(context):
if aborted is not None:
raise aborted
assert response is not None
if satisfied(response.features):
return response.features

Expand Down
Loading
Loading