A simple websocket/msgpack messaging library that provides pub/sub and rpc patterns.
Project description
wsmsg Package
Overview
Websockets are used, and the server will accept all authenticated connections to an endpoint root (default '/'). Authentication is via an HTTP bearer token on initial connection. Each node attaches to an endpoint specified in the connection URI. For example: wss://<host>:<port>/<endpoint>. Channels receive (and queue) messages. Subscriptions bind an endpoint to one or more channels. Subjects can be used to filter messages in a subscription. Message delivery is at most once, with flow control and no guaranteed persistence.
Flow Diagram
Publish Channel Routing Endpoint
-------+--------+-----------+-------------
/ Queue -> Forward \ / Queue -> Node
Node < Queue -> Forward X / Node
\ Queue -> Forward / \ Queue < Node
\ Node
Message Format
Messages are sent/received as a msgpack array
- uid: bytes
Binary-safe unique id of message - recommand ULID as bytes - channel: str
Channel this message will be sent to - subject: str
Subject of message - reference: bytes
Binary-safe id of message that this references - data: any
Subject-defined message data
Python
Message = NamedTuple('Message', uid=bytes, channel=str, subject=str, reference=bytes, data=any)
class Message(NamedTuple):
uid: bytes
channel: str
subject: str
reference: bytes
data: any
def __repr__(self):
elements = [
str(ulid.from_bytes(self.uid)),
self.channel,
self.subject,
str(ulid.from_bytes(self.reference)) if self.reference else None,
repr(self.data)
]
return "Message(%s)" % ', '.join(elements)
System Messages
The 'system' channel has ultimate priority (100) over all other channels.
If there are any queued messages in the system queue, no other channels will forward.
Commands
- ack
Acknowledge message receipt, not forwarded.
Ack'd message uid should be in the refercence field.
Next message for connection will not be dequeued or sent
until previous message was ackowledged. - subcribe
- channel (endpoint: str, channel_pattern: str, subject_pattern: str)
Subscribe endpoint to channel pattern with subject pattern filter. - message (endpoint: str, message_uid: bytes, expires: int)
Subscribe endpoint to any messages that reference the message uid until the expiration (unix timestamp).
- channel (endpoint: str, channel_pattern: str, subject_pattern: str)
- unsubscribe (endpoint: str, channel_pattern: str)
- channel (endpoint: str, channel_pattern: str)
Unsubscribe endpoint from patterns matching pattern.
Includes any subject filters. - message (endpoint: str, message_uid: bytes)
Unsubscribe endpoint to any messages that reference the message uid.
- channel (endpoint: str, channel_pattern: str)
- pause
Pause all message transmissions.
Messages can still be received and system commands processed. - resume
Resume message transmissions. - stop
Cleanly stop and shutdown server.
Pending messages are persisted. - channel
- create (channel: str)
Create channel queue. - size (channel: str, size: int)
Set max size of channel queue. Old messages are dropped when full - priority (channel: str, priority: int)
Set channel priority. Higher is better, range 0-99. Default of zero.
- create (channel: str)
- endpoint
- create (endpoint: str)
Create endpoint queue - size (endpoint: str, size: int)
Set max size of endpoint queue. Old messages are dropped when full
- create (endpoint: str)
System Events
- stopping Server is shutting down (in one second), no more messages will be delivered
- node
- joined (endpoint, uid)
- left (endpoint, uid)
- endpoint
- created (endpoint)
- sized (channel: str)
- channel
- created (channel: str)
- sized (channel: str)
- prioritized (channel: str)
Misc
Router
- Until stop
- Wait for next message
- For each subscriber
- Push message to node Queue
Route Map?
class Route NamedTuple('Route', router=Router, dest=set[UUID])
RouteMap = dict[UUID, Route]
routes: RouteMap =
{
<channel>: <route>,
...
}
while True:
msg = channel.get()
for node in routes[channel].dest:
forward_message(msg, node)
Examples
Task Queue
By connecting all nodes to one endpoint, each node will get a fair shair of the messages in the 'order' channel with the subject 'submitted'.
All nodes connects to ws://<host>:<port>/group1
node1 -> Message(ulid.new().bytes, 'system', 'subscribe.channel', None, ('group1', 'orders', 'submitted'))
node2 -> Message(ulid.new().bytes, 'system', 'subscribe.channel', None, ('group1', 'orders', 'submitted'))
node3 -> Message(ulid.new().bytes, 'system', 'subscribe.channel', None, ('group1', 'orders', 'submitted'))
Fanout
By using one endpoint per node, each node will get a copy of the messages in the 'order' channel with the subject 'submitted'.
Each node connects to ws://<host>:<port>/<node id>
node1 -> Message(ulid.new().bytes, 'system', 'subscribe.channel', None, (<node id>, 'orders', 'submitted'))
node2 -> Message(ulid.new().bytes, 'system', 'subscribe.channel', None, (<node id>, 'orders', 'submitted'))
node3 -> Message(ulid.new().bytes, 'system', 'subscribe.channel', None, (<node id>, 'orders', 'submitted'))
Direct Delivery
If each node subscribes to a unique addressable channel, messages can be routed to each node by the sender.
To send a message to a single node, use channel 'node.' with any subject.
To send a message to all nodes, additionally subscribe all nodes to a 'nodes' channel and send to that.
Each node connects to ws://<host>:<port>/<node id>
node1 -> Message(ulid.new().bytes, 'system', 'subscribe.channel', None, (<node id>, 'node.<node id>', '*'))
node2 -> Message(ulid.new().bytes, 'system', 'subscribe.channel', None, (<node id>, 'node.<node id>', '*'))
node3 -> Message(ulid.new().bytes, 'system', 'subscribe.channel', None, (<node id>, 'node.<node id>', '*'))
RPC
Advantages: Less messages to send, supports heavy RPC traffic through server.
Disadvantages: Need to know what reply channel and subjects will be used, large number of channel subscriptions can slow down message forwarding.
Connect to ws://<host>:<port>/<id>
node -> Message(ulid.new().bytes, 'system', 'subscribe.channel', None, (<id>, 'order', 'created'))
node -> Message(ulid.new().bytes, 'system', 'subscribe.channel', None, (<id>, 'order', 'exception'))
...
# Generate uid
Register Queue with RX task/thread for reference=uid
node -> Message(uid, 'order', 'submitted', None, None, Check)
...
node <- Message(<message uid>, 'order', 'created', uid, None, <check uid>)
RX task/thread matches and puts message(s) in Queue
...
Get result(s) with timeout from Queue
Resend if needed or show error on exception
Unregister Queue from RX task/thread
RPC - Streamlined
Advantages: message subscription matching is much faster, RTT should be lower. No need to know what reply subject may be.
Disadvantages: more messages to send, total throughput may be lower if heavy use of RPC.
Notes: If you set the message subscription expiration, no need to unsubscribe.
server -> Message(ulid.new().bytes, 'system', 'subscribe.channel', None, (<endpoint>, 'rpc', 'order.*'))
server -> Message(ulid.new().bytes, 'system', 'channel.priority', None, ('rpc', 99))
server -> Message(ulid.new().bytes, 'system', 'endpoint.size', None, (<endpoint>, 10))
...
Connect to ws://<host>:<port>/<id>
...
Generate uid
Register Queue with RX task/thread for reference=uid
node -> Message(ulid.new().bytes, 'system', 'subscribe.message', None, (<id>, uid, <timeout>))
node -> Message(uid, 'rpc', 'order.get.request', None, Check)
...
node <- Message(<message uid>, 'rpc', 'order.get.reply', uid, None, <check uid>)
RX task/thread matches and puts message(s) in Queue
...
Get result(s) with timeout from Queue
Resend if needed or show error on exception
Unregister Queue from RX task/thread
Project details
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file wsmsg-0.12.tar.gz.
File metadata
- Download URL: wsmsg-0.12.tar.gz
- Upload date:
- Size: 13.8 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/3.8.0 pkginfo/1.8.2 readme-renderer/33.0 requests/2.27.1 requests-toolbelt/0.9.1 urllib3/1.26.8 tqdm/4.63.0 importlib-metadata/4.11.2 keyring/23.5.0 rfc3986/2.0.0 colorama/0.4.4 CPython/3.9.2
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
c3662c6c5d7f2535906cba4e1639d7035df3f0c8d4f464841eda0a2a7c965911
|
|
| MD5 |
24623dd17ee0fad16b667b75aab5336e
|
|
| BLAKE2b-256 |
f51bdb76376c6c33a2ce79c4e33d3ee660b6e736758becd1b47a7235258d849f
|
File details
Details for the file wsmsg-0.12-py3-none-any.whl.
File metadata
- Download URL: wsmsg-0.12-py3-none-any.whl
- Upload date:
- Size: 12.7 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/3.8.0 pkginfo/1.8.2 readme-renderer/33.0 requests/2.27.1 requests-toolbelt/0.9.1 urllib3/1.26.8 tqdm/4.63.0 importlib-metadata/4.11.2 keyring/23.5.0 rfc3986/2.0.0 colorama/0.4.4 CPython/3.9.2
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
183789449eac20eea6772e30afe66028abbc7135ab56891da6a8ff284422ba1f
|
|
| MD5 |
cf94317967b54f3fb218a7bc038b3d30
|
|
| BLAKE2b-256 |
6f78a1283973770651fa5ea410bd697e499678ab1a98db47649c9bf06faa52d4
|