pumpIn
Collect this flow in scope as a long-lived pump that a throw cannot kill — neither one raised in body nor one raised by the flow — reporting either through onFailure.
Use this instead of onEach { … }.launchIn(scope) for any pump that has to keep running for the life of a session and whose flow or body reaches consumer-authored or peer-supplied code.
Return
the pump's Job. After PumpFailure.UPSTREAM it completes normally — the pump is over, but it ended with a diagnosis rather than a SIGABRT.
Parameters
the scope the pump runs in. Cancelling it still cancels the pump: a genuine cancellation of this job propagates through both guards untouched.
invoked with which half failed and the throwable. Required, not defaulted — a pump that absorbs in silence is what this class of defect is made of, and every caller should have to decide what to do about a dead pump. :kuilt-core is logger-free by contract, so this is the only signal there is. Best-effort and non-suspending.
what this pump is called in a coroutine dump or census — see the section above. Make it distinct per pump instance, not per call site: where several pumps of one kind run side by side, qualify it with whatever tells them apart ("composite-ply-peers[$plyId]"), so a census can group by kind and still name the instance.
the per-item work.
Samples
runTest {
val applied = mutableListOf<String>()
val reported = mutableListOf<PumpFailure>()
// A consumer-authored flow that hands over one item the body cannot apply, and then fails outright.
val updates = flow {
emit("apply-me")
emit("i-will-not-apply")
error("…and then the flow itself gave up")
}
val pump = updates.pumpIn(
scope = backgroundScope,
// ITEM: that update was lost, the pump lives. UPSTREAM: the pump is over — say so, loudly.
onFailure = { half, _ -> reported += half },
// What a coroutine census calls this pump when it is the one that wedged.
name = "sample-updates",
) { update ->
if (update == "i-will-not-apply") error("this update could not be applied")
applied += update
}
pump.join()
assertEquals(listOf("apply-me"), applied)
assertEquals(listOf(PumpFailure.ITEM, PumpFailure.UPSTREAM), reported)
}