From AI prototype to credible information system.Explore Go Live Fintech
Technical

BPM + Integration: The Two-Layer Architecture Behind Priostack

ArchitectureEngineeringDeep dive

A common question when evaluating Priostack: "What happens when I need more than a job worker fetching tasks over REST? What if I need deduplication, content-based routing, or message correlation across multiple process instances?"

The answer is a two-layer architecture. Layer 1 is the BPM runtime - BPMN, DMN, CMMN execution with a clean REST job API. It is the complete runtime for the majority of use cases. Layer 2 is an optional set of Enterprise Integration Pattern (EIP) models from the qubit-core library that you run inside your own Go service with the SDK's camel runtime: no broker, and nothing extra to host at Priostack. It is not a hosted Priostack feature. The two layers are entirely independent by design: the integration layer knows nothing about processes, and the process runtime knows nothing about the integration layer.

This article walks through both layers in depth: what each one gives you, where the boundary sits, and - critically - how the integration layer makes idempotency, deduplication, and message orchestration a configuration problem rather than an infrastructure one.

1. The stack at a glance

┌──────────────────────────────────────────────────────────────┐ │ YOUR WORKERS / SERVICES │ │ (any language · any host · REST over HTTP) │ └───────────────────────────┬──────────────────────────────────┘ │ POST /api/v1/jobs/activate │ POST /api/v1/jobs/{key}/complete │ POST /api/v1/jobs/{key}/fail ▼ ┌──────────────────────────────────────────────────────────────┐ │ LAYER 2 - Integration patterns (optional, your service) │ │ │ │ Message Channel Content-Based Router Splitter │ │ Message Endpoint Message Filter Aggregator │ │ Message Translator Correlation Identifier Pipes & Filters │ │ │ │ • idempotency • content-based routing │ │ • deduplication • fan-out / fan-in │ │ • correlation • schema translation │ └───────────────────────────┬──────────────────────────────────┘ │ calls the same REST API ▼ ┌──────────────────────────────────────────────────────────────┐ │ LAYER 1 - BPM runtime (always) │ │ │ │ BPMN execution DMN decisions CMMN cases │ │ Token semantics Incident handling │ │ Instance state (ACTIVE / INCIDENT / COMPLETED / TERMINATED) │ └──────────────────────────────────────────────────────────────┘
The dependency rule, and why it is enforced The integration layer has no dependency on the process runtime at all. It knows about one thing - a message - and nothing about BPMN, decisions, cases or architecture models. That is not tidiness for its own sake: it is what lets you adopt Layer 2 for an integration problem that has nothing to do with process execution, and what guarantees that adding it can never change how your processes run.

2. Layer 1 - The BPM runtime

Layer 1 is the engine. It covers three process notations:

NotationWhat it modelsTypical use case
BPMN 2.0 Sequential and parallel process flows with tasks, gateways, events Approval workflows, order processing, service orchestration
DMN 1.3 Decision tables and FEEL expressions Credit scoring, eligibility rules, routing decisions
CMMN 1.1 Case management with discretionary tasks and sentries Fraud investigation, patient journeys, legal case handling

The REST job API

The fundamental interaction model is pull-based. Workers are external services - a Python microservice, a Node script, a Go binary, anything that speaks HTTP. They operate a three-step loop:

# 1. Activate  -  take waiting jobs from the engine (answers at once)
POST /api/v1/jobs/activate
     { "type": "credit-check", "maxJobsToActivate": 1 }
→    { "jobs": [ { "key", "type", "processInstanceKey", "variables", ... } ] }

# 2. Execute  -  your logic runs locally, fully isolated
# The job stays active until you complete or fail it; nothing expires it.

# 3. Complete or fail
POST /api/v1/jobs/{key}/complete   { "variables": { "approved": true } }  → 204
POST /api/v1/jobs/{key}/fail       { "errorMessage": "timeout", "retries": 2 }

No broker. No persistent connection. No SDK required. An activated job goes to one worker only, and nothing takes it back by itself: if the worker dies between steps 2 and 3, the job stays active until a worker releases it with a fail call that keeps its retries, or the server restarts. So a worker releases what it holds when it shuts down, and handlers are idempotent on the job key.

Instance state and observability

Every process instance has a state maintained by the engine: ACTIVE, INCIDENT, COMPLETED, or TERMINATED. An instance goes to INCIDENT when a worker fails a job with no retries left, and also on other run failures such as a decision that cannot be evaluated; only the first also writes a record to GET /api/v1/incidents, and GET /api/v1/process-instances?filter.state=INCIDENT lists them all. From the dashboard you can see which element the instance stopped at and inspect its variables. There is no REST call to retry an INCIDENT instance today, and cancelling one answers 422 because it is not active.

What Layer 1 alone gives you For the large majority of BPM use cases - approval flows, decision-driven routing, case management - Layer 1 is the complete runtime. You do not need Layer 2. Add it when you have a specific integration requirement, not as a default.

3. Layer 2 - The integration layer

Layer 2 implements nine patterns from the Hohpe & Woolf EIP catalogue, under the book's own names. Each component is a Go struct in ideaswave.com/qubit/core/pkg/integration that you register in a catalogue and compose into pipelines; in practice the SDK's camel runtime builds them from a Camel Spring-DSL (XML) routes file. Only the filter and aggregate steps run in that runtime today; the other patterns are declared in the model and carried out by your own code. Conditions and mappings are written in FEEL, the same expression language your decision tables already use.

Message Channel

A named conduit with two delivery modes: point-to-point, where each message reaches exactly one consumer, or publish-subscribe, where each message reaches every subscriber. A channel is typed to a declared item definition, so the schema a channel carries is part of the topology rather than an assumption each consumer makes privately.

Message Endpoint

Connects an external service to a channel, referencing the service by its identity in your ArchiMate model rather than by a hostname. Inbound endpoints are entry points. Because the reference is to the architecture model, "which services can put a message on this channel?" is a question your model answers.

Content-Based Router

Routes using FEEL conditions evaluated against the message payload. Routes are prioritised, the first matching route wins, and an empty condition is the catch-all default. Exactly one branch fires per message.

router  claim-router

  priority 10   →  channel high-value      when   amount > 10000
  priority  5   →  channel auto-approve    when   amount <= 500
  default       →  channel standard

Message Filter

Drops messages that do not satisfy a FEEL predicate. Non-matching messages are discarded and never reach the engine - useful for idempotency tokens (drop already-seen correlation keys) and for schema validation (drop malformed payloads before they raise incidents deep inside a process).

Aggregator

The most useful pattern for idempotency and fan-in. The Aggregator groups related messages by a correlation expression, holds them until a completion condition is satisfied, and then releases the batch as a single message. (The model has a timeout field, but the camel runtime refuses an aggregate that declares one.)

aggregator  order-results

  group by      orderId                 group messages sharing this key
  release when  count(items) >= 3       the completion condition
  emit to       channel order-complete

In a fan-in scenario: three upstream services each report a result for the same order. The results arrive at the Aggregator under the same correlation key, and only when the third arrives does one combined message go on, for example to start or complete one step of a process. Note that racing workers on one job do not happen: the engine hands an activated job to one worker only, and a repeated completion answers 404.

Correlation Identifier

Holds the mapping from a correlation key to the process instance waiting on it. The canonical use case is a BPMN receive task: the task suspends the instance and registers the key it is waiting for. When a message arrives with a matching key, the match returns that entry and removes it in the same atomic step, so a second message with the same key matches nothing. That is first-match-wins deduplication rather than a check in every sender that might send a duplicate. The engine waits on message catch events this way (zeebe:subscription correlationKey), but the hosted API cannot publish a message yet: on priostack.com a waiting receive task is offered to workers as a job of type message:<name>, completed without a key check.

correlation identifier   key expression: paymentId

  receive task activates   →  register  paymentId  →  instance, resume point
  confirmation arrives     →  match     paymentId  →  entry found, removed
                                                       instance resumes
  duplicate arrives        →  match     paymentId  →  nothing; ignored

Splitter

Produces many output messages from one input: a FEEL expression evaluates to a list and one message is emitted per element. Use it for fan-out - a single order confirmation that needs to trigger a fulfilment task, a notification task, and an audit task in parallel, each as an independent process instance.

Message Translator

Declares a schema transformation between two item definitions as a FEEL mapping expression, which separates integration concerns - field renaming, unit conversion, restructuring - from process logic. The process receives a correctly shaped payload regardless of what the upstream service emitted, and the mapping is a declaration you can inspect rather than code buried in an adapter.

Pipes and Filters

An ordered chain of the components above. A message enters at the first step and passes through each in order - a typical chain being Filter, then Translator, then Router, then Aggregator - with the output of one feeding the input of the next. If a filter rejects the message, execution stops there and nothing downstream sees it.

pipeline  payment-ingest

  step 1   filter       duplicate-filter
  step 2   translator   schema-translator
  step 3   router       value-router
  step 4   aggregator   result-aggregator

4. Idempotency without a broker

The typical argument against REST polling for job workers is that it requires idempotency to be implemented by each worker individually. On Priostack two workers never race on one job: an activated job goes to one worker and is never handed out again by itself. Duplicates come from elsewhere: a worker that did the work but crashed before completing (whoever releases the job runs the work again), or an external event delivered twice.

Layer 2, running in your own service in front of the API, handles the external side. Here is the whole stack:

MechanismWhereWhat it prevents
Aggregator Layer 2 Related events are held and passed on once, as one group
Correlation Identifier Layer 2 A second message with the same key matches nothing and is ignored
Message Filter Layer 2 Already-seen idempotency tokens dropped before they reach the API
Single activation Layer 1 An activated job goes to one worker only; a repeated completion answers 404 and changes nothing
Point-to-point channel Layer 2 Each message is delivered to exactly one consumer by channel semantics

Workers keep one habit: make each handler idempotent on the job key, so work that runs again after a crash does no harm. Beyond that they can stay stateless.

Practical consequence You can write a worker in Python and deploy it as a Lambda function or a Kubernetes Job. Key its side effects on the job key, release the jobs it cannot finish with a fail call that keeps their retries, and put a filter in front of external events that may arrive twice. A duplicated completion call is harmless: the second one answers 404.

5. When to use which layer

Layer 1 only - when to stop here

Your workers are reliable services (internal microservices, not lambdas), job durations are short (under 60 seconds), you have a small number of concurrent instances (under a few thousand), and you do not need cross-process message correlation. This covers most BPM consultant use cases: approval flows, credit decisions, HR onboarding, fraud escalation. Layer 1 is the complete system.

Layer 2 - when you actually need it

Workers are ephemeral (Lambdas, spot instances, serverless), you have high-frequency job completion from multiple parallel workers, you need content-based routing between process variants, you are correlating external messages (webhooks, payment confirmations, IoT events) in your own service before they reach a process, or you are building fan-out/fan-in patterns (one order triggers three parallel tracks that must all complete before proceeding).

The decision is not permanent. You start with Layer 1 and add integration components incrementally when a specific requirement surfaces. There is no migration and no data re-modelling. Layer 2 runs in your own service, so adding it is a deploy of that service, not of Priostack, and the processes already running are unaffected.

6. The scale story

REST polling has a real ceiling. At tens of thousands of concurrent workers with sub-second polling intervals, the activate endpoint becomes a bottleneck. This is acknowledged honestly: for extremely high-frequency job throughput, a message queue (Kafka, RabbitMQ, SQS) in front of the activate endpoint is the right call.

But that ceiling is much higher than most assume, and the integration layer raises it further. Consider:

  • A worker polling every second with an average job duration of 10 seconds generates approximately one activate request every 10 seconds while busy.
  • An Aggregator can turn many small upstream events into one instance start.
  • A Message Filter drops invalid or duplicate events before they reach the API, saving requests against the limit of 300 per minute per IP.

For the use cases Priostack is designed for - BPM consultant tooling, enterprise approval flows, architecture-to-execution workflows - the ceiling is never reached. The typical deployment has dozens to hundreds of concurrent instances, not millions. You add a broker when you have a broker-scale problem. Not before.

The honest summary Layer 1 alone is the simple, correct choice for most BPM workflows. Layer 2 is not a complexity tax you pay by default - it is a set of tools available when a specific integration requirement surfaces. The two layers are independent, composable, and connected by a shared substrate that means you never pay an impedance mismatch cost at the boundary.

If you want to see the integration layer in action, the EIP Pipelines article walks through two complete topologies - a big-data corpus-building pipeline and a fast-data anomaly-routing pipeline - built entirely with these nine primitives.

Priostack Engineering

Architecture-to-execution stack for technical transformation teams. Questions and feedback welcome on Telegram.