Route SaaS platform events from SQS to different destinations based on tenant identity and message priority, using multi-condition rules with mixed transport targets.
A SaaS platform receives events from an SQS queue. Each event carries:
x-tenant header identifying the tenant (e.g., enterprise, startup-42)$.priority field (integer 1–10, where 10 is most critical)The routing rules are:
enterprise – dedicated MQTT topic enterprise/criticaltenants/{x-tenant}/events (address template)This scenario demonstrates multi-condition AND logic, address templates combined with resolver rules, mixed transport routing (SQS to MQTT + SQS), and default fallback bindings.
flowchart LR
subgraph AWS
Q["SQS Queue\nplatform-events"]
HQ["SQS Queue\nhigh-priority"]
CQ["SQS Queue\ncatch-all"]
end
subgraph GoBridge
R[Receiver\nsqs-in]
Proc["Processor Chain\ntenant-validator\n+ resolver"]
Route[Route\npriority-dispatch]
Res["Resolver\ntype: rules\nmulti-condition"]
end
subgraph MQTT Broker
T1["enterprise/critical"]
T2["tenants/startup-42/events"]
T3["tenants/acme-corp/events"]
end
Q --> R
R --> Route
Route --> Proc
Proc --> Res
Res -->|"priority>8 AND enterprise"| T1
Res -->|"priority>8"| HQ
Res -->|"x-tenant exists"| T2
Res -->|"x-tenant exists"| T3
Res -->|"default (no match)"| CQ
style Route fill:#f96,stroke:#333
style Res fill:#ff9,stroke:#333
style GoBridge fill:#eef,stroke:#333
bridge:
id: tenant-priority-router
sessions:
- id: mqtt-conn
transport: mqtt
options:
session:
broker_url: tcp://mqtt.platform.local:1883
client_id: tenant-router-01
keep_alive: 60
receivers:
- id: sqs-in
transport: sqs
options:
queue_url: https://sqs.us-west-1.amazonaws.com/123456789/platform-events
region: us-west-1
wait_time_seconds: 20
max_messages: 10
senders:
- id: mqtt-enterprise
session_id: mqtt-conn
options:
sender:
qos: 1
- id: mqtt-tenants
session_id: mqtt-conn
options:
sender:
qos: 1
- id: sqs-high-priority
transport: sqs
options:
queue_url: https://sqs.us-west-1.amazonaws.com/123456789/high-priority
region: us-west-1
batch_size: 10
- id: sqs-catch-all
transport: sqs
options:
queue_url: https://sqs.us-west-1.amazonaws.com/123456789/catch-all
region: us-west-1
bindings:
- id: to-enterprise-critical
sender_id: mqtt-enterprise
address: "enterprise/critical"
- id: to-high-priority
sender_id: sqs-high-priority
address: high-priority
- id: to-tenant-topic
sender_id: mqtt-tenants
address: "tenants/{x-tenant}/events"
- id: to-catch-all
sender_id: sqs-catch-all
address: catch-all
routes:
- id: priority-dispatch
receiver_id: sqs-in
delivery_mode: direct_hold
dispatch_mode: single
bindings:
- to-enterprise-critical
- to-high-priority
- to-tenant-topic
- to-catch-all
processors: [tenant-validator]
resolver:
type: rules
default_binding: to-catch-all
rules:
- binding_id: to-enterprise-critical
match:
- field: $.priority
operator: gt
value: 8
- field: header.x-tenant
operator: eq
value: "enterprise"
- binding_id: to-high-priority
match:
- field: $.priority
operator: gt
value: 8
- binding_id: to-tenant-topic
match:
- field: header.x-tenant
operator: exists
value: true
policy:
max_in_flight: 200
on_permanent_failure: dlq
stores:
outbox: { type: memory, options: { acknowledge_volatile: true } }
dlq: { type: memory, options: { acknowledge_volatile: true } }
Each rule’s match array uses AND logic – all conditions must be true for the rule to select its binding.
Rule 1 has two conditions:
- binding_id: to-enterprise-critical
match:
- field: $.priority
operator: gt
value: 8
- field: header.x-tenant
operator: eq
value: "enterprise"
Both $.priority > 8 AND header.x-tenant == "enterprise" must be true. A message from tenant enterprise with priority 7 does not match this rule. A message from tenant startup-42 with priority 9 does not match either – both conditions are required.
Rules are evaluated top to bottom. The ordering matters:
x-tenant header that did not match the priority rules above.x-tenant header or with malformed payloads where $.priority extraction fails.If rule 2 were placed before rule 1, enterprise critical messages would be consumed by the shared high-priority queue and never reach the dedicated enterprise topic.
Binding to-tenant-topic uses an address template:
- id: to-tenant-topic
sender_id: mqtt-tenants
address: "tenants/{x-tenant}/events"
When the resolver selects this binding, the RenderAddress function replaces {x-tenant} with the value of the x-tenant header from the envelope. A message with header x-tenant: startup-42 is published to MQTT topic tenants/startup-42/events.
This combines two GoBridge features:
Rule 3 ensures the x-tenant header exists before the template is rendered. If the header were missing, RenderAddress would return an error, and the message would fail. The exists condition acts as a guard.
This route sends to both MQTT and SQS senders from a single receiver:
| Binding | Transport | Destination |
|---|---|---|
to-enterprise-critical |
MQTT | enterprise/critical (static topic) |
to-high-priority |
SQS | high-priority queue (static) |
to-tenant-topic |
MQTT | tenants/{x-tenant}/events (dynamic topic) |
to-catch-all |
SQS | catch-all queue (static) |
The MQTT senders share a session (mqtt-conn). The SQS senders each have their own connection because SQS is stateless at the transport level.
The route uses direct_hold for simplicity. In production, consider the durability requirements for each binding:
| Binding | Recommended Mode | Rationale |
|---|---|---|
to-enterprise-critical |
shared_outbox |
Critical messages must not be lost. Outbox provides at-least-once delivery even if MQTT is temporarily unavailable. |
to-high-priority |
direct_hold |
SQS is durable. The SQS sender uses SendMessage which is itself durable. |
to-tenant-topic |
shared_outbox |
Per-tenant MQTT topics benefit from outbox durability. |
to-catch-all |
direct_hold |
Catch-all is best-effort. Direct delivery is sufficient. |
A single route can only have one delivery_mode. To use different modes per destination, split into multiple routes or accept the trade-off of a single mode. When shared_outbox is configured, all bindings benefit from outbox durability, at the cost of additional latency and storage.
The processors: [tenant-validator] runs before the resolver evaluates rules. This ensures:
sequenceDiagram
participant SQS as SQS Queue
participant R as Receiver (sqs-in)
participant P as Processor (tenant-validator)
participant Route as RouteRunner
participant Res as RuleResolver
participant MQTT as MQTT Sender (mqtt-enterprise)
participant Broker as MQTT Broker
SQS->>R: ReceiveMessage (x-tenant: enterprise)
R->>Route: Envelope{Headers: {x-tenant: enterprise}, Payload: {priority: 9}}
Route->>P: Process(envelope)
P->>P: Validate tenant "enterprise" -- OK
P-->>Route: envelope (unchanged)
Route->>Res: Resolve(envelope)
Note over Res: Rule 1: $.priority > 8 -- YES
Note over Res: Rule 1: header.x-tenant == "enterprise" -- YES
Note over Res: Rule 1: ALL conditions match
Res-->>Route: DispatchPlan{BindingID: "to-enterprise-critical", Address: "enterprise/critical"}
Route->>MQTT: Send(envelope, subject="enterprise/critical")
MQTT->>Broker: PUBLISH enterprise/critical (QoS 1)
Broker-->>MQTT: PUBACK
MQTT-->>Route: OK
Route->>R: ACK (DeleteMessage)
Register the tenant validator programmatically:
package main
import (
"context"
"log/slog"
"github.com/mariotoffia/gobridge/bridge"
cfgparser "github.com/mariotoffia/gobridge/config/parser"
"github.com/mariotoffia/gobridge/ports"
"github.com/mariotoffia/gobridge/adapters/mqtt/transport/paho"
"github.com/mariotoffia/gobridge/adapters/aws/transport/sqs"
"github.com/mariotoffia/gobridge/processors/tenant"
)
func main() {
logger := slog.Default()
// Register each linked adapter's config decoder; ParseFile requires a
// non-nil registry. myTenantValidator is a user-supplied hook.
reg := ports.NewRegistry()
_ = paho.Register(reg)
_ = sqs.Register(reg)
cfg, _ := cfgparser.ParseFile("bridge.yaml", cfgparser.FormatAuto, reg)
tenantProc, _ := tenant.New(tenant.Config{
Name: "tenant-validator",
TenantHeader: "x-tenant",
RequireTenant: false, // allow catch-all for messages without tenant
}, tenant.WithValidator(myTenantValidator))
sup, _ := bridge.NewBuilder(cfg, bridge.WithLogger(logger)).
RegisterTransportFactory("mqtt", paho.NewFactory(logger)).
RegisterTransportFactory("sqs", sqs.NewFactory(logger)).
RegisterProcessor("tenant-validator", tenantProc).
Build(context.Background())
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
sup.Start(ctx)
// ... wait for signal ...
sup.Stop(ctx)
}
Note that RequireTenant is set to false. Messages without the x-tenant header are allowed through the processor and fall to the catch-all binding via the resolver’s default_binding. If RequireTenant were true, messages without the header would be rejected by the processor and never reach the resolver.
The three field types used in this scenario:
| Pattern | Example | Resolves To |
|---|---|---|
$.priority |
JSON payload path | Envelope.Payload unmarshaled, then priority key |
header.x-tenant |
Header lookup | Envelope.Headers["x-tenant"] |
subject |
Envelope subject | Envelope.Subject (not used in this scenario) |
When multiple rules reference $. paths, the JSON payload is parsed only once per envelope. The evalContext caches the parsed map[string]any and reuses it across all condition evaluations. In this scenario, rules 1 and 2 both access $.priority, but the payload is unmarshaled once.
The gt operator converts both the extracted payload value and the condition value to float64 before comparing. JSON numbers are parsed via json.Number to avoid precision loss. String-encoded numbers (e.g., "9") are also supported.
If a condition fails to evaluate (e.g., payload is not valid JSON, field path does not exist), the condition is treated as non-matching. The rule is skipped, and evaluation continues to the next rule. This means:
$.priority: "not-a-number" fails the gt comparison, skips rules 1 and 2, and falls through to rule 3 or the default.x-tenant header, matches rule 3.If routing is based solely on the tenant header without priority logic, use header_map instead of rules:
resolver:
type: header_map
header_key: x-tenant
default_binding: to-catch-all
header_map:
enterprise: to-enterprise-critical
startup-42: to-tenant-a
acme-corp: to-tenant-b
This is simpler but cannot express multi-condition logic or numeric comparisons.
Protect downstream MQTT from a flood of events from a single tenant:
routes:
- id: priority-dispatch
processors: [tenant-validator, tenant-cb]
cbProc := circuitbreaker.New("tenant-cb", circuitbreaker.Config{
FailureThreshold: 5,
SuccessThreshold: 3,
ResetTimeout: 30 * time.Second,
HalfOpenMaxProbes: 2,
}, circuitbreaker.WithKeyExtractor(circuitbreaker.HeaderKey("x-tenant")))
Each tenant gets an independent circuit breaker. If tenant startup-42 causes 5 consecutive failures, only that tenant’s breaker opens. Enterprise messages continue flowing.
If the tenant ID is in the JSON payload instead of a header, use a transform processor to extract it before the resolver runs:
extractProc, _ := transform.New(transform.Config{
Name: "extract-tenant",
Mappings: []transform.FieldMapping{
{Source: "$.metadata.tenant_id", Target: "header.x-tenant"},
},
})
routes:
- id: priority-dispatch
processors: [extract-tenant, tenant-validator]
The transform runs first, copies $.metadata.tenant_id into the x-tenant header, and then the tenant validator and resolver can use the header as normal.
Extend the rules to support three priority tiers:
resolver:
type: rules
default_binding: to-catch-all
rules:
- binding_id: to-enterprise-critical
match:
- field: $.priority
operator: gt
value: 8
- field: header.x-tenant
operator: eq
value: "enterprise"
- binding_id: to-critical
match:
- field: $.priority
operator: gt
value: 8
- binding_id: to-elevated
match:
- field: $.priority
operator: gt
value: 5
- binding_id: to-tenant-topic
match:
- field: header.x-tenant
operator: exists
value: true
Messages with priority 6–8 go to the to-elevated binding. Messages with priority 1–5 that have a tenant header go to the per-tenant topic. The ordering ensures more specific rules are evaluated first.
binding_id values in rules must reference bindings in the route’s bindings array. The default_binding must also be listed. Invalid references cause a config validation error at startup.CompileMatchRules is called once at startup. Regex patterns are compiled and cached. Condition evaluators are safe for concurrent use across goroutines.RenderAddress uses single-pass substitution. Rendered values are never re-expanded, preventing header-value injection. Missing placeholders cause an error, not a silent drop.