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 to5672users: 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 configurationvhost: Virtual host to connect to; messages are all constrained to this namespace.consumers: Dict of names toConsumerThreadobjects. AConsumerThreadwill accept a connection to a particularconnection, aqueue, and acallback_funcconnection: MQ connection to thevhostspecified above.queue: Queue to monitor within thevhost. Avhostmay handle multiple queues.callback_func: Function to call when a message arrives in thequeue
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
partandis_finalkeys. partis defined as a non-negative integer (the first response will specify0).- 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
partandis_finalfields 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.
-
To disable async consumers for a particular class/object, set the class-attribute
async_consumers_enabledtoFalse:from neon_mq_connector import MQConnector class MQConnectorChild(MQConnector): async_consumers_enabled = False
-
To disable the use of async consumers at runtime, set the
MQ_ASYNC_CONSUMERSenvvar toFalseexport 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)
| File | Size | Uploaded | |
|---|---|---|---|
| neon_mq_connector-0.10.1a3.tar.gz | 34.2 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| 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
|