Skip to content
AllenNeuralDynamicsPublic

About

adapter pattern implemented as zmq client/server

Topics

Resources

Stars

3 stars

Watchers

1 watching

Forks

Repository files navigation

one-liner 🐍 ⟵ 🚂 ⟵ 🐍

License

a ZMQ-based Router pattern for interacting with remote python objects.

High level features:

  • Remote execution of Python code
  • Streaming of periodically called functions with the ability to enable/disable them
  • caching and queuing options for receiving only relevant data from RouterClient
  • Multiple topologies supported between RouterServer and RouterClient communication including
    • between processes on the same PC
    • across multiple PCs.
    • Single RouterServers to many RouterClients.
    • Cascading RouterServers: (i.e: a RouterServer can forward to another RouterServer).

Why build this?

The router pattern provides a bridge between distinct applications.

For our use cases we mostly apply this pattern to separate instrument control code from its respective GUI.

This separation provides several advantages:

  • GUIs can be developed independently of standalone projects.
  • GUIs can run on separate processes or on different machines altogether, providing more flexibility where some machines are resource-constrained.
  • Failures are siloed. A GUI can crash independent of the application code it is interfacing with.

Implementation-wise, this package acts as a lightweight wrapper around zmq, delegating most of the heavy-lifting to existing zmq solutions while hiding the messy details. Interface-wise, this package tries not to commit the user to adopting a specific architecture and supports various ways of attaching existing code so it can be "routed" elsewhere. one_liner place nicely in multi-threaded and multiprocessing-based Python projects.

Package Installation with uv

To install and develop this package, from the project root directory, run:

uv sync

To install all optional dependencies to play with the examples, run:

uv sync --group examples

Package Installation with Pip

To install this package, from the project root directory, run

pip install .

To install in editable mode, in the root directory, run:

pip install -e .

To install with supplementary dependencies for running the examples, run:

pip install -e .[examples]

Quick Reference

Vocabulary

Term Description
Server A Python process with functions and objects made available for remote communication. For instance, a collection of hardware drivers running on a lab computer.
Streams Servers can be configured to automatically call functions at a given frequency (e.g., every second) and post their results through one-liner for consumption.
Client A Python process that connects to a one-liner server, sends commands to server functions/objects, and receives data from server streams.
RPC Remote Procedure Call, an architecture wherein a client can "call a function", but the actual function logic gets run somewhere else, in a different process, or on a different computer.
ZMQ ZeroMQ, a messaging library that is used as the communication layer by one-liner

RouterClient

Method Description
call_by_name(call_name: str, args, kwargs, ...) Tell the server to call a function (the server must be configured with a named function).
call(obj_name, attr_name, ...) Tell the server to call a method on one of the objects in its namespace.
configure_stream(name, ...) Configure the client to either buffer data from a stream or only hold onto the most recent packet.
get_stream(name) Get data from a stream.
enable_stream(name) Tell the server to begin calling the configured stream function.
disable_stream(name) Tell the server to stop calling the configured stream function.
get_stream_configurations() Ask the server for a description of all its streams.

RouterServer

Attributes Description
instances a dict of objects to make available through one-liner
Methods Description
run(block=False) Start the server RPC listener and broadcaster.
add_named_call(call_name, obj_name, attr_name, ...) Set up a function that can be called by name from the client.
add_stream(stream_name, frequency_hz, obj_name, attr_name, ...) Set up a function to be called automatically at a given frequency, with results sent over one-liner.
add_stream_from_callable(stream_name, frequency_hz, obj_name, attr_name, ...) Same as add_stream, but for a bare function instead of an object method.
add_zmq_stream(name, address, enabled, ...) Add a stream output that, instead of calling a function periodically, connects to a ZMQ PUB socket and forwards messages.
get_stream_fn(name) Add a stream output tha, instead of the Server automatically calling a function periodically, lets application code run at its own rate and periodically push data into the stream by calling the function returned by get_stream_fn.
enable_stream(name) Begin calling the stream function at its configured frequency.
disable_stream(name) Stop calling the stream function at its configured frequency.
remove_stream(name) Remove a stream from the server configuration.
get_version() Get the one_liner version the server is running.
close() Close server resources.

Quickstart

Remote Function Execution

In the PC acting as the server:

from one_liner.server import RouterServer
from time import sleep

class Horn:
  def beep(self):
    return "Boop!"
  
my_horn = Horn()

server = RouterServer(instances={"my_horn": my_horn})
server.run()

while True: # Nothing to do! Object control happens in another thread.
    time.sleep(0.001)

In the PC acting as the client:

from one_liner.client import RouterClient

client = RouterClient()
result = client.call("my_horn", "beep") # call method w/ args & kwargs; return the result.

Streaming Data

There are three ways to stream data from a RouterServer to one or more RouterClient objects.

Periodic Streaming

In the PC acting as the server:

from one_liner.server import RouterServer
import cv2

video = cv2.VideoCapture(0) # Get the first available camera.

def get_frame():
  return video.read()[1] # just get the frame.

server = RouterServer()
server.add_stream_from_callable("live_video", # name of the stream
                                30,  # How fast to call this function.
                                get_frame) # function to call.
server.run(block=True)  # if block=False, run in an internal thread. That's it!

In the PC acting as the client:

from one_liner.client import RouterClient
import cv2
import zmq

client = RouterClient()
client.configure_stream("live_video")

while True:
    try:
        timestamp, frame = client.get_stream("live_video")
    except zmq.Again:
        continue
    cv2.imshow("live_video from RouterServer", frame)
    cv2.waitKey(1) # Required short wait (1-ms).

That's it!

Application-Controlled Streaming

Instead of having RouterServer periodically call a function in a thread, it's also possible to have your native application send data via a handler function.

In the PC acting as the server:

from one_liner.server import RouterServer
import cv2

video = cv2.VideoCapture(0) # Get the first available camera.

server = RouterServer()
send_frame = server.get_broadcast_fn("live_video") # Create the handler function.
server.run()  # don't block.

while True:
    send_frame(video.read()[1])  # Read and send camera frames as fast as possible.

In the PC acting as the client--there are no changes! It's just the same example client code as before.

Streaming from an existing ZMQ Socket

It's possible to relay data from an existing ZMQ Socket as if it were another stream coming from a RouterServer.

TODO: see the relay_zmq_video_stream example folder for now.

Handling Received Data

There are two ways to handle received data: If you need every message, use 'queue', where messages will fill up an internal client-side FIFO buffer. If you only need to work off of the latest message, use 'cache', where only the most recently available message will be available on the client side. In the PC acting as the client:

client = RouterClient()
client.configure_stream("live_video", storage_type="cache")  # or 'queue' (default)

Note that configuring a stream also enables it.

Controlling Data Streams

Streamed data can be enabled or disabled such that messages from the server will or will not be sent to all clients for a given stream. To do this, simply enable or disable an existing stream by name:

client.enable_stream("live_video")  # the connected Router will not send this strea 
# ...
client.disable_stream("live_video")  # the connected Router will send the stream

Init from Config

It's possible to drive the creation of remote function calls and streams from a config dictionary provided that the object instance also exists in the instance dict.

Recall that, to create a RouterServer, we can optionally pass in a dictionary of object instances. Now we additionally pass in a config.

config = {
    "named_calls":
    {
        "set_axis_position":
            {
              "obj_name": "my_horn",
              "attr_name": "beep",
              "args": [],   # If args are empty, we can also omit this field.
              "kwargs": {}, # If kwargs are empty, we can also omit this field.
            }
    },
    "periodic_streams": {}
}

server = RouterServer(instances={"my_horn": my_horn}, config=config)
server.run()

Now, just like before, a connected client can client.call("my_horn", "beep") like before.

Restricting "Write" Access

RPCs exposed from a RouterClient support an optional annotation to mark whether the function implementation alters state. Marking a function as get hints that the function implements "read-only-like" behavior, while a set hints that the function applies "write-like" behavior that mutates the remote object's state. These get/set annotations only apply to named calls (i.e: not streams).

Write Token

Functions marked as set can only be accessed by one RouterClient with the "write token." After creating a RouterClient, you can request the write token with get_write_token(). Only one RouterClient instance at a time can hold the write token, giving it access to the set methods. Conceptually, write-token based write access implements mutual exclusion similar to a mutex lock.

server = RouterServer()
server.run()

client = RouterClient()
client.get_write_token()  # This client can now call functions marked with `set`
client.release_write_token()  # This client no longer has the write token

The write token can be acquired by force, causing any previously existing write token to be invalidated where further calls to set-annotated functions will raise a PermissionError.

client2 = RouterClient()
client2.get_write_token(force=True)  # This client booted other clients off.

Implementation Details

Bird's Eye View

High level, the RouterServer and RouterClient support two ways of sending and receiving data remotely.

It is also OK to connect multiple RouterClients to a single RouterServer like so:

Warning

There are no restrictions for connecting multiple clients at this level. Any restrictions or limitations on what functions can be called when multiple clients are connected needs to be applied at a higher level.

Streamer

Streaming is done by aggregating all calls of the same frequency and creating one thread per frequency. Streaming using threads simplifies the problem of scheduling when certain functions would be call using strategies like a priority queue. Because sockets are explicitly not threadsafe, each streamer thread instead creates its own PUB socket. To simplify connections on the client side, streams are aggregated together through a single socket such that the client only needs to know one address. Socket-to-socket communication is done using zmq's same-process shared memory implementation (inproc). Creating this proxy makes the system threadsafe.

Relaying data from another ZMQ Socket

It's also possible to stream data from an existing zmq socket (including another RouterServer). This is done with a zmq proxy.

By relaying data from a completely separate zmq socket, it is possible to cascade RouterServers.

Interfacing with external ZMQ Sockets

It should be possible to interface one-liner with other packages communicating with ZMQ with some knowledge about the underlying implementation details.

Both (1) the names of streams and (2) the obj_names of a remote procedure call are represented as ZMQ topics. While these topic names are represented as strings from the user side, they are converted to null-terimated strings under the hood before sending and receiving them. Doing so prevents topic collisions.

Example: topics new_frame and new_frame_2. When the data from topic new_frame starts with _2, that data would be erroneously routed to topic new_frame_2. To prevent this situation, all topics are represented as null-terminated strings.

Package/Project Management

This project utilizes uv to handle installing dependencies as well as setting up environments for this project. It replaces tool like pip, poetry, virtualenv, and conda.

This project also uses tox for orchestrating multiple testing environments that mimics the github actions CI/CD so that you can test the workflows locally on your machine before pushing changes.

Documentation

Environment Setup

To install with supplementary dependencies for creating local docs, run:

With uv:

uv sync --group docs

with pip:

pip install -e .[docs]

Building the Docs

To update the documentation and generate HTML files, from the project root directory, run

uv run mkdocs build

To view the docs locally, run:

uv run mkdocs serve

About

adapter pattern implemented as zmq client/server

Topics

Resources

Stars

3 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages