-
-
Notifications
You must be signed in to change notification settings - Fork 13
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Broker-clients example and renamed broker run() to serve_forever().
- Loading branch information
Showing
3 changed files
with
100 additions
and
5 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,95 @@ | ||
""" | ||
+--------+ +--------+ +-------------+ | ||
| |--- ping -->| |--- ping -->| | | ||
| client | | broker | | echo client | | ||
| |<-- pong ---| |<-- pong ---| | | ||
+--------+ +--------+ +-------------+ | ||
""" | ||
|
||
import time | ||
import asyncio | ||
import mqttools | ||
|
||
|
||
BROKER_PORT = 10008 | ||
|
||
|
||
async def start_client(): | ||
client = mqttools.Client('localhost', BROKER_PORT) | ||
|
||
while True: | ||
try: | ||
await client.start() | ||
break | ||
except: | ||
print('Client start failed. Retrying...') | ||
await asyncio.sleep(0.2) | ||
|
||
return client | ||
|
||
|
||
async def client_main(): | ||
"""Publish the current time to /ping and wait for the echo client to | ||
publish it back on /pong, with a one second interval. | ||
""" | ||
|
||
client = await start_client() | ||
await client.subscribe('/pong') | ||
|
||
while True: | ||
print() | ||
message = str(int(time.time())).encode('ascii') | ||
print(f'client: Publishing {message} on /ping.') | ||
client.publish('/ping', message) | ||
topic, message = await client.messages.get() | ||
print(f'client: Got {message} on {topic}.') | ||
|
||
if topic is None: | ||
print('Client connection lost.') | ||
break | ||
|
||
await asyncio.sleep(1) | ||
|
||
|
||
async def echo_client_main(): | ||
"""Wait for the client to publish to /ping, and publish /pong in | ||
response. | ||
""" | ||
|
||
client = await start_client() | ||
await client.subscribe('/ping') | ||
|
||
while True: | ||
topic, message = await client.messages.get() | ||
print(f'echo_client: Got {message} on {topic}.') | ||
|
||
if topic is None: | ||
print('Echo client connection lost.') | ||
break | ||
|
||
print(f'echo_client: Publishing {message} on /pong.') | ||
client.publish('/pong', message) | ||
|
||
|
||
async def broker_main(): | ||
"""The broker, serving both clients, forever. | ||
""" | ||
|
||
broker = mqttools.Broker('localhost', BROKER_PORT) | ||
await broker.serve_forever() | ||
|
||
|
||
async def main(): | ||
await asyncio.gather( | ||
broker_main(), | ||
echo_client_main(), | ||
client_main() | ||
) | ||
|
||
|
||
asyncio.run(main()) |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters