Give Unfinished Work an Owner

Compare Async, a transient VM queue and an account-dependent MQ consumer with explicit acknowledgement ownership.

A caller can finish before every downstream activity finishes. That reduces response waiting, but it raises a more important question: who owns the unfinished work after the response has gone?

Two local mechanisms make deliberately modest promises about that, and their modesty is the lesson; a broker placed between the outbox publisher and a consumer makes a larger one. Broker access is an optional account-dependent path — the local service remains complete without it.

Publisher: Stable event identity: Durable intent. Queue + consumer: Receive and process: Retain ACK token. Receiver then ACK: Effect first: Redelivery can repeat. Transient VM and account-dependent MQ have different guarantees.

Open this checkpoint in ACB

Stop the previous local application, then use File → Open Folder to open book/checkpoints/20-async-and-queues. Open src/main/mule/app.xml and use Flow List to select the named flow for each example. Start the controlled dependency in a separate terminal from the companion root with python3 book/stubs/server.py. Choose Run and Debug → Run Mule Application and wait for deployment. Save canvas edits and use Save and Hot-deploy to Local Runtime before repeating a request.

Run python3 book/run.py verify 20 from the companion root to exercise this running checkpoint with its synthetic fixtures. The verifier supplies requests and checks results; it does not start the ACB application. Keep the editor on this checkpoint while reading a failure so an old deployment cannot supply a misleading answer.

Start work without joining its result

The cheapest way to stop waiting is to run the branch beside the main path and answer the caller straight away. The question that follows is what the answer is then entitled to say.

Example 052 — Start a local asynchronous notice

The complete source is in checkpoints/20-async-and-queues/src/main/mule/app.xml.

Choose Flow List → async-notice. Expand Async and select its Logger: Level INFO, Message Local asynchronous notice. The Transform Message after Async returns JSON {startedLocally: true}. The Listener accepts POST /lab/async.

Keep Transform Message outside Async. Moving it into the scope would change which branch constructs the caller’s answer.

The Async scope runs its processors separately from the main path. The response says startedLocally, not delivered or persisted — the logger is useful for observing execution, but the application has not stored a durable job that survives process loss. An Async failure cannot retroactively change the HTTP response already sent, so the branch has to own its own failure.

Avoid sharing mutable objects across that boundary, or treating asynchronous variable changes as a return value: copying an event does not deep-copy a mutable payload, and a shared Java object is not made safe by putting one reference on another execution path. If the caller needs a result, use an explicit request-response path. If the business needs recovery, establish durable ownership before acknowledging acceptance. The decision that matters is whether the side effect is dispensable — not whether it sounds like background work.

VM Publish retains the local queue name and payload content

VM Publish retains the local queue name and payload content.

Put a message on a local queue

A queue moves the handoff out of the flow and gives the unfinished work a name of its own. Whether it also gives that work durability is a separate question, and here the answer is no.

Example 053 — Publish a local VM notice

The complete source is in checkpoints/20-async-and-queues/src/main/mule/app.xml.

Choose Flow List → publish-local-notice and select VM Publish. Its connection is Local_VM, Queue Name is order-notices, and Content is expression payload. Open the connection to inspect its queue definition: Queue Type is TRANSIENT.

The following Set Variable has Name httpStatus and numeric Value 202. Transform Message returns JSON {queuedLocally: true}. Select the Listener to check POST /lab/notices and its Response status expression vars.httpStatus default 200.

Example 054 — Consume a local VM notice

The complete source is in checkpoints/20-async-and-queues/src/main/mule/app.xml.

Choose Flow List → consume-local-notice. This flow starts with VM Listener, not HTTP Listener. It uses Local_VM and Queue Name order-notices. The following File Write uses Order_Files, Path notice.json and Content expression output application/json --- payload.

Use the canvas Flow List to move between the publisher and consumer. They are separate flows: a line on the publisher canvas does not continue into the consumer.

The complete checkpoint declares the order-notices queue as TRANSIENT. POST /lab/notices returns 202 after local publication. The listener writes the notice as JSON; GET /lab/notices reads the file. Use a fresh notice identity when testing — an old file can otherwise impersonate a new delivery.

VM gives this application a source/consumer boundary, but the selected transient queue does not promise restart durability. Its notice file is a diagnostic result, not the order’s authoritative outbox. The consuming flow also starts a new event, so anything it needs has to be in the message — depending on the original HTTP attributes after a queue boundary is a hidden coupling. Queue persistence and transaction support depend on the connector configuration and deployment; do not infer them from the existence of a queue name. A transient queue stays useful inside one running application, when lost work can be regenerated or another durable system remains responsible for it.

Move publication across a broker boundary

A broker puts the message outside the application replica, which is the first arrangement in this chapter that could outlive the process that accepted the order. It also moves the interesting decision to acknowledgement: the moment the consumer says the message is finished with.

Example 055 — Publish and acknowledge an MQ order event

Source: platform/mq/src/main/mule/app.xml.

Open book/platform/mq as a separate ACB project. Its README lists the account prerequisites. In src/main/mule/app.xml, inspect these two flows without starting the subscriber until the sandbox queue and credentials exist:

Canvas locationComponent settings
publish-order-event sourceHTTP Listener, connection Publish_Listener, POST /events, loopback port 18886
Publisher processorMQ Publish, MQ_Config, Destination ${mq.queue}, Body expression write(payload, 'application/json')
consume-order-event sourceMQ Subscriber, same destination, Acknowledgement Mode MANUAL
First consumer processorSet Variable, Name ackToken, Value attributes.ackToken
Consumer RequestReceiver_HTTP, POST /events, Body output application/json --- read(payload, 'application/json'), header Content-Type: application/json
Last consumer processorMQ Ack, MQ_Config, Ack Token expression vars.ackToken
Consumer Error handlerOn Error Propagate ANY, then Logger at ERROR

In Global Configurations, MQ_Config takes URL ${mq.url}, Client ID ${mq.clientId} and Client Secret ${mq.clientSecret}. Receiver_HTTP uses ${receiver.host}, ${receiver.port} and response timeout 2000. Supply real secrets through protected runtime configuration; leave placeholders in shared source and screenshots.

This separate project accepts a local event on port 18886, publishes its JSON representation to the configured Anypoint MQ destination, and subscribes in manual acknowledgement mode. It saves the acknowledgement token before the HTTP call replaces message attributes. The ACK follows successful receiver processing.

To connect it to the outbox, point the dispatcher’s publication requester at port 18886. Keep the subscriber’s receiver on port 18882. Pointing the subscriber back at the publication endpoint would create a loop instead of an effect.

The README supplies the required queue, client properties and topology. This adapter has not been deployed against an MQ account. Its source is a concrete starting project, and broker delivery, timeout, ACK failure and redelivery remain acceptance work. Local HTTP or VM results cannot establish those broker behaviors.

Manual acknowledgement makes the ordering of operations visible: receive, establish the effect, then ACK. The dangerous sequence is business commit, crash, no ACK — a second delivery is then correct broker behavior, and the receiver’s event identity is the only thing preventing a duplicate effect. ACK-first has the opposite risk: work lost if processing subsequently fails. Anypoint MQ acknowledgement.

Poison messages and available capacity

A malformed event can fail on every delivery. Repeating it indefinitely consumes receiver capacity and obscures useful work. Define a delivery limit, dead-letter or quarantine destination, retained identity and an operator procedure for repair and replay. Moving a message to a dead-letter queue is an operational state, not fulfillment.

Classify failures before choosing delay and concurrency. A temporary network failure may merit retry; an invalid schema needs repair; an expired credential can affect every message. Rapidly retrying all three with the same policy increases load without increasing certainty.

Backpressure means unfinished work influences admission or consumption rate. Bound consumer concurrency and dependency connections, then measure queue age and destination saturation. A broker is a buffer — its retention and capacity remain finite, so a short HTTP response does not mean the system is keeping up if pending work grows faster than completion. Adding consumers against a saturated destination can make the outage worse.

Try it

1. Interpret 202. Does Publish a local VM notice establish survival across a runtime restart?

Show answer

No. The configured queue is transient and the response promises local queueing. Test and declare a durable handoff separately if restart survival is required.

2. Lose an ACK. Why does manual acknowledgement still need receiver idempotency?

Show answer

Processing can commit before the acknowledgement reaches the broker. Redelivery must carry the same stable business/event identity so another attempt does not repeat the effect.

3. Diagnose growing age. Queue depth is steady, but the oldest message keeps aging. What should an operator inspect?

Show answer

Look for a poison or repeatedly failing message, partition or ordering blockage, consumer ownership and destination failures. Aggregate depth alone can conceal work that never reaches a terminal state.

A queue handles a stream of messages. A finite import has another useful shape: a known input whose individual records need processing and a completion report. That is the next chapter’s Batch Job.

Next: Track Each Record in an Import

Comments