discoveryRoster

fun discoveryRoster(sources: List<PeerDiscoverySource>, scope: CoroutineScope, onSourceFailure: (PeerDiscoverySource, Throwable) -> Unit = { _, _ -> }): StateFlow<Set<Tag>>

Merge several discovery feeds into one live roster of who this peer can currently see.

A phone browsing for nearby games over both Bonjour (mDNS) and Apple Multipeer has two separate feeds of "someone just appeared" / "someone just left". This folds all of them — sources — into a single set that a lobby UI can render directly, so you don't hand-write the merge each time. Every PeerDiscoverySource.discoveries event adds a peer, every PeerDiscoverySource.departures event removes one, keyed on Tag.peerKey.

The returned StateFlow claims only this peer's current best view — nothing more. It is not an agreement, a vote, or a decision about who hosts. Two peers folding the same feeds can hold different rosters at the same instant (a feed lags, a departure hasn't propagated, the same physical peer carries different Tag.peerKeys across transports), so this view is not a safe election input. Pick a host from us.tractat.kuilt.core.Seam.peers once connected — never from a discovery roster. See docs/discovery-bootstrap.md.

One feed's failure costs you that feed, and nothing else. Each source's discoveries() and departures() are isolated separately, so a transport that dies mid-session degrades to "that feed stopped reporting" rather than taking the whole fold down with it. Without the isolation a single throw cancels the merged flow, kills the one backing coroutine, and freezes this StateFlow at its last value permanently — every healthy source silently stops being observed too, with the roster still looking plausible (#1904). A failed feed cannot be restarted: a source is asked for its flow exactly once.

A failed feed's peers linger — the roster is best-view, not liveness. When a source's discoveries() dies, the peers it already contributed stay in the set. They are not synthesised away, because nothing observed them leave: dropping them would assert a departure this peer never saw, and would turn one transient transport fault into a mass exodus in the UI. This is exactly the ghost caveat below, reached by a second route — a failed feed is a feed that can no longer report departures — so the same reading applies: those entries are what a dead feed last knew, not a claim that anyone is still there.

Ghost caveat — the roster is add-only over a source that returns emptyFlow(). PeerDiscoverySource.departures has no default: a source with no leave signal (a fixed-roster test fake, a platform stub, a browse API that only reports arrivals) must return emptyFlow() explicitly. Over such a source a discovered peer is never removed — it lingers as a ghost long after it is gone, and the set only grows. This is a real limitation of what that feed can tell you, not a bug in the fold. The caveat applies to exactly those sources: read a source's departures() body, and if it is emptyFlow(), expect a stale roster from it.

The fold runs on scope: its backing coroutine is launched eagerly there and is cancelled when scope is cancelled. scope is required — pass the caller's scope (in a test, a backgroundScope bound to the test clock); this function never spins up a dispatcher of its own.

Parameters

onSourceFailure

invoked with the source whose feed died and the throwable that killed it. :kuilt-core is logger-free by contract, so this callback is the only signal a dead feed can produce — absent it the isolation above would be silent, which would be a worse diagnosis than the freeze it replaces (a freeze at least reaches the scope's CoroutineExceptionHandler). The throwable is passed whole rather than summarised: it came out of the consumer's own source implementation, so its trace is the diagnosis — including which of the two feeds died, which is why that is not a separate parameter. Best-effort and non-suspending; whatever it throws is absorbed, since a consumer's logger must never be able to kill the fold this exists to protect. Defaulted for the same reason us.tractat.kuilt.core.fabric.acceptPump's is — this is a diagnostic hook on an existing primitive, not a new obligation on its callers. Whether a failed source should be observable in the roster's own type, so a lobby can render "mDNS is down", is a separate design question this deliberately does not answer.

Samples

runTest {
    // Two transports, each a feed of "appeared" / "left" events.
    val mdnsPeers = MutableSharedFlow<Tag>(extraBufferCapacity = 8)
    val mdnsGone = MutableSharedFlow<String>(extraBufferCapacity = 8)
    val mdns = object : PeerDiscoverySource {
        override val kind = DiscoveryKind.Mdns
        override fun discoveries(): Flow<Tag> = mdnsPeers
        override fun departures(): Flow<String> = mdnsGone
    }
    val multipeer = object : PeerDiscoverySource {
        override val kind = DiscoveryKind.Multipeer
        override fun discoveries(): Flow<Tag> = emptyFlow() // idle in this sample
        override fun departures(): Flow<String> = emptyFlow() // idle too — nothing to leave
    }

    // One StateFlow the lobby UI renders directly — no hand-rolled merge.
    val roster = discoveryRoster(listOf(mdns, multipeer), backgroundScope)
    runCurrent()

    mdnsPeers.emit(InMemoryTag("alice"))
    mdnsPeers.emit(InMemoryTag("bob"))
    runCurrent()
    check(roster.value.map { it.peerKey }.toSet() == setOf("alice", "bob"))

    // A departure removes the peer, keyed on Tag.peerKey.
    mdnsGone.emit("alice")
    runCurrent()
    check(roster.value.map { it.peerKey }.toSet() == setOf("bob"))
}
runTest {
    val mdnsPeers = MutableSharedFlow<Tag>(extraBufferCapacity = 8)
    val multipeerPeers = MutableSharedFlow<Tag>(extraBufferCapacity = 8)
    val mdns = object : PeerDiscoverySource {
        override val kind = DiscoveryKind.Mdns

        // A real source fails from inside its own flow — a callbackFlow whose jmdns listener
        // throws. This one fails when the sample sends it the sentinel, so the sample controls when.
        override fun discoveries(): Flow<Tag> = mdnsPeers.transform { tag ->
            if (tag.peerKey == "boom") error("jmdns IO error") else emit(tag)
        }

        override fun departures(): Flow<String> = emptyFlow()
    }
    val multipeer = object : PeerDiscoverySource {
        override val kind = DiscoveryKind.Multipeer
        override fun discoveries(): Flow<Tag> = multipeerPeers
        override fun departures(): Flow<String> = emptyFlow()
    }

    // kuilt-core is logger-free, so this callback is the only signal a dead feed can produce.
    val dead = mutableListOf<DiscoveryKind>()
    val roster = discoveryRoster(
        listOf(mdns, multipeer),
        backgroundScope,
        onSourceFailure = { source, _ -> dead += source.kind },
    )
    runCurrent()

    mdnsPeers.emit(InMemoryTag("alice"))
    multipeerPeers.emit(InMemoryTag("bob"))
    runCurrent()
    check(roster.value.map { it.peerKey }.toSet() == setOf("alice", "bob"))

    mdnsPeers.emit(InMemoryTag("boom")) // the mDNS feed dies here
    runCurrent()
    check(dead == listOf(DiscoveryKind.Mdns))

    // Multipeer carries on, and alice lingers — nothing ever observed her leave.
    multipeerPeers.emit(InMemoryTag("carol"))
    runCurrent()
    check(roster.value.map { it.peerKey }.toSet() == setOf("alice", "bob", "carol"))
}