Skip to main content
Pre-release

This release is a pre-release and may not be stable for production use.

Neon MQ Connector

The Neon MQ Connector is an MQ interface for microservices.

Configuration

By default, this package will use ovos-config to read default configuration. In general, configuration should be passed to the MQConnector object at init.

Legacy Configuration

A global configuration for the MQ Connector may be specified at ~/.config/neon/mq_config.json. This configuration file may contain the following keys:

  • server: The hostname or IP address of the MQ server to connect to. If left blank, this defaults to "localhost"
  • port: The port used by the MQ server. If left blank, this defaults to 5672
  • users: A mapping of service names to credentials. Note that not all users will have permissions required to access each service.
{
  "server": "localhost",
  "port": 5672,
  "users": {
    "<service_name>": {
      "username": "<username>",
      "password": "<password>"
    }
  }
}

Services

The MQConnector class should be extended by a class providing some specific service. Service classes will specify the following parameters.

  • service_name: Name of the service, used to identify credentials in configuration
  • vhost: Virtual host to connect to; messages are all constrained to this namespace.
  • consumers: Dict of names to ConsumerThread objects. A ConsumerThread will accept a connection to a particular connection, a queue, and a callback_func
    • connection: MQ connection to the vhost specified above.
    • queue: Queue to monitor within the vhost. A vhost may handle multiple queues.
    • callback_func: Function to call when a message arrives in the queue

Callback Functions

A callback function should have the following signature:

def handle_api_input(self,
                     channel: pika.channel.Channel,
                     method: pika.spec.Basic.Return,
                     properties: pika.spec.BasicProperties,
                     body: bytes):
    """
        Handles input requests from MQ to Neon API

        :param channel: MQ channel object (pika.channel.Channel)
        :param method: MQ return method (pika.spec.Basic.Return)
        :param properties: MQ properties (pika.spec.BasicProperties)
        :param body: request body (bytes)
    """

Generally, body should be decoded into a dict, and that dict should contain message_id. The message_id should be included in the body of any response to associate the response to the request. A response may be sent via:

 channel.queue_declare(queue='<queue>')

 channel.basic_publish(exchange='',
                       routing_key='<queue>',
                       body=<data>,
                       properties=pika.BasicProperties(expiration='1000')
                       )

Where <queue> is the queue to which the response will be published, and data is a bytes response (generally a base64-encoded dict).

Multi-part Responses

A callback function may choose to publish multiple response messages so the client may receive partial responses as they are being generated. If multiple responses will be returned, the following requirements must be met:

  • Each response must be a dict with part and is_final keys.
  • part is defined as a non-negative integer (the first response will specify 0).
  • The final response must specify is_final=True
  • The final response must be complete and MUST NOT require the client to handle partial responses

Client Requests

Most client applications will interact with services via send_mq_request. This function will return a dict response to the input message.

Multi-part Responses

A caller may optionally include a stream_callback argument which may receive partial responses if supported by the service generating the response. The stream_callback will always be called with the final result that is returned by send_mq_request. Keep in mind that the timeout param passed to send_mq_request applies to the full response, so the timeout value should reflect the longest time it will take for a final response to be generated, plus some margin.

Note: Responses may be received out of order, so the client is responsible for monitoring the part and is_final fields as needed

Asynchronous Consumers

By default, async-based consumers handling based on pika.SelectConnection will be used

Override use of async consumers

There are a few methods to disable use of async consumers/subscribers.

  1. To disable async consumers for a particular class/object, set the class-attribute async_consumers_enabled to False:

    from neon_mq_connector import MQConnector
    
    class MQConnectorChild(MQConnector):
       async_consumers_enabled = False
    
  2. To disable the use of async consumers at runtime, set the MQ_ASYNC_CONSUMERS envvar to False

    export MQ_ASYNC_CONSUMERS=false
    

Metadata

Release files for neon-mq-connector 0.10.1a3

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

Source distribution (sdist)

Source distribution for neon-mq-connector 0.10.1a3
File Size Uploaded
neon_mq_connector-0.10.1a3.tar.gz 34.2 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for neon-mq-connector 0.10.1a3
File Interpreter ABI Platform
neon_mq_connector-0.10.1a3-py3-none-any.whl Python 3 none any Details

Total release size: 89.4 kB

Release files / neon_mq_connector-0.10.1a3.tar.gz

Download URL neon_mq_connector-0.10.1a3.tar.gz
Size 34.2 kB
Tags Source
SHA-256 checksum
How to use checksums
f924575b05e34993de5f37c664b27ac681af11730a16761a4a56c3dfc5b3f46c
BLAKE2b-256 checksum
How to use checksums
5a436340080041d1c1643ce85ec448877cf5ea5cecbe867b064da880a8477122
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.13.14

Release files / neon_mq_connector-0.10.1a3-py3-none-any.whl

Download URL neon_mq_connector-0.10.1a3-py3-none-any.whl
Size 55.2 kB
Tags Python 3
SHA-256 checksum
How to use checksums
7b7889279b775fd6fa1c856896ba2148fbb609acb9f92423337eaae4e4c07efc
BLAKE2b-256 checksum
How to use checksums
e6765a46c041cadd58f6dff92e2298592a5ce26e47e5e51fe785ca38c13c1023
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.13.14

Release history Release notifications | RSS feed

This release

0.10.1a3 This release

2 release files

0.10.0

2 release files

0.9.0

2 release files

0.8.0

2 release files

0.7.1

2 release files

0.7.0

1 release file

0.6.0

1 release file

0.5.1

1 release file

0.5.0

1 release file

0.4.0

1 release file

0.3.1

1 release file

0.3.0

1 release file

0.2.0

1 release file

0.1.2

1 release file

0.1.0

1 release file

0.0.4

2 release files

0.0.3

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