No broker required

Add integration pipelines without a broker

Layer 2 is the set of Enterprise Integration Pattern models in the qubit-core library, run inside your own Go service by the SDK's camel runtime. It is not a hosted Priostack feature: nothing is deployed to priostack.com. No Kafka, no RabbitMQ.

Start building → Read the docs →

9 integration patterns, built in

Every pattern is a typed Go struct in ideaswave.com/qubit/core/pkg/integration. The runnable form is the SDK's camel runtime (ideaswave.com/qubit/sdk/camel), which reads your routes from a Camel Spring-DSL XML document. Today it runs the filter and aggregate steps; routers, translators and splitters are declared in the model and carried out by your own code.

📡

MessageChannel

A named conduit, point-to-point or publish-subscribe

🔌

MessageEndpoint

Connects a process or service to a channel, inbound or outbound

🔀

MessageRouter

Route messages to different channels via FEEL expressions

🔍

MessageFilter

Drops a message whose FEEL predicate is false (the camel runtime answers 204)

📦

Aggregator

Holds messages grouped by a correlation expression until a completion condition releases them as one batch

✂️

Splitter

Split a list variable into parallel messages for fan-out

🔄

MessageTranslator

Map and transform message fields using FEEL expressions

⛓️

Pipeline

An ordered chain of EIP components (Pipes and Filters)

🎯

CorrelationContext

Matches an incoming message to one waiting process instance by a FEEL key

5 Idempotency Mechanisms

Where a duplicate is stopped, and what keeps the record of it. None of the five needs an external broker.

Mechanism Scope Storage When to use
CorrelationContext Per context In memory A match removes the entry, so a second message with the same key resumes nothing
Job key Per job Priostack engine A job completes once; a repeated completion answers 404, so make handlers idempotent on the job key
Database upsert Application Your database Business-level idempotency
MessageFilter Route None Drop messages that fail a FEEL predicate before they reach business logic
Aggregator Route In memory A group is released once, when its completion condition fires

Works exactly as you expect

The model is Go structs; the runnable routes are one Camel XML document plus the Go handlers (beans) you register. The three tabs below form one program.

<beans xmlns="http://www.springframework.org/schema/beans">
<camelContext xmlns="http://camel.apache.org/schema/spring" id="loans">
  <!-- Filter: only gold-tier applications continue; the others answer 204 -->
  <route id="gold">
    <from uri="direct:gold-in"/>
    <filter>
      <simple>tier = "gold"</simple>
      <to uri="direct:svc"/>
    </filter>
  </route>
  <!-- Aggregate: hold a portfolio's applications (202) until the last one -->
  <route id="batch">
    <from uri="direct:batch-in"/>
    <aggregate>
      <correlationExpression><simple>portfolio_id</simple></correlationExpression>
      <completionPredicate><simple>last = true</simple></completionPredicate>
      <to uri="direct:svc"/>
    </aggregate>
  </route>
  <!-- REST routes enter a component through direct:<route id>-in -->
  <route id="loans.apply">
    <from uri="rest:post:/api/loans"/>
    <to uri="direct:gold-in"/>
    <to uri="bean:loans?method=apply"/>
  </route>
  <route id="loans.batch">
    <from uri="rest:post:/api/loans/batch"/>
    <to uri="direct:batch-in"/>
    <to uri="bean:loans?method=batch"/>
  </route>
</camelContext>
</beans>
package main

import (
    _ "embed"
    "log"
    "net/http"

    "ideaswave.com/qubit/sdk/camel"
)

//go:embed routes.xml
var routesXML []byte

func main() {
    rt, err := camel.New(routesXML, nil) // nil: no JSON-Schema validation
    if err != nil {
        log.Fatal(err)
    }
    rt.RegisterBean("loans.apply", applyHandler) // bean:loans?method=apply
    rt.RegisterBean("loans.batch", batchHandler) // bean:loans?method=batch

    mux := http.NewServeMux()
    guard := func(h http.HandlerFunc) http.HandlerFunc { return h } // your auth middleware
    if unbound := rt.Bind(mux, guard); len(unbound) > 0 {
        log.Printf("routes with no bean (they answer 501): %v", unbound)
    }
    log.Fatal(http.ListenAndServe(":8080", mux))
}

func applyHandler(w http.ResponseWriter, r *http.Request) {
    // r.Body is an application that passed the filter.
    w.WriteHeader(http.StatusAccepted)
}

func batchHandler(w http.ResponseWriter, r *http.Request) {
    // r.Body is a JSON array: every application of one portfolio.
    w.WriteHeader(http.StatusAccepted)
}
// The model the XML is parsed into: plain structs in package integration.
import "ideaswave.com/qubit/core/pkg/integration"

reg := integration.NewRegistry()
_ = reg.RegisterFilter(&integration.MessageFilter{
    ID:        "gold-only",
    Name:      "Gold tier only",
    Predicate: `tier = "gold"`, // FEEL; false means the message is dropped
    ChannelID: "svc",
})
reg.LookupFilter("gold-only").PassesVars(map[string]any{"tier": "silver"}) // false

// CorrelationContext matches a message to ONE waiting instance, once.
cc := integration.NewCorrelationContext("orderId")
cc.Register("A-17", "run_1ce49bbc5260b39423bd", "wait_payment")
cc.MatchVars(map[string]any{"orderId": "A-17"}) // the entry, now removed
cc.MatchVars(map[string]any{"orderId": "A-17"}) // nil: a second message resumes nothing

Layer 2 vs External Brokers

External brokers solve hard distribution problems - but most enterprise workflows don't need them. Layer 2 removes the accidental complexity.

vs Kafka

Kafka adds 3+ brokers and days of config

Kafka requires 3+ brokers, ZooKeeper or KRaft, topic management, consumer group coordination, and a dedicated ops team. It solves large-scale log distribution - not in-process workflow messaging.

Layer 2: in your own process, no broker to run
vs RabbitMQ

RabbitMQ needs a server before line one of code

RabbitMQ needs a dedicated server, an AMQP client library, exchange and binding configuration, and a connection string in every service. Infrastructure overhead precedes the first message.

Layer 2: the SDK camel runtime in your Go service, no server to install
vs AWS SQS

SQS adds latency and a billing line-item per message

SQS charges per API call, adds 50–200 ms of network latency on every message, and requires AWS credentials, SDK configuration, and IAM policies before you can enqueue anything.

Layer 2: in-process, no per-message fee

Up and running in 4 steps

One routes document and the Go handlers it names. No broker, no YAML manifests.

1

Declare your routes

Write a Camel Spring-DSL document: rest: endpoints, plus the <filter> and <aggregate> routes they enter through direct:<route id>-in.

2

Register your beans

Call rt.RegisterBean("loans.apply", handler) for each bean: step. A route whose bean is missing still binds and answers 501.

3

Bind to your mux

Call rt.Bind(mux, guard) once, then serve. Filtered messages answer 204, held ones 202, and a released group reaches the bean as one JSON array.

4

Start Priostack processes from a bean

To hand a message to a BPMN process, call POST /api/v1/process-instances from the bean with the message as variables (one credit per instance).

Layer 2 starter templates

Common shapes to build with the camel runtime. They are patterns to follow in your own service, not one-click deployments: nothing in the console deploys or runs them.

Aggregator pattern

Batch Loan Processing

Collect loan applications by portfolio ID and hand the batch to your risk bean once the completion condition is met.

MessageRouter pattern

Risk Level Routing

Declare the routes from loan decisions to auto-approval, manual review or decline channels with FEEL conditions; your code dispatches them.

MessageFilter pattern

Event Screening

Drop webhook events that fail a FEEL predicate before they enter the processing pipeline.

MessageChannel pub-sub

Fan-out Notifications

Declare a publish-subscribe channel for audit, email and analytics consumers; your beans do the delivery.

MessageTranslator + Pipeline

Pipeline Transform

Validate a request body against a JSON Schema (validate: step) before it reaches the bean that maps and delivers it.

Common questions

If your question isn't here, reach out to the team via the console or the community forum.

Can I use Layer 2 without Layer 1 (BPMN)?

Yes. The model is package integration (ideaswave.com/qubit/core/pkg/integration, module ideaswave.com/qubit/core v1.10.0) and the runtime is package camel (ideaswave.com/qubit/sdk/camel, module ideaswave.com/qubit/sdk v1.10.0). A camel runtime needs no BPMN process: it serves the REST routes your Camel document declares.

Does it require a message broker?

No. The camel runtime runs inside your Go process: a request enters through your http.ServeMux, passes the route's filter, aggregate, validate and log steps, and reaches your bean. There is no broker to deploy or manage.

Can Layer 2 talk to Kafka or RabbitMQ?

Not directly: there is no broker adapter. A bean is ordinary Go code, so it can call any Kafka or AMQP client library you choose.

How do I persist messages across restarts?

They are not persisted. The messages an aggregate step holds live in memory and are lost when the process stops; a group holds at most 10,000 messages and a route 100,000, beyond which the runtime answers 503. An aggregate that declares a completionTimeout is refused (its route answers 501). Keep anything that must survive a restart in your own store.

What is the throughput?

No benchmark is published for the camel runtime. Measure it with your own routes and payloads.

Layer 2 Trailblazers

Layer 2 runs in your own service, so priostack.com does not see your pipelines. Tell us what you built at support@priostack.com.