pumpIn

fun <T> Flow<T>.pumpIn(scope: CoroutineScope, onFailure: (PumpFailure, Throwable) -> Unit, name: String, body: suspend (T) -> Unit): Job

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

scope

the scope the pump runs in. Cancelling it still cancels the pump: a genuine cancellation of this job propagates through both guards untouched.

onFailure

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.

name

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.

body

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)
}