AMQP Trigger Node
Overview
The AMQP Trigger Node automatically initiates MaestroHub pipelines when messages arrive on a subscribed AMQP 1.0 queue or topic. Unlike the AMQP Publish connector node which sends messages from inside a running pipeline, the AMQP Trigger starts new pipeline executions in response to incoming messages — enabling event-driven automation against any standard AMQP 1.0 broker (ActiveMQ Artemis, Azure Service Bus, Solace, IBM MQ, RabbitMQ 4.x).
Core Functionality
What It Does
AMQP Trigger enables real-time, event-driven pipeline execution by:
1. Event-Driven Pipeline Execution Start pipelines automatically when messages arrive on the address configured in the underlying Subscribe function, without manual intervention or polling.
2. Automatic Subscription Management The connector's subscription manager opens and tears down receiver links automatically as pipelines are enabled, disabled, or updated — no manual subscription bookkeeping.
3. At-Least-Once Settlement Each message is accepted (removed from the queue) only after the pipeline has durably captured it. A transient dispatch failure releases the message for broker redelivery; a permanent failure rejects it so the broker dead-letters it. Messages queued while the pipeline is disabled or MaestroHub is down are held by the broker (durable queues) and replayed on reconnect — nothing is lost.
4. Competing-Consumer Load Balancing AMQP 1.0 queues natively support competing consumers: multiple replicas subscribing to the same queue split deliveries between them, scaling consumption horizontally with at-least-once semantics.
5. Credit-Based Flow Control The Subscribe function's Prefetch bounds how many unsettled messages are in flight at once, so a slow pipeline backpressures the broker instead of buffering unboundedly.
6. Message Payload and Metadata Passthrough
The message body and its properties (address, messageId, correlationId, contentType, subject, replyTo, application properties) are passed to downstream nodes via the $trigger variable.
Reconnection Handling
When an AMQP connection is lost and restored:
- Connection Lost: the client detects the dropped link/connection and surfaces it to the runtime's health check
- Reconnect Loop: the platform's reconnect worker re-dials with backoff (a session-level broker error is treated the same way — the connection self-heals within seconds)
- Subscriptions Restored: receiver links for every enabled trigger are re-established on the fresh connection
- Replay: messages the durable queue held during the outage are delivered once the subscription is back — at-least-once, so downstream pipelines should be idempotent
| Scenario | Behaviour |
|---|---|
| Brief network interruption | Automatic reconnect + resubscribe within seconds; durable-queue messages sent during the gap are replayed |
| Broker restart | Subscriptions automatically restored when the broker is back |
| Pipeline disabled | Subscription is closed; messages accumulate on the (durable) queue and are processed after re-enable |
| MaestroHub restart | All triggers for enabled pipelines are restored on startup |
Configuration Options
Basic Information
| Field | Type | Description |
|---|---|---|
| Node Label | String (Required) | Display name for the node on the pipeline canvas |
| Description | String (Optional) | Explains what this trigger initiates |
Parameters
| Parameter | Type | Default | Required | Constraints | Description |
|---|---|---|---|---|---|
| Connection ID | string | "" | Yes | -- | AMQP connection profile to use. |
| Function ID | string | "" | Yes | -- | Subscribe function within the connection. Only Subscribe functions are listed. |
| Trigger Mode | select | "always" | No | always / onChange | always: Trigger on every message. onChange: Only trigger when the payload differs from the previously received payload on the same address. |
| Enabled | boolean | true | No | -- | Enable/disable the trigger. When disabled, the subscription is not opened. |
| Dedup Max Keys | number | 1000 | If onChange | 1–10,000 | Maximum distinct addresses tracked for change detection. LRU eviction beyond the limit. |
| State TTL | select | No expiry | No | 1h / 6h / 12h / 24h / 72h / 168h | Shown in the panel, and selectable only in Enterprise, but the trigger does not apply it. Change-detection state is kept in memory and resets when MaestroHub restarts, so the first message after a restart always fires. |
The selected function must be an AMQP Subscribe function. Publish, Request, Receive, and Browse functions cannot drive a trigger node. The address and prefetch live on the Subscribe function, not on the trigger node.
The trigger hashes the payload per address — two queues consumed on the same connection each get their own onChange slot.
Settings
Execution Settings
| Setting | Options | Default | Description |
|---|---|---|---|
| Timeout (seconds) | number | Pipeline default | Maximum execution time for this node (1–600). Leave empty for pipeline default. |
| Retry on Timeout | Pipeline Default / Enabled / Disabled | Pipeline Default | Whether to retry the node if it times out. |
| Retry on Fail | Pipeline Default / Enabled / Disabled | Pipeline Default | Whether to retry on failure. |
| On Error | Pipeline Default / Stop Pipeline / Continue Execution | Pipeline Default | Behavior when the node fails after all retries. |
Output Data Structure
When an AMQP message triggers pipeline execution, the following data is available to downstream nodes via the $trigger variable.
Output Format
{
"_metadata": {
"type": "amqp_trigger",
"connectionId": "bf29be94-fc0a-4dc4-8e5c-092f1b74eb4b",
"functionId": "aef374c3-aa2b-454e-aabc-5657faac5950",
"protocol": "amqp",
"address": "/queues/orders",
"messageId": "9f6d2c…",
"correlationId": "ORD-12345",
"contentType": "application/json",
"subject": "order-created",
"timestamp": "2026-08-19T10:30:00.123456789Z"
},
"result": {
"orderId": "ORD-12345",
"customer": "acme",
"amount": 99.99
}
}
Accessing Message Data
In downstream nodes, use the $trigger variable:
| Field | Expression | Description |
|---|---|---|
| Message Payload | $trigger.result | The message body (parsed JSON object or raw value) |
| Address | $trigger._metadata.address | The queue/topic address the message was delivered from |
| Message ID | $trigger._metadata.messageId | The sender's message identifier (when set) |
| Correlation ID | $trigger._metadata.correlationId | Correlation identifier for request/response tracking (when set) |
| Content Type | $trigger._metadata.contentType | MIME content type stamped by the sender (when set) |
| Subject | $trigger._metadata.subject | The AMQP subject property (when set) |
| Reply-To | $trigger._metadata.replyTo | Reply address set by the sender (present for request/reply senders) |
| Application Properties | $trigger._metadata.attr_<key> | The sender's application properties, prefixed with attr_ — the same prefix every connector uses for the message's own attributes. A key with a hyphen needs bracket syntax |
| Connection ID | $trigger._metadata.connectionId | The AMQP connection profile used |
| Function ID | $trigger._metadata.functionId | The Subscribe function that received the message |
| Trigger Type | $trigger._metadata.type | Always amqp_trigger |
| Timestamp | $trigger._metadata.timestamp | When MaestroHub received the message (RFC 3339, UTC, nanosecond precision) |
If your messages contain JSON, access nested fields directly — $trigger.result.orderId.
Validation Rules
Connection ID
- Must be provided and non-empty
- Must reference a valid AMQP connection profile
- Error: "Connection is required"
Function ID
- Must be provided and non-empty
- Must reference a valid AMQP Subscribe function belonging to the selected connection
- Error: "Subscribe Function is required"
Usage Examples
ERP Event Ingest
Scenario: React to business events an ERP publishes to an Azure Service Bus queue.
Configuration:
- Label: ERP Order Events
- Connection: Service Bus (amqps endpoint)
- Function: Subscribe to
orders-out, prefetch 10 - Trigger Mode: always
Downstream Processing: parse $trigger.result, enrich via a REST call, publish the normalized event into the Unified Namespace.
Load-Balanced Work Queue
Scenario: Distribute jobs on an Artemis queue across MaestroHub replicas.
Configuration: the same pipeline enabled on multiple replicas, all subscribing to the same queue — AMQP's competing-consumer semantics deliver each message to exactly one replica, and prefetch bounds each replica's in-flight work.
Change-Only State Relay
Scenario: A device publishes its full state on every heartbeat, but downstream should only run when the state actually changes.
Configuration: Trigger Mode = onChange — repeated identical payloads on the same address are suppressed; any changed payload fires.
On RabbitMQ 4.x the Subscribe function's address must use the v2 grammar (/queues/<name>), and the queue must already exist — see the AMQP connector guide.