Skip to main content

codecov

talus (noun) - ta·​lus | ˈtā-ləs: a slope formed especially by an accumulation of rock debris; Occasional habitat of the pika.

A wrapper for connecting to RabbitMQ which constrains clients to a single purpose channel (producer or consumer) with healing for intermittent connectivity.

Features

  • Guided separation of connections for producers and consumers

  • Re-establish connections to the server when lost

  • Constrained interface to support simple produce / consume use cases for direct exchanges

Installation

pip install talus

Examples

Creating a consumer which listens on a queue, processes valid messages and publishes as part of processing

Uses default connection parameters and connection retryer expecting a rabbitmq server running in its default configuration.

from talus import DurableConsumer
from talus import DurableProducer
from talus import ConnectionRetryerFactory
from talus import ConsumerConnectionParameterFactory, ProducerConnectionParameterFactory
from talus import MessageProcessorBase
from talus import ConsumeMessageBase, PublishMessageBase, MessageBodyBase
from talus import Queue
from talus import Exchange
from talus import Binding
from typing import Type

##########################
# Consumer Configurations#
##########################
# Configure messages that will be consumed
class ConsumeMessageBody(MessageBodyBase):
    objectName: str
    bucket: str

class ConsumeMessage(ConsumeMessageBase):
    message_body_cls: Type[ConsumeMessageBody] = ConsumeMessageBody

# Configure the queue the messages should be consumed from
inbound_queue = Queue(name="inbound.q")


###########################
# Producer Configurations #
###########################
# Configure messages that will be produced
class ProducerMessageBody(MessageBodyBase):
    key: str
    code: str

class PublishMessage(PublishMessageBase):
    message_body_cls: Type[ProducerMessageBody] = ProducerMessageBody
    default_routing_key: str = "outbound.message.m"

# Configure the queues the message should be routed to
outbound_queue_one = Queue(name="outbound.one.q")
outbound_queue_two = Queue(name="outbound.two.q")


# Configure the exchange and queue bindings for publishing (Publish Message -> Outbound Queues)
publish_exchange = Exchange(name="outbound.exchange") # Direct exchange by default
bindings = [Binding(queue=outbound_queue_one, message=PublishMessage),
            Binding(queue=outbound_queue_two, message=PublishMessage)] # publishing PublishMessage will route to both queues.


############################
# Processor Configurations #
############################

# Configure a message processor to handle the consumed messages
class MessageProcessor(MessageProcessorBase):
    def process_message(self, message: ConsumeMessage):
        print(message)
        outbound_message = PublishMessage(
            body=ProducerMessageBody(
                key=message.body.objectName,
                code="newBucket",
                conversationId=message.body.conversationId,
            )
        )  # crosswalk the values from the consumed message to the produced message
        self.producer.publish(outbound_message)
        print(outbound_message)


# Actually Connect and run the consumer
def main():
    """Starts a listener which will consume messages from the inbound queue and publish messages to the outbound queues."""
    with DurableProducer(
        queue_bindings=bindings,
        publish_exchange=publish_exchange,
        connection_parameters=ProducerConnectionParameterFactory(),
        connection_retryer=ConnectionRetryerFactory(),
    ) as producer:
        with DurableConsumer(
            consume_queue=inbound_queue,
            connection_parameters=ConsumerConnectionParameterFactory(),
            connection_retryer=ConnectionRetryerFactory(),
        ) as consumer:
            message_processor = MessageProcessor(message_cls=ConsumeMessage, producer=producer)
            consumer.listen(message_processor)


if __name__ == "__main__":
    # First message to consume
    class InitialMessage(PublishMessageBase):
        message_body_cls: Type[
            ConsumeMessageBody] = ConsumeMessageBody
        default_routing_key: str = "inbound.message.m"

    initial_message_bindings = [Binding(queue=inbound_queue, message=InitialMessage)]

    with DurableProducer(
            queue_bindings=initial_message_bindings,
            publish_exchange=publish_exchange,
            connection_parameters=ProducerConnectionParameterFactory(),
            connection_retryer=ConnectionRetryerFactory(),
    ) as producer:
        producer.publish(InitialMessage(body={"objectName": "object", "bucket": "bucket"}))
    # Consume the message and process it
    main()

Metadata

Release files for talus 1.3.4

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

Source distribution (sdist)

Source distribution for talus 1.3.4
File Size Uploaded
talus-1.3.4.tar.gz 22.1 kB Details

Built distribution (wheel)

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

Total release size: 47.6 kB

Release files / talus-1.3.4.tar.gz

Download URL talus-1.3.4.tar.gz
Size 22.1 kB
Tags Source
SHA-256 checksum
How to use checksums
47637da515f6a75664598ca639d3fbb28dc1d63e1215d68a7378d527fc112c27
BLAKE2b-256 checksum
How to use checksums
0cf775a71111ffe2215053b283ce8f33c53a67ec12704be028f38d9263f8c780
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/6.1.0 CPython/3.11.13

Release files / talus-1.3.4-py3-none-any.whl

Download URL talus-1.3.4-py3-none-any.whl
Size 25.5 kB
Tags Python 3
SHA-256 checksum
How to use checksums
81264e0a318a3c25feafb7bb71a48dcad0c83998aeb2122f8a9f4e475fa005ed
BLAKE2b-256 checksum
How to use checksums
553ad552b02826874059aa618028a2fc9e77fd4c2ed48c93a5ef636aeeeab29d
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/6.1.0 CPython/3.11.13

Release history Release notifications | RSS feed

This release

1.3.4 This release

2 release files

1.3.3

2 release files

1.3.2

2 release files

1.3.1

2 release files

1.3.0

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.2.1

1 release file

0.2.0

1 release file

0.1.9

1 release file

0.1.8

1 release file

0.1.7

1 release file

0.1.6

1 release file

0.1.5

1 release file

0.1.4

1 release file

0.1.2

1 release file

0.1.1

1 release file

0.1.0

1 release file

0.0.2

1 release file

0.0.1

1 release file

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