Skip to content
Open
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
38 changes: 19 additions & 19 deletions src/frequenz/sdk/actor/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -380,9 +380,9 @@ async def _run(self) -> None:

async def main() -> None: # (2)!
# (4)!
input_channel: Broadcast[str] = Broadcast("Input to Actor1")
middle_channel: Broadcast[str] = Broadcast("Actor1 -> Actor2 stream")
output_channel: Broadcast[str] = Broadcast("Actor2 output")
input_channel: Broadcast[str] = Broadcast(name="Input to Actor1")
middle_channel: Broadcast[str] = Broadcast(name="Actor1 -> Actor2 stream")
output_channel: Broadcast[str] = Broadcast(name="Actor2 output")

input_sender = input_channel.new_sender()
output_receiver = output_channel.new_receiver()
Expand Down Expand Up @@ -460,8 +460,7 @@ async def main() -> None: # (2)!
```python title="select.py"
import asyncio

from frequenz.channels import Broadcast, Receiver, Sender
from frequenz.channels.util import select, selected_from
from frequenz.channels import Broadcast, Receiver, Sender, select, selected_from
from frequenz.sdk.actor import Actor, run


Expand All @@ -481,14 +480,14 @@ def __init__(
async def _run(self) -> None: # (2)!
async for selected in select(self._receiver_1, self._receiver_2): # (10)!
if selected_from(selected, self._receiver_1): # (11)!
print(f"Received from receiver_1: {selected.value}")
await self._output.send(selected.value)
if not selected.value: # (12)!
print(f"Received from receiver_1: {selected.message}")
await self._output.send(selected.message)
if not selected.message: # (12)!
break
elif selected_from(selected, self._receiver_2): # (13)!
print(f"Received from receiver_2: {selected.value}")
await self._output.send(selected.value)
if not selected.value: # (14)!
print(f"Received from receiver_2: {selected.message}")
await self._output.send(selected.message)
if not selected.message: # (14)!
break
else:
assert False, "Unknown selected channel"
Expand All @@ -497,9 +496,9 @@ async def _run(self) -> None: # (2)!


# (3)!
input_channel_1 = Broadcast[bool]("input_channel_1")
input_channel_2 = Broadcast[bool]("input_channel_2")
echo_channel = Broadcast[bool]("echo_channel")
input_channel_1 = Broadcast[bool](name="input_channel_1")
input_channel_2 = Broadcast[bool](name="input_channel_2")
echo_channel = Broadcast[bool](name="echo_channel")

echo_actor = EchoActor( # (4)!
input_channel_1.new_receiver(),
Expand Down Expand Up @@ -555,23 +554,23 @@ async def main() -> None: # (6)!
10. The [`select()`][frequenz.channels.select] function will get the first message
available from the two channels. The order in which they will be handled is
unknown, but in this example we assume that the first message will be from
`input_channel_1` (`True`) and the second from `input_channel_1` (`False`).
`input_channel_1` (`True`) and the second from `input_channel_2` (`False`).

11. The [`selected_from()`][frequenz.channels.selected_from] function will return
`True` for the `input_channel_1` receiver. `selected.value` holds the received
`True` for the `input_channel_1` receiver. `selected.message` holds the received
message, so `"Received from receiver_1: True"` will be printed and `True` will be
sent to the `output` channel.

12. Since `selected.value` is `True`, the loop will continue, going back to the
12. Since `selected.message` is `True`, the loop will continue, going back to the
[`select()`][frequenz.channels.select] function.

13. The [`selected_from()`][frequenz.channels.selected_from] function will return
`False` for the `input_channel_1` receiver and `True` for the `input_channel_2`
receiver. The message stored in `selected.value` will now be `False`, so
receiver. The message stored in `selected.message` will now be `False`, so
`"Received from receiver_2: False"` will be printed and `False` will be sent to the
`output` channel.

14. Since `selected.value` is `False`, the loop will break.
14. Since `selected.message` is `False`, the loop will break.

15. The [`_run()`][_run] method will finish normally and the actor will be stopped, so
the [`run()`][frequenz.sdk.actor.run] function will return.
Expand All @@ -589,6 +588,7 @@ async def main() -> None: # (6)!
```
Received from receiver_1: True
Received from receiver_2: False
EchoActor finished
Received message=True
Received message=False
```
Expand Down
Loading