fix(cw): make waterfall revision atomic and keep clear from resurrecting old data

Two races shared the same cause: pushSamples runs on the audio capture thread
while clear() runs on the Compose main thread.

1. Lost redraw notifications

Both paths did `_revision.value += 1`. That expands to get -> add -> set and is
not atomic. A controlled two-thread probe (20k increments each, five runs)
lost up to 6,402 increments / 16%; using StateFlow.update lost zero. Since
revision is the Canvas's only redraw signal, every lost update can leave the
waterfall showing stale rows. If both writes land on the same number, StateFlow
sees no value change and notifies nobody.

Use `_revision.update { it + 1 }` in both paths.

2. Clear resurrected pre-clear audio

pushSamples copies pending audio under the lock, deliberately performs FFT
outside it, then reacquires the lock to append rows. The exact interleaving:

  audio thread: take old audio, start FFT
  main thread:  user taps Clear -> rows/pending empty
  audio thread: old FFT completes -> appends old rows again

The display becomes empty then immediately redraws the audio the user cleared.
A deterministic thread probe reproduced old rows after clear. Add a generation
counter protected by the same lock: pushSamples records it before FFT and drops
the computed rows when clear incremented it meanwhile. Fixed probe remains empty.

Verification: :feature:cw:compileReleaseKotlin + full :core:domain:test BUILD
SUCCESSFUL; grep confirms no non-atomic revision increments remain.
This commit is contained in:
mckero committed 2026-08-14 16:31:25 +00:00
1 parent a7a6d70d40
commit 19ca5205fb
1 file changed
+17 -2
@@ -31,6 +31,7 @@ import com.rtbishop.look4sat.core.domain.cw.CwDeepSpectrogram
import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.flow.update
/** /**
* Rolling spectrogram history for the waterfall display. * Rolling spectrogram history for the waterfall display.
@@ -46,6 +47,10 @@ class CwWaterfallState(private val historyRows: Int = 96) {
private val lock = Any() private val lock = Any()
private val rows = ArrayDeque<FloatArray>(historyRows) private val rows = ArrayDeque<FloatArray>(historyRows)
private val pending = ArrayList<Float>(CwDeepSpectrogram.SAMPLE_RATE) private val pending = ArrayList<Float>(CwDeepSpectrogram.SAMPLE_RATE)
// Incremented by clear(). A pushSamples call records the generation before
// doing FFT outside the lock and discards its result if a clear occurred in
// the meantime, otherwise pre-clear audio would reappear after the button tap.
private var generation = 0L
/** /**
* Bumped on every change so Compose knows to redraw. * Bumped on every change so Compose knows to redraw.
@@ -66,6 +71,7 @@ class CwWaterfallState(private val historyRows: Int = 96) {
chunk, sampleRate, CwDeepSpectrogram.SAMPLE_RATE chunk, sampleRate, CwDeepSpectrogram.SAMPLE_RATE
) )
val audio: FloatArray val audio: FloatArray
val generationAtStart: Long
synchronized(lock) { synchronized(lock) {
pending.ensureCapacity(pending.size + resampled.size) pending.ensureCapacity(pending.size + resampled.size)
for (sample in resampled) pending.add(sample) for (sample in resampled) pending.add(sample)
@@ -73,25 +79,34 @@ class CwWaterfallState(private val historyRows: Int = 96) {
if (pending.size < CwDeepSpectrogram.FFT_LENGTH) return if (pending.size < CwDeepSpectrogram.FFT_LENGTH) return
audio = FloatArray(pending.size) { pending[it] } audio = FloatArray(pending.size) { pending[it] }
pending.clear() pending.clear()
generationAtStart = generation
} }
// FFT outside the lock; only the append below needs exclusivity. // FFT outside the lock; only the append below needs exclusivity.
val computed = CwDeepSpectrogram.compute(audio) val computed = CwDeepSpectrogram.compute(audio)
synchronized(lock) { synchronized(lock) {
// Drop the result when the user cleared the display while this FFT
// was running: those samples belong to the discarded history.
if (generation != generationAtStart) return
for (row in computed) { for (row in computed) {
if (rows.size >= historyRows) rows.removeFirst() if (rows.size >= historyRows) rows.removeFirst()
rows.addLast(row) rows.addLast(row)
} }
} }
_revision.value += 1 // update {} not `value += 1`: this runs on the audio capture thread while
// clear() runs on the main thread, and `+=` is a non-atomic
// read-modify-write. A lost increment means a dropped redraw, and two
// writes landing on the same value make StateFlow report no change at all.
_revision.update { it + 1 }
} }
fun clear() { fun clear() {
synchronized(lock) { synchronized(lock) {
generation++
rows.clear() rows.clear()
pending.clear() pending.clear()
} }
_revision.value += 1 _revision.update { it + 1 }
} }
} }