#419 · bluetape4k-projects · 2.0.0 · Sequence

NATS JetStream Flow: one cold contract, two intake paths

Play the ConsumerContext pull and JetStream push paths as one bounded Flow<Message>, then follow caller-owned acknowledgement, redelivery, drop, cancellation, and terminal completion.

Choose a delivery story

BRANCH × 5

Sequence: NATS JetStream Flow: one cold contract, two intake paths

Step 0 / 9
Sequence progress1 · Declare a cold Flow
Caller / collectorack · nak · term policy
cold Flow<Message>single bounded contract
NATS adapterowns only its handle
JetStream / consumerserver delivery + state
  1. JetStream.consumeAsFlow(...) or ConsumerContext.consumeAsFlow(...)
    1 · Declare a cold Flowwaiting
  2. collect → runInterruptible(Dispatchers.IO)
    2 · Collection opens the adapter handlewaiting
  3. nextMessage(receiveTimeout) → send(message)
    3 · Pull and push convergewaiting
  4. emit(Message)
    5 · Deliver without inventing policywaiting
  5. business work → message.ack()
    6 · Caller acknowledges successwaiting
  6. message.nak() / ack wait expires → redelivery
    7 · Nak or unack makes redelivery visiblewaiting
  7. droppedCount ↑ → NatsConsumerFlowException
    8 · Push drop becomes a typed failurewaiting
  8. cancel → interrupt receive → close / unsubscribe
    9 · Cancellation closes owned handleswaiting
  9. message.term() or consumer end → Flow completion
    10 · Terminal choice and completionwaiting

1 · Declare a cold Flow

waiting

JetStream.consumeAsFlow(...) or ConsumerContext.consumeAsFlow(...)

What happens

Guardrail

What follows

Visible signal