Skip to main content

Dead Simple Message Queue

What it does

Part mail room, part bulletin board, dsmq is a central location for sharing messages between processes, even when they are running on computers scattered around the world.

Its defining characteristic is bare-bones simplicity.

How to use it

Install

pip install dsmq

Create a dsmq server

As in src/dsmq/example_server.py

from dsmq.server import serve

serve(host="127.0.0.1", port=30008)

Connect a client to a dsmq server

As in src/dsmq/example_put_client.py

from dsmq.client import connect

mq = connect(host="127.0.0.1", port=12345)

Add a message to a queue

As in src/dsmq/example_put_client.py

topic = "greetings"
msg = "hello world!"
mq.put(topic, msg)

Read a message from a queue

As in src/dsmq/example_get_client.py

topic = "greetings"
msg = mq.get(topic)

Spin up and shut down a dsmq in its own process

A dsmq server doesn't come with a built-in way to shut itself down. It can be helpful to have it running in a separate process that can be managed

import multiprocessing as mp

p_mq = mp.Process(target=serve, args=(config.MQ_HOST, config.MQ_PORT))
p_mq.start()

p_mq.join()
# or 
p_mq.kill()
p_mq.close()

Demo

  1. Open 3 separate terminal windows.
  2. In the first, run src/dsmq/server.py as a script.
  3. In the second, run src/dsmq/example_put_client.py.
  4. In the third, run src/dsmq/example_get_client.py.

Alternatively, you can run them all at once with src/dsmq/demo.py.

How it works

Expected behavior and limitations

  • Many clients can read messages of the same topic. It is a one-to-many publication model.

  • A client will not be able to read any of the messages that were put into a queue before it connected.

  • A client will get the oldest message available on a requested topic. Queues are first-in-first-out.

  • Messages older than a certain age (typically 600 seconds) will be deleted from the queue.

  • Put and get operations are fairly quick--less than 100 $\mu$s of processing time plus any network latency--so it can comfortably handle requests at rates of hundreds of times per second. But if you have several clients reading and writing at 1 kHz or more, you may overload the queue.

  • The queue is backed by an in-memory SQLite database. If your message volumes get larger than your RAM, you will reach an out-of-memory condition.

API Reference

[source]

serve(host="127.0.0.1", port=30008)

Kicks off the mesage queue server. This process will be the central exchange for all incoming and outgoing messages.

  • host (str), IP address on which the server will be visible and
  • port (int), port. These will be used by all clients. Non-privileged ports are numbered 1024 and higher.

connect(host="127.0.0.1", port=30008)

Connects a client to an existing message queue server.

  • host (str), IP address of the server.
  • port (int), port on which the server is listening.
  • returns a DSMQClientSideConnection object.

DSMQClientSideConnection class

This is a convenience wrapper, to make the get() and put() functions easy to write and remember. It's under the hood only, not meant to be called directly.

put(topic, msg)

Puts msg into the queue named topic. If the queue doesn't exist yet, it is created.

  • msg (str), the content of the message.
  • topic (str), name of the message queue in which to put this message.

get(topic)

Get the oldest eligible message from the queue named topic. The client is only elgibile to receive messages that were added after it connected to the server.

  • topic (str)
  • returns str, the content of the message. If there was no eligble message in the topic, or the topic doesn't yet exist, returns "".

get_latest(topic)

Get the most recent eligible message from the queue named topic. All the messages older than that in the queue become ineligible and never get seen by the client.

  • topic (str)
  • returns str, the content of the message. If there was no eligble message in the topic, or the topic doesn't yet exist, returns "".

get_wait(topic)

A variant of get() that retries a few times until it gets a non-empty message. Adjust internal values _n_tries and _initial_retry to change how persistent it will be.

  • topic (str)
  • returns str, the content of the message. If there was no eligble message in the topic after the allotted number of tries, or the topic doesn't yet exist, returns "".

shutdown_server()

Gracefully shut down the server, through the client connection.

close()

Gracefully shut down the client connection.

Testing

Run all the tests in src/dsmq/tests/ with pytest, for example

uv run pytest

Performance characterization

Time typical operations on your system with the script at src/dsmq/tests/performance_suite.py

Release files for dsmq 1.4.1

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for dsmq 1.4.1
File Size Uploaded
dsmq-1.4.1.tar.gz 18.9 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for dsmq 1.4.1
File Interpreter ABI Platform
dsmq-1.4.1-py3-none-any.whl Python 3 none any Details

Total release size: 30.7 kB

Release files / dsmq-1.4.1.tar.gz

Download URL dsmq-1.4.1.tar.gz
Size 18.9 kB
Tags Source
SHA-256 checksum
How to use checksums
cc565d965b511be8a8f5591706858606137b3a7c78477361598c22e5e51b9adf
BLAKE2b-256 checksum
How to use checksums
2facbacc3a36b147f6780a01c309fdf41688281a0b159a0151031caaec7feca6
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via uv/0.5.1

Release files / dsmq-1.4.1-py3-none-any.whl

Download URL dsmq-1.4.1-py3-none-any.whl
Size 11.8 kB
Tags Python 3
SHA-256 checksum
How to use checksums
c87e48475ee7f301b37f6fdcfb1b5c35b3585c5bbca3b8bf093a18931da2eba3
BLAKE2b-256 checksum
How to use checksums
e1684a3ca1e73fad97db133d82306ee1193fa0dacd33997099585e4b8fdda8dc
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via uv/0.5.1

Release history Release notifications | RSS feed

This release

1.4.1 This release

2 release files

1.4.0

2 release files

1.3.1

2 release files

1.3.0

2 release files

1.2.4

2 release files

1.2.3

2 release files

1.2.2

2 release files

1.2.1

2 release files

1.2.0

2 release files

1.1.0

2 release files

1.0.0

2 release files

0.7.0

2 release files

0.6.0

2 release files

0.5.0

2 release files

0.4.0

2 release files

0.3.0

2 release files

0.2.0

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page