Skip to main content
Version: 3.0 (next)

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:

  1. Connection Lost: the client detects the dropped link/connection and surfaces it to the runtime's health check
  2. 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)
  3. Subscriptions Restored: receiver links for every enabled trigger are re-established on the fresh connection
  4. 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
ScenarioBehaviour
Brief network interruptionAutomatic reconnect + resubscribe within seconds; durable-queue messages sent during the gap are replayed
Broker restartSubscriptions automatically restored when the broker is back
Pipeline disabledSubscription is closed; messages accumulate on the (durable) queue and are processed after re-enable
MaestroHub restartAll triggers for enabled pipelines are restored on startup

Configuration Options​

Basic Information​

FieldTypeDescription
Node LabelString (Required)Display name for the node on the pipeline canvas
DescriptionString (Optional)Explains what this trigger initiates

Parameters​

ParameterTypeDefaultRequiredConstraintsDescription
Connection IDstring""Yes--AMQP connection profile to use.
Function IDstring""Yes--Subscribe function within the connection. Only Subscribe functions are listed.
Trigger Modeselect"always"Noalways / onChangealways: Trigger on every message. onChange: Only trigger when the payload differs from the previously received payload on the same address.
EnabledbooleantrueNo--Enable/disable the trigger. When disabled, the subscription is not opened.
Dedup Max Keysnumber1000If onChange1–10,000Maximum distinct addresses tracked for change detection. LRU eviction beyond the limit.
State TTLselectNo expiryNo1h / 6h / 12h / 24h / 72h / 168hShown 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.
Function Requirement

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.

onChange dedup scope is per-address

The trigger hashes the payload per address — two queues consumed on the same connection each get their own onChange slot.


Settings​

Execution Settings

SettingOptionsDefaultDescription
Timeout (seconds)numberPipeline defaultMaximum execution time for this node (1–600). Leave empty for pipeline default.
Retry on TimeoutPipeline Default / Enabled / DisabledPipeline DefaultWhether to retry the node if it times out.
Retry on FailPipeline Default / Enabled / DisabledPipeline DefaultWhether to retry on failure.
On ErrorPipeline Default / Stop Pipeline / Continue ExecutionPipeline DefaultBehavior 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:

FieldExpressionDescription
Message Payload$trigger.resultThe message body (parsed JSON object or raw value)
Address$trigger._metadata.addressThe queue/topic address the message was delivered from
Message ID$trigger._metadata.messageIdThe sender's message identifier (when set)
Correlation ID$trigger._metadata.correlationIdCorrelation identifier for request/response tracking (when set)
Content Type$trigger._metadata.contentTypeMIME content type stamped by the sender (when set)
Subject$trigger._metadata.subjectThe AMQP subject property (when set)
Reply-To$trigger._metadata.replyToReply 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.connectionIdThe AMQP connection profile used
Function ID$trigger._metadata.functionIdThe Subscribe function that received the message
Trigger Type$trigger._metadata.typeAlways amqp_trigger
Timestamp$trigger._metadata.timestampWhen MaestroHub received the message (RFC 3339, UTC, nanosecond precision)
Accessing Nested Payload Data

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.

RabbitMQ address format

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.