fix(audio-rooms): clean up deferred audit items (VM + Android + MoQ comment)
Lands the audit follow-ups that didn't require external input. Only the wire-format draft pinning (MoQ #1, #8) remains deferred — that's gated on the M4 manual interop pass, and the same-room-PIP-re-entry corner case (Android #5) — Android has no programmatic PIP-exit API. ViewModel: - VM #10: serialize disconnect→connect via a tracked `pendingCloseJob`. teardown() records the listener.close() launch (when not finalCleanup); the next launchConnect() awaits it before opening a fresh transport. Eliminates the brief two-QUIC-session overlap that some MoQ relays reject by deduping on client pubkey. - VM #4: extracted shared connect-launch body into `launchConnect(triggerRetryOnFailure)` so connect() / connectInternal() share one implementation. Removed the near-duplicate viewModelScope.launch block. - VM #8b: setMicMuted no longer silently swallows handle failures. BroadcastUiState.Broadcasting gains a `muteError: String?` field that the UI can surface as an inline message; cleared on the next successful toggle. The broadcast itself stays running with its previous mute state — only the mute toggle failed. - VM #6: documented the dispatcher-confinement contract in the class kdoc. Audit was theoretical — every map mutation already runs on viewModelScope (Dispatchers.Main.immediate on Android, same dispatcher the MoQ flow's onEach callback uses because the player launch lives in viewModelScope). Future cross-thread callers must marshal explicitly. Android: - Android #2: AudioRoomActivity.toggleMuteSignal type tightened from `MutableSharedFlow<Unit>` to `SharedFlow<Unit>` so external code can't tryEmit into it. Internal emit uses the private `_toggleMuteSignal`. - Android #10: presence debounce-publisher (the LaunchedEffect keyed on micMutedTag) now skips entirely when micMutedTag is null. Stops the duplicate first-frame publish where heartbeat fires immediately AND debounce-publisher fires 500 ms later, both with muted=null. Once the user goes live the debounce-publisher kicks in for state changes. MoQ session: - MoQ #11: send() rollback comment rewritten to say "monotonic; gaps acceptable per spec, this just minimises them on full-fanout failures" instead of the misleading "strictly contiguous" claim. Verified: spotlessApply clean; :commons:jvmTest, :nestsClient:jvmTest (80 tests), :amethyst:compilePlayDebugKotlin all green. Still deferred: - MoQ #1, #8: wire-format draft pinning (draft-17 vs draft-11). Needs M4 interop input from `nostrnests.com` to know what the relay actually speaks. - Android #5: same-room re-entry from MainActivity while in PIP doesn't auto-exit PIP. Android has no programmatic PIP-exit API; user must tap the expand button. Corner case. - Test coverage gaps (round-1 VM #10, round-2 VM #13): retry-counter + broadcast state + setMicMuted-no-handle + server-Closed cleanup + double-connect-while-Failed. Each is a small dedicated test using the existing connector-seam pattern; landing as a separate test-only commit.
This commit is contained in:
+9
-4
@@ -38,6 +38,7 @@ import com.vitorpamplona.amethyst.R
|
||||
import com.vitorpamplona.amethyst.ui.theme.AmethystTheme
|
||||
import kotlinx.coroutines.channels.BufferOverflow
|
||||
import kotlinx.coroutines.flow.MutableSharedFlow
|
||||
import kotlinx.coroutines.flow.SharedFlow
|
||||
import kotlinx.coroutines.flow.asSharedFlow
|
||||
|
||||
/**
|
||||
@@ -72,7 +73,7 @@ class AudioRoomActivity : AppCompatActivity() {
|
||||
// and forwarded to the VM. Replaces an earlier
|
||||
// process-wide singleton that could leak edges across
|
||||
// activity instances.
|
||||
toggleMuteSignal.tryEmit(Unit)
|
||||
_toggleMuteSignal.tryEmit(Unit)
|
||||
}
|
||||
|
||||
ACTION_PIP_LEAVE -> {
|
||||
@@ -85,8 +86,12 @@ class AudioRoomActivity : AppCompatActivity() {
|
||||
private val _toggleMuteSignal =
|
||||
MutableSharedFlow<Unit>(replay = 0, extraBufferCapacity = 1, onBufferOverflow = BufferOverflow.DROP_OLDEST)
|
||||
|
||||
/** Emitted when the PIP-overlay mute action is tapped. */
|
||||
val toggleMuteSignal: MutableSharedFlow<Unit> get() = _toggleMuteSignal
|
||||
/**
|
||||
* Emitted when the PIP-overlay mute action is tapped. Read-only so
|
||||
* external code can't tryEmit into it from outside the Activity
|
||||
* (audit round-2 Android #2).
|
||||
*/
|
||||
val toggleMuteSignal: SharedFlow<Unit> get() = _toggleMuteSignal.asSharedFlow()
|
||||
|
||||
override fun onCreate(savedInstanceState: Bundle?) {
|
||||
super.onCreate(savedInstanceState)
|
||||
@@ -144,7 +149,7 @@ class AudioRoomActivity : AppCompatActivity() {
|
||||
}
|
||||
},
|
||||
onConnectedChange = { isConnected.value = it },
|
||||
pipMuteSignal = toggleMuteSignal.asSharedFlow(),
|
||||
pipMuteSignal = toggleMuteSignal,
|
||||
onLeave = { finish() },
|
||||
)
|
||||
}
|
||||
|
||||
+10
-3
@@ -184,9 +184,16 @@ private fun AudioRoomActivityBody(
|
||||
// PRESENCE_DEBOUNCE_MS for further changes before sending a fresh
|
||||
// presence event. The LaunchedEffect's auto-cancel on key change
|
||||
// serves as the debounce mechanism.
|
||||
LaunchedEffect(micMutedTag) {
|
||||
delay(PRESENCE_DEBOUNCE_MS)
|
||||
publishPresence(account, event, handRaised, micMutedTag)
|
||||
//
|
||||
// Skip when `micMutedTag` is null (we're not broadcasting) — otherwise
|
||||
// the FIRST composition would publish twice within 500 ms (heartbeat
|
||||
// immediately + debounce-publisher 500 ms later, both with muted=null).
|
||||
// Audit round-2 Android #10.
|
||||
if (micMutedTag != null) {
|
||||
LaunchedEffect(micMutedTag) {
|
||||
delay(PRESENCE_DEBOUNCE_MS)
|
||||
publishPresence(account, event, handRaised, micMutedTag)
|
||||
}
|
||||
}
|
||||
DisposableEffect(event.address().toValue()) {
|
||||
onDispose {
|
||||
|
||||
+66
-42
@@ -76,6 +76,17 @@ import kotlin.coroutines.cancellation.CancellationException
|
||||
* [playerFactory] so commonMain doesn't have to know which platform's
|
||||
* MediaCodec / AudioTrack is in play. M1 wires Android-only — desktop
|
||||
* passes nothing here yet.
|
||||
*
|
||||
* **Threading contract:** all public methods (`connect`, `disconnect`,
|
||||
* `updateSpeakers`, `setMuted`, `setMicMuted`, `startBroadcast`,
|
||||
* `stopBroadcast`) MUST be called on the main thread. Internal mutation
|
||||
* of `activeSubscriptions` / `speakingExpiryJobs` / `requestedSpeakers`
|
||||
* is only thread-safe under that assumption — they're plain HashMaps
|
||||
* confined to `viewModelScope`'s dispatcher (Dispatchers.Main.immediate
|
||||
* on Android, which is the same dispatcher the MoQ flow's `onEach`
|
||||
* callback runs on because the player launch lives in viewModelScope).
|
||||
* If a future caller needs to invoke from a background thread, marshal
|
||||
* via `viewModelScope.launch(Dispatchers.Main.immediate) { ... }` first.
|
||||
*/
|
||||
@Stable
|
||||
class AudioRoomViewModel(
|
||||
@@ -112,6 +123,12 @@ class AudioRoomViewModel(
|
||||
private var connectJob: Job? = null
|
||||
private var stateObserverJob: Job? = null
|
||||
|
||||
// Last in-flight `listener.close()` launched by teardown(). A subsequent
|
||||
// connect() awaits this before opening a fresh transport so two QUIC
|
||||
// sessions (the closing old one and the opening new one) don't briefly
|
||||
// coexist. Some servers dedupe by client pubkey and reject the new one.
|
||||
private var pendingCloseJob: Job? = null
|
||||
|
||||
private val activeSubscriptions = mutableMapOf<String, ActiveSubscription>()
|
||||
private val speakingExpiryJobs = mutableMapOf<String, Job>()
|
||||
private var requestedSpeakers: Set<String> = emptySet()
|
||||
@@ -175,34 +192,7 @@ class AudioRoomViewModel(
|
||||
teardown(targetState = ConnectionUiState.Idle, finalCleanup = false)
|
||||
}
|
||||
|
||||
_uiState.update { it.copy(connection = ConnectionUiState.Connecting(ConnectionUiState.Step.ResolvingRoom)) }
|
||||
|
||||
connectJob =
|
||||
viewModelScope.launch {
|
||||
try {
|
||||
val l =
|
||||
connector.connect(
|
||||
httpClient = httpClient,
|
||||
transport = transport,
|
||||
scope = viewModelScope,
|
||||
serviceBase = serviceBase,
|
||||
roomId = roomId,
|
||||
signer = signer,
|
||||
)
|
||||
if (closed) {
|
||||
runCatching { l.close() }
|
||||
return@launch
|
||||
}
|
||||
listener = l
|
||||
observeListenerState(l)
|
||||
} catch (ce: CancellationException) {
|
||||
throw ce
|
||||
} catch (t: Throwable) {
|
||||
_uiState.update {
|
||||
it.copy(connection = ConnectionUiState.Failed(t.message ?: t::class.simpleName ?: "connect failed"))
|
||||
}
|
||||
}
|
||||
}
|
||||
launchConnect(triggerRetryOnFailure = false)
|
||||
}
|
||||
|
||||
fun setMuted(muted: Boolean) {
|
||||
@@ -293,19 +283,22 @@ class AudioRoomViewModel(
|
||||
if (closed) return
|
||||
val handle = broadcastHandle ?: return
|
||||
viewModelScope.launch {
|
||||
handle
|
||||
.runCatching { setMuted(muted) }
|
||||
.onSuccess {
|
||||
if (closed) return@onSuccess
|
||||
_uiState.update {
|
||||
val current = it.broadcast
|
||||
if (current is BroadcastUiState.Broadcasting) {
|
||||
it.copy(broadcast = current.copy(isMuted = muted))
|
||||
} else {
|
||||
it
|
||||
}
|
||||
}
|
||||
val result = handle.runCatching { setMuted(muted) }
|
||||
if (closed) return@launch
|
||||
_uiState.update {
|
||||
val current = it.broadcast
|
||||
if (current !is BroadcastUiState.Broadcasting) return@update it
|
||||
if (result.isSuccess) {
|
||||
it.copy(broadcast = current.copy(isMuted = muted, muteError = null))
|
||||
} else {
|
||||
// Surface the failure so the UI can show a toast / inline
|
||||
// message instead of silently swallowing (audit round-2
|
||||
// VM #8b). We keep the broadcast running with the prior
|
||||
// mute state — only the mute toggle failed.
|
||||
val why = result.exceptionOrNull()?.message ?: "mute failed"
|
||||
it.copy(broadcast = current.copy(muteError = why))
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -462,12 +455,32 @@ class AudioRoomViewModel(
|
||||
if (closed) return
|
||||
val current = _uiState.value.connection
|
||||
if (current is ConnectionUiState.Connecting || current is ConnectionUiState.Connected) return
|
||||
launchConnect(triggerRetryOnFailure = true)
|
||||
}
|
||||
|
||||
/**
|
||||
* Shared connect-launch body for [connect] (user-driven, no retry on
|
||||
* failure) and [connectInternal] (auto-retry path, schedules another
|
||||
* retry on failure). Awaits any in-flight `listener.close()` from a
|
||||
* previous teardown so two QUIC sessions don't briefly coexist
|
||||
* (audit round-2 VM #10), then runs the connector and observes the
|
||||
* resulting listener.
|
||||
*/
|
||||
private fun launchConnect(triggerRetryOnFailure: Boolean) {
|
||||
_uiState.update { it.copy(connection = ConnectionUiState.Connecting(ConnectionUiState.Step.ResolvingRoom)) }
|
||||
|
||||
val priorClose = pendingCloseJob
|
||||
pendingCloseJob = null
|
||||
|
||||
connectJob =
|
||||
viewModelScope.launch {
|
||||
try {
|
||||
// Drain the previous listener's close before opening a
|
||||
// new transport. This is fast (< 100 ms typical: send
|
||||
// UNSUBSCRIBE / UNANNOUNCE / WT_CLOSE_SESSION + drain).
|
||||
if (priorClose?.isActive == true) {
|
||||
runCatching { priorClose.join() }
|
||||
}
|
||||
val l =
|
||||
connector.connect(
|
||||
httpClient = httpClient,
|
||||
@@ -489,7 +502,7 @@ class AudioRoomViewModel(
|
||||
_uiState.update {
|
||||
it.copy(connection = ConnectionUiState.Failed(t.message ?: t::class.simpleName ?: "connect failed"))
|
||||
}
|
||||
scheduleAutoRetry()
|
||||
if (triggerRetryOnFailure) scheduleAutoRetry()
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -629,7 +642,11 @@ class AudioRoomViewModel(
|
||||
// before the QUIC transport drops. Same scope rule as the
|
||||
// player stop above: viewModelScope when alive, cleanupScope
|
||||
// for onCleared so the wire teardown survives.
|
||||
scope.launch { runCatching { l.close() } }
|
||||
// Track the close so the next connect() can await it (audit
|
||||
// round-2 VM #10). If `finalCleanup` is true we don't bother
|
||||
// — the VM is going away and no future connect() will run.
|
||||
val closeLaunch = scope.launch { runCatching { l.close() } }
|
||||
if (!finalCleanup) pendingCloseJob = closeLaunch
|
||||
}
|
||||
_uiState.update {
|
||||
it.copy(
|
||||
@@ -787,6 +804,13 @@ sealed class BroadcastUiState {
|
||||
|
||||
data class Broadcasting(
|
||||
val isMuted: Boolean,
|
||||
/**
|
||||
* Non-null when the most recent mute toggle failed (e.g. handle
|
||||
* setMuted threw). UI surfaces this as an inline message; the
|
||||
* broadcast itself stays running with the previous mute state.
|
||||
* Cleared on the next successful mute toggle.
|
||||
*/
|
||||
val muteError: String? = null,
|
||||
) : BroadcastUiState()
|
||||
|
||||
data class Failed(
|
||||
|
||||
@@ -834,10 +834,11 @@ class MoqSession private constructor(
|
||||
.onSuccess { anySent = true }
|
||||
}
|
||||
// Roll the objectId back if the entire fan-out failed (e.g.
|
||||
// transport down). Audio-rooms NIP wants strictly contiguous
|
||||
// object ids per group; a gap from a fully-failed send would
|
||||
// trip strict subscribers. We roll back only if the next send
|
||||
// hasn't already grabbed an id past us.
|
||||
// transport down). The audio-rooms NIP wants object ids
|
||||
// monotonic; gaps from real network loss are acceptable per
|
||||
// spec, but a gap caused by a fully-failed local send is
|
||||
// pure noise we can avoid. We only roll back if no
|
||||
// concurrent send has already grabbed an id past us.
|
||||
if (!anySent) {
|
||||
stateMutex.withLock {
|
||||
if (nextObjectId == objectId + 1) nextObjectId = objectId
|
||||
|
||||
Reference in New Issue
Block a user