flexaudio 0.3.0

Flexible cross-platform audio capture (microphone, system loopback, per-process) for Linux, Windows, and macOS with a unified API.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
1059
1060
1061
1062
1063
1064
1065
1066
1067
1068
1069
1070
1071
1072
1073
1074
1075
1076
1077
1078
1079
1080
1081
1082
1083
1084
1085
1086
1087
1088
1089
1090
1091
1092
1093
1094
1095
1096
1097
1098
1099
1100
1101
1102
1103
1104
1105
1106
1107
1108
1109
1110
1111
1112
1113
1114
1115
1116
1117
1118
1119
1120
1121
1122
1123
1124
1125
1126
1127
1128
1129
1130
1131
1132
1133
1134
1135
1136
1137
1138
1139
1140
1141
1142
1143
1144
1145
1146
1147
1148
1149
1150
1151
1152
1153
1154
1155
1156
1157
1158
1159
1160
1161
1162
1163
1164
1165
1166
1167
1168
1169
1170
1171
1172
1173
1174
1175
1176
1177
1178
1179
1180
1181
1182
1183
1184
1185
1186
1187
1188
1189
1190
1191
1192
1193
1194
1195
1196
1197
1198
1199
1200
1201
1202
1203
1204
1205
1206
1207
1208
1209
1210
1211
1212
1213
1214
1215
1216
1217
1218
1219
1220
1221
1222
1223
1224
1225
1226
1227
1228
1229
1230
1231
1232
1233
1234
1235
1236
1237
1238
1239
1240
1241
1242
1243
1244
1245
1246
1247
1248
1249
1250
1251
1252
1253
1254
1255
1256
1257
1258
1259
1260
1261
1262
1263
1264
1265
1266
1267
1268
1269
1270
1271
1272
1273
1274
1275
1276
1277
1278
1279
1280
1281
1282
1283
1284
1285
1286
1287
1288
1289
1290
1291
1292
1293
1294
1295
1296
1297
1298
1299
1300
1301
1302
1303
1304
1305
1306
1307
1308
1309
1310
1311
1312
1313
1314
1315
1316
1317
1318
1319
1320
1321
1322
1323
1324
1325
1326
1327
1328
1329
1330
1331
1332
1333
1334
1335
1336
1337
1338
1339
1340
1341
1342
1343
1344
1345
1346
1347
1348
1349
1350
1351
1352
1353
1354
1355
1356
1357
1358
1359
1360
1361
1362
1363
1364
1365
1366
1367
1368
1369
1370
1371
1372
1373
1374
1375
1376
1377
1378
1379
1380
1381
1382
1383
1384
1385
1386
1387
1388
1389
1390
1391
1392
1393
1394
1395
1396
1397
1398
1399
1400
1401
1402
1403
1404
1405
1406
1407
1408
1409
1410
1411
1412
1413
1414
1415
1416
1417
1418
1419
1420
1421
1422
1423
//! A composite backend that mixes mic + system into one stream: [`CompositeBackend`].
//!
//! It owns two child backends, mic and system, normalizes each child's audio to the
//! internal canonical form (48 kHz/stereo), then sums them with per-side gain. It
//! appears as a single backend to [`Stream`](crate::Stream). The Stream itself is
//! unchanged, so seq/PTS, the watchdog, pause, global gain, and switch_source all work
//! as before.
//!
//! # Thread layout
//! - Child backend RT threads: each only pushes to its dedicated child RawRing (the
//!   existing backends remain untouched).
//! - Mixer thread (one, `flexaudio-mix`): pops child rings → converts each with its
//!   [`Normalizer`] to 48 kHz/stereo → sums aligned frames using per-side gain (clamped
//!   to ±1.0) → pushes to the real sink. Mic and system use separate clocks and can
//!   differ by several to hundreds of ppm, so drift correction ([`LinearStitcher`] +
//!   [`DriftController`]) compensates by slightly resampling only the system side. This
//!   is not an RT thread, so heap allocation is allowed (but scratch buffers are reused
//!   to avoid steady-state allocations in the loop).

use std::panic::AssertUnwindSafe;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::thread::{self, JoinHandle};
use std::time::{Duration, Instant};

use flexaudio_core::backend::{CaptureBackend, RawSink};
use flexaudio_core::clock::monotonic_now_ns;
use flexaudio_core::normalizer::Normalizer;
use flexaudio_core::raw_ring::{raw_ring, RawConsumer};
use flexaudio_core::types::{Error, OutputFormat, Result, CHANNELS, SAMPLE_RATE};

use crate::stream::RAW_RING_SAMPLES;

/// If one side supplies nothing for longer than this, continue mixing and fill the
/// missing samples with silence (0.0). Normalization emits only 20 ms chunks, so arrival
/// jitter of 2–3 chunks is considered normal. Continue the overall recording even if
/// supply stops for longer, for example when nothing is playing on the system side.
const STARVATION_FILL_THRESHOLD: Duration = Duration::from_millis(60);

/// Maximum size of one side's normalized FIFO, in f32 samples. This is a safety limit:
/// discard the oldest samples above about 500 ms (48 kHz × 2 channels × 0.5 s = 48,000).
/// Drift correction ([`DriftController`]) handles rate differences between child clocks,
/// so this limit should not normally be reached while correction is working. Keep it as
/// a last resort for anomalies beyond its ±500 ppm range, such as runaway supply on one
/// side.
const FIFO_MAX_SAMPLES: usize = 48_000;

/// How long to wait when neither side has data to mix (same approach as stream.rs's
/// capture thread).
const IDLE_SLEEP: Duration = Duration::from_millis(2);

/// Range in which drift correction can adjust the system-side read ratio r (1.0 ± 500
/// ppm). Consumer-device clock drift is usually tens to a few hundred ppm, so this
/// covers it. Within this range, linear interpolation distortion is negligible because
/// interpolation points remain very close to the source samples.
const DRIFT_RATIO_LIMIT: f64 = 500e-6;

/// Interval for reevaluating the ratio, in f32 samples of mixed output: once per 100 ms
/// of output. This is five 20 ms normalization chunks, coarser than the arrival
/// granularity but much finer than the drift timescale (on the order of minutes).
const DRIFT_UPDATE_INTERVAL_SAMPLES: usize = (SAMPLE_RATE as usize / 10) * CHANNELS as usize;

/// EMA coefficient for the FIFO level difference. With a 100 ms update interval, the
/// time constant is about one second. It smooths chunk-arrival jitter (sawtooth FIFO
/// changes at 20 ms granularity) while still tracking drift changes.
const DRIFT_EMA_ALPHA: f64 = 0.1;

/// P-control gain. A 200 ms FIFO difference (19,200 samples) reaches the 500 ppm
/// clamp limit. This gentle gain keeps the steady-state difference for real drift
/// (tens to hundreds of ppm) at tens to a little over a hundred milliseconds (drift ÷
/// gain), well below the 500 ms safety limit.
const DRIFT_GAIN: f64 = DRIFT_RATIO_LIMIT / 19_200.0;

/// Maximum ratio change per update (slew): 20 ppm per 100 ms, so even a full clamp-range
/// change takes five seconds. This prevents FIFO measurement noise from abruptly
/// changing the ratio and causing pitch fluctuations.
const DRIFT_SLEW_PER_UPDATE: f64 = 20e-6;

/// Composite backend that sums the mic and system child backends in the internal
/// canonical form.
///
/// [`native_format`](CaptureBackend::native_format) always returns the internal
/// canonical form `(48000, 2)`, making Stream's first resampler stage effectively a
/// pass-through. Children are injected through the constructor (tests can pass mocks).
/// The facade's `build_backend` constructs the real children.
pub(crate) struct CompositeBackend {
    mic: Box<dyn CaptureBackend>,
    system: Box<dyn CaptureBackend>,
    mic_gain: f32,
    system_gain: f32,
    /// Stop signal for the mixer thread. Replace it with a new Arc on every start so it
    /// cannot be confused with a leftover from an old thread.
    stopping: Arc<AtomicBool>,
    /// Mixer thread handle. `Some` means it is running.
    mixer: Option<JoinHandle<()>>,
}

impl CompositeBackend {
    /// Create with two injected children and per-side gains. The caller (facade's
    /// `build_backend`) must validate that gains are finite and non-negative.
    pub(crate) fn new(
        mic: Box<dyn CaptureBackend>,
        system: Box<dyn CaptureBackend>,
        mic_gain: f32,
        system_gain: f32,
    ) -> Self {
        Self {
            mic,
            system,
            mic_gain,
            system_gain,
            stopping: Arc::new(AtomicBool::new(false)),
            mixer: None,
        }
    }
}

impl CaptureBackend for CompositeBackend {
    fn native_format(&self) -> (u32, u16) {
        // Mixing always uses the internal canonical form. Stream's first stage is
        // effectively a pass-through.
        (SAMPLE_RATE, CHANNELS)
    }

    fn start(&mut self, sink: RawSink) -> Result<()> {
        // A second start while running is a no-op (CaptureBackend contract).
        if self.mixer.is_some() {
            return Ok(());
        }

        // Return immediately if mic fails to start. If system fails, stop mic before
        // returning the error; never report success with only one side running.
        let mic_lane = start_child(&mut self.mic)?;
        let system_lane = match start_child(&mut self.system) {
            Ok(lane) => lane,
            Err(e) => {
                stop_child(&mut self.mic);
                return Err(e);
            }
        };

        // Start the mixer thread. Create a fresh stop flag on every start so the flag
        // from a previous stop is not reused.
        self.stopping = Arc::new(AtomicBool::new(false));
        let stopping = self.stopping.clone();
        let mic_gain = self.mic_gain;
        let system_gain = self.system_gain;
        let mixer = thread::Builder::new()
            .name("flexaudio-mix".into())
            .spawn(move || {
                run_mixer(mic_lane, system_lane, mic_gain, system_gain, sink, stopping);
            })
            .map_err(|e| Error::Backend(format!("spawn mix thread: {e}")));
        match mixer {
            Ok(handle) => {
                self.mixer = Some(handle);
                Ok(())
            }
            Err(e) => {
                // If the thread cannot start, stop both children and return an error;
                // never leave only one side running.
                stop_child(&mut self.mic);
                stop_child(&mut self.system);
                Err(e)
            }
        }
    }

    fn stop(&mut self) {
        // Set the stop flag → join the mixer thread → stop both children. This is
        // idempotent: if not running, only the child stops run, and those are idempotent
        // by contract too.
        self.stopping.store(true, Ordering::SeqCst);
        if let Some(h) = self.mixer.take() {
            let _ = h.join();
        }
        stop_child(&mut self.mic);
        stop_child(&mut self.system);
    }
}

impl Drop for CompositeBackend {
    fn drop(&mut self) {
        // Do not leave the mixer thread or children running if dropped without stop.
        self.stop();
    }
}

/// All capture state for one child side (child ring consumer, normalizer, and
/// normalized FIFO).
struct ChildLane {
    consumer: RawConsumer,
    /// Child native format → internal canonical form (48 kHz/stereo). Since output is
    /// fixed to the canonical form, the second stage is a pass-through.
    normalizer: Normalizer,
    /// FIFO of normalized samples (48 kHz/stereo interleaved).
    fifo: Vec<f32>,
    /// Time when this side last supplied normalized samples (for starvation detection).
    last_supply: Instant,
}

impl ChildLane {
    /// Pop from the child ring, normalize, and append completed output to the FIFO.
    ///
    /// If the FIFO exceeds [`FIFO_MAX_SAMPLES`], discard the oldest samples as a guard
    /// against unbounded growth. Return `Err` on a rubato processing failure so the
    /// caller can stop the mixer thread.
    fn ingest(&mut self, scratch: &mut [f32]) -> Result<()> {
        let got = self.consumer.pop_slice(scratch);
        if got == 0 {
            return Ok(());
        }
        // The sink does not use pts (the wiring layer handles it separately by
        // contract), but pass monotonic now as the normalizer's anchor.
        self.normalizer.push(&scratch[..got], monotonic_now_ns())?;
        let mut supplied = false;
        while let Some((chunk, _pts)) = self.normalizer.pop_chunk() {
            self.fifo.extend_from_slice(&chunk);
            supplied = true;
        }
        if supplied {
            self.last_supply = Instant::now();
            if self.fifo.len() > FIFO_MAX_SAMPLES {
                let excess = self.fifo.len() - FIFO_MAX_SAMPLES;
                self.fifo.drain(..excess);
            }
        }
        Ok(())
    }

    /// Whether this side has supplied nothing for at least
    /// [`STARVATION_FILL_THRESHOLD`].
    fn is_starved(&self, now: Instant) -> bool {
        now.duration_since(self.last_supply) >= STARVATION_FILL_THRESHOLD
    }
}

/// Fine resampler for the system lane (linear interpolation stitcher).
///
/// Use mic as the reference clock (it is natural to align the recording timeline with the
/// person's voice on the mic side). Read only the system FIFO at ratio r using linear
/// interpolation to compensate for the rate difference between child clocks. Keep only
/// the fractional read position as state.
struct LinearStitcher {
    /// Fractional read position from the FIFO's first frame (in frames, [0, 1)).
    frac: f64,
}

impl LinearStitcher {
    fn new() -> Self {
        Self { frac: 0.0 }
    }

    /// Number of output frames that can be interpolated at `ratio` when the FIFO has
    /// `fifo_frames` frames. The zero-based k-th output uses the two frames around
    /// position `frac + k×ratio`, so output is possible only while that position does
    /// not exceed the final frame F-1 (the exact final frame is valid because the right-hand
    /// interpolation term has zero weight).
    fn producible(&self, fifo_frames: usize, ratio: f64) -> usize {
        if fifo_frames == 0 {
            return 0;
        }
        let span = (fifo_frames - 1) as f64 - self.frac;
        if span < 0.0 {
            return 0;
        }
        (span / ratio) as usize + 1
    }

    /// Read `out_frames` from the system FIFO using linear interpolation at `ratio`,
    /// appending to `out` in interleaved form. Discard consumed frames from the FIFO and
    /// carry the fractional position forward so phase remains continuous across block
    /// boundaries. The caller must ensure `out_frames <= producible(...)`.
    fn pull(&mut self, fifo: &mut Vec<f32>, ratio: f64, out_frames: usize, out: &mut Vec<f32>) {
        let ch = CHANNELS as usize;
        let frames = fifo.len() / ch;
        debug_assert!(out_frames <= self.producible(frames, ratio));
        for k in 0..out_frames {
            let pos = self.frac + k as f64 * ratio;
            let left = pos as usize;
            // Clamp the right edge only when the position is exactly on the final frame
            // (the weight is then zero, so the interpolation result is unchanged).
            let right = (left + 1).min(frames - 1);
            let w = (pos - left as f64) as f32;
            for c in 0..ch {
                let a = fifo[left * ch + c];
                let b = fifo[right * ch + c];
                out.push(a + w * (b - a));
            }
        }
        // Compute the next read position. Discard consumed whole frames and keep only
        // the fractional part.
        let end = self.frac + out_frames as f64 * ratio;
        let consumed = (end as usize).min(frames);
        self.frac = end - consumed as f64;
        fifo.drain(..consumed * ch);
    }

    /// Reset phase when the starvation path flushes the entire system FIFO.
    fn reset(&mut self) {
        self.frac = 0.0;
    }
}

/// Feedback controller that sets read ratio r from the FIFO level difference (EMA + P
/// control + slew).
///
/// Every [`DRIFT_UPDATE_INTERVAL_SAMPLES`] of mixed output, take an exponential moving
/// average of the post-consumption FIFO level difference (mic - system). If the
/// difference is negative (system is accumulating), increase r to consume it faster; if
/// positive, decrease r. P control is sufficient because the level difference itself
/// integrates the rate difference, so proportional correction balances it at a finite
/// level.
struct DriftController {
    /// Current read ratio r, clamped around 1.0 by ±[`DRIFT_RATIO_LIMIT`].
    ratio: f64,
    /// Exponential moving average of the level difference (mic_len - system_len), in
    /// f32 samples.
    ema_diff: f64,
    /// Mixed output samples since the last update.
    pending_samples: usize,
}

impl DriftController {
    fn new() -> Self {
        Self {
            ratio: 1.0,
            ema_diff: 0.0,
            pending_samples: 0,
        }
    }

    /// Record mixed output progress and update the ratio once it reaches
    /// [`DRIFT_UPDATE_INTERVAL_SAMPLES`]. Pass post-consumption levels; pre-consumption
    /// levels would include the samples consumed in this iteration in the difference.
    fn on_output(&mut self, samples: usize, mic_len: usize, system_len: usize) {
        self.pending_samples += samples;
        if self.pending_samples >= DRIFT_UPDATE_INTERVAL_SAMPLES {
            self.pending_samples = 0;
            self.update(mic_len, system_len);
        }
    }

    /// Update the EMA of the level difference and adjust the ratio by one step using P
    /// control + slew.
    fn update(&mut self, mic_len: usize, system_len: usize) {
        let diff = mic_len as f64 - system_len as f64;
        self.ema_diff += DRIFT_EMA_ALPHA * (diff - self.ema_diff);
        let target = (1.0 - DRIFT_GAIN * self.ema_diff)
            .clamp(1.0 - DRIFT_RATIO_LIMIT, 1.0 + DRIFT_RATIO_LIMIT);
        let step = (target - self.ratio).clamp(-DRIFT_SLEW_PER_UPDATE, DRIFT_SLEW_PER_UPDATE);
        self.ratio += step;
    }
}

/// Drift correction state (stitcher + ratio controller), held once per mixer thread.
struct DriftCorrection {
    stitcher: LinearStitcher,
    controller: DriftController,
}

impl DriftCorrection {
    fn new() -> Self {
        Self {
            stitcher: LinearStitcher::new(),
            controller: DriftController::new(),
        }
    }
}

/// Start one child: create a dedicated child RawRing (same capacity as stream.rs) and
/// call `start` with a [`RawSink`] in the child's native format. Return a [`ChildLane`]
/// on success.
///
/// Convert a panic from the child's `start` into [`Error::Backend`] with catch_unwind
/// (same intent as stream.rs's start_backend_catching: prevent cascading panics in the
/// mixer thread or caller).
fn start_child(child: &mut Box<dyn CaptureBackend>) -> Result<ChildLane> {
    let (rate, channels) = child.native_format();
    if rate == 0 || channels == 0 {
        return Err(Error::InvalidArg(
            "mix child native_format must have non-zero rate and channels".into(),
        ));
    }
    let (producer, consumer) = raw_ring(RAW_RING_SAMPLES);
    let sink = RawSink::new(producer, rate, channels);
    match std::panic::catch_unwind(AssertUnwindSafe(|| child.start(sink))) {
        Ok(Ok(())) => {}
        Ok(Err(e)) => return Err(e),
        Err(_) => return Err(Error::Backend("mix child panicked during start()".into())),
    }
    // Child native format → internal canonical form (48 kHz/stereo). The Normalizer's
    // output is fixed to the canonical form, so the second stage is always pass-through.
    let normalizer = Normalizer::new(
        rate,
        channels,
        OutputFormat {
            sample_rate: SAMPLE_RATE,
            channels: CHANNELS,
        },
    )
    .inspect_err(|_| {
        // If the normalizer cannot be created, stop the child before returning an error
        // so no started child is left running.
        stop_child(child);
    })?;
    Ok(ChildLane {
        consumer,
        normalizer,
        fifo: Vec::with_capacity(FIFO_MAX_SAMPLES),
        last_supply: Instant::now(),
    })
}

/// Call the child's `stop` inside catch_unwind so a panic does not propagate (same
/// intent as stream.rs's stop_backend_catching).
fn stop_child(child: &mut Box<dyn CaptureBackend>) {
    let _ = std::panic::catch_unwind(AssertUnwindSafe(|| child.stop()));
}

/// Mixer thread entry point.
///
/// First wait for initial supply from both sides with [`prime_lanes`] (up to the
/// starvation threshold), then ingest each child (pop → normalize → FIFO), sum aligned
/// frames using per-side gain, and push them to the real sink. If one side supplies
/// nothing for at least [`STARVATION_FILL_THRESHOLD`], fill its missing data with
/// silence and continue (recording keeps flowing even when the system side is silent).
/// If neither side has data, sleep for [`IDLE_SLEEP`].
///
/// A normalization failure (a theoretical rubato failure) ends the loop. No more samples
/// will flow, so Stream's watchdog detects the stall and reopens the backend.
fn run_mixer(
    mut mic: ChildLane,
    mut system: ChildLane,
    mic_gain: f32,
    system_gain: f32,
    mut sink: RawSink,
    stopping: Arc<AtomicBool>,
) {
    // Scratch buffers for popping (child ring capacity) and mixed output. Reuse them in
    // the loop.
    let mut scratch = vec![0.0f32; RAW_RING_SAMPLES];
    let mut mixed: Vec<f32> = Vec::with_capacity(FIFO_MAX_SAMPLES);
    // Drift correction state for child clocks (fresh on each start so the previous
    // recording's ratio is not carried over).
    let mut drift = DriftCorrection::new();

    // Child threads may start at different times, so wait for both sides to begin
    // flowing (up to the starvation threshold) before mixing. This prevents the start
    // of a recording from containing only one side.
    if !prime_lanes(&mut mic, &mut system, &mut scratch, &stopping) {
        return;
    }

    loop {
        if stopping.load(Ordering::SeqCst) {
            break;
        }

        if mic.ingest(&mut scratch).is_err() || system.ingest(&mut scratch).is_err() {
            // Stop mixing if normalization fails; let the watchdog reopen the backend.
            return;
        }

        let pushed = mix_and_push(
            &mut mic,
            &mut system,
            mic_gain,
            system_gain,
            &mut drift,
            &mut sink,
            &mut mixed,
        );

        if !pushed {
            thread::sleep(IDLE_SLEEP);
        }
    }
}

/// Prime the lanes before mixing. Child backend threads may start at different times,
/// so poll with [`IDLE_SLEEP`] until both FIFOs receive their first normalized samples.
/// Without this, starvation fill could start while the slower side's ring is empty,
/// leaving only one side at the start of the recording.
///
/// The wait is limited by [`STARVATION_FILL_THRESHOLD`]. If one side never supplies
/// data (for example, nothing is playing on the system side), begin mixing after the
/// same timeout used for ordinary starvation. Stop waiting when a stop signal arrives.
/// Return `false` on normalization failure so the caller can end the mixer thread.
fn prime_lanes(
    mic: &mut ChildLane,
    system: &mut ChildLane,
    scratch: &mut [f32],
    stopping: &AtomicBool,
) -> bool {
    let start = Instant::now();
    while !stopping.load(Ordering::SeqCst) {
        if mic.ingest(scratch).is_err() || system.ingest(scratch).is_err() {
            return false;
        }
        if (!mic.fifo.is_empty() && !system.fifo.is_empty())
            || start.elapsed() >= STARVATION_FILL_THRESHOLD
        {
            break;
        }
        thread::sleep(IDLE_SLEEP);
    }
    true
}

/// Take the mixable amount from both FIFOs, sum with per-side gain (clamped to ±1.0),
/// and push to the sink. Return `true` if anything was pushed.
///
/// How much to take:
/// - Data on both sides (steady state): the minimum of the mic inventory and the amount
///   interpolable from system. Mic is consumed at its native rate; [`LinearStitcher`]
///   reads system at ratio r using linear interpolation to compensate for child-clock
///   drift. [`DriftController`] adjusts r based on post-consumption levels.
/// - Data on one side and starvation on the other (no supply for at least 60 ms): mix
///   all available data against silence (fill missing samples with 0.0). Do not apply
///   correction on this path; retain the existing starvation semantics.
/// - Otherwise (both empty and the other side is not yet starved): do nothing and wait
///   for data to align.
fn mix_and_push(
    mic: &mut ChildLane,
    system: &mut ChildLane,
    mic_gain: f32,
    system_gain: f32,
    drift: &mut DriftCorrection,
    sink: &mut RawSink,
    mixed: &mut Vec<f32>,
) -> bool {
    let ch = CHANNELS as usize;
    let ratio = drift.controller.ratio;
    let steady_frames =
        (mic.fifo.len() / ch).min(drift.stitcher.producible(system.fifo.len() / ch, ratio));
    if steady_frames > 0 {
        // Steady-state path: read the system side at ratio r with linear interpolation,
        // then mix the sides.
        mixed.clear();
        drift
            .stitcher
            .pull(&mut system.fifo, ratio, steady_frames, mixed);
        let count = steady_frames * ch;
        for (i, out) in mixed.iter_mut().enumerate() {
            *out = (mic.fifo[i] * mic_gain + *out * system_gain).clamp(-1.0, 1.0);
        }
        mic.fifo.drain(..count);
        drift
            .controller
            .on_output(count, mic.fifo.len(), system.fifo.len());
        // As in stream.rs capture, monotonic now is sufficient for pts (the sink handles
        // it separately by contract).
        sink.push(mixed, monotonic_now_ns());
        return true;
    }

    // Starvation path (no correction; retain existing semantics).
    let now = Instant::now();
    let (mic_take, system_take) = if !mic.fifo.is_empty() && system.is_starved(now) {
        // System stopped supplying: output all mic data. Flush any remaining fraction
        // (at most one frame) that lacked a right interpolation frame, then reset phase.
        drift.stitcher.reset();
        (mic.fifo.len(), system.fifo.len())
    } else if mic.fifo.is_empty() && !system.fifo.is_empty() && mic.is_starved(now) {
        drift.stitcher.reset();
        (0, system.fifo.len())
    } else {
        return false;
    };

    let count = mic_take.max(system_take);
    mixed.clear();
    for i in 0..count {
        let m = if i < mic_take { mic.fifo[i] } else { 0.0 };
        let s = if i < system_take { system.fifo[i] } else { 0.0 };
        mixed.push((m * mic_gain + s * system_gain).clamp(-1.0, 1.0));
    }
    mic.fifo.drain(..mic_take);
    system.fifo.drain(..system_take);

    // As in stream.rs capture, monotonic now is sufficient for pts (the sink handles it
    // separately by contract).
    sink.push(mixed, monotonic_now_ns());
    true
}

#[cfg(test)]
mod tests {
    use super::*;
    use flexaudio_core::raw_ring::raw_ring;
    use std::sync::atomic::AtomicU32;

    /// Test-only child backend that supplies constant-amplitude (DC) samples as fast as
    /// possible.
    ///
    /// Sine waves (MockBackend) introduce phase issues when mixing two sources, so use
    /// DC signals for deterministic verification of the mix. If `feed_for` is set, feed
    /// for that duration and then stop pushing while keeping the thread alive, to
    /// reproduce starvation on one side.
    ///
    /// Supply uses a saturating strategy rather than real-time pacing (10 ms sleeps):
    /// push as fast as the child ring accepts data, yielding for 1 ms before retrying if
    /// it is full. Real-time pacing can fall behind due to coarse sleep granularity on a
    /// slow test machine. When the mixer wakes, one FIFO may be empty, causing a
    /// starvation-filled single-side chunk to make the mix-value check scheduler-
    /// dependent. Saturating supply keeps both FIFOs nonempty whenever the mixer wakes,
    /// limiting the check to mix arithmetic (scheduling resilience is covered by
    /// mix_survives_one_side_starvation). With DC, drops and partial writes when full do
    /// not affect the values.
    struct ConstBackend {
        sample_rate: u32,
        channels: u16,
        value: f32,
        feed_for: Option<Duration>,
        running: Arc<AtomicBool>,
        handle: Option<JoinHandle<()>>,
    }

    impl ConstBackend {
        fn new(value: f32, feed_for: Option<Duration>) -> Self {
            Self {
                sample_rate: 48_000,
                channels: 2,
                value,
                feed_for,
                running: Arc::new(AtomicBool::new(false)),
                handle: None,
            }
        }
    }

    impl CaptureBackend for ConstBackend {
        fn native_format(&self) -> (u32, u16) {
            (self.sample_rate, self.channels)
        }

        fn start(&mut self, mut sink: RawSink) -> Result<()> {
            if self.running.load(Ordering::SeqCst) {
                return Ok(());
            }
            self.running.store(true, Ordering::SeqCst);
            let running = self.running.clone();
            let sample_rate = self.sample_rate;
            let channels = self.channels as usize;
            let value = self.value;
            let feed_for = self.feed_for;
            let handle = thread::Builder::new()
                .name("flexaudio-const-gen".into())
                .spawn(move || {
                    let frames_per_block = (sample_rate as usize / 100).max(1); // About 10 ms
                    let block = vec![value; frames_per_block * channels];
                    let start = Instant::now();
                    while running.load(Ordering::SeqCst) {
                        let feeding = feed_for.is_none_or(|d| start.elapsed() < d);
                        if !feeding {
                            // Once the feed duration ends, supply nothing further to
                            // reproduce one-sided starvation. Sleep while waiting for
                            // the stop signal.
                            thread::sleep(Duration::from_millis(5));
                            continue;
                        }
                        // Saturating supply: push again immediately if all data fit;
                        // yield for 1 ms and retry if the ring was full (push is
                        // nonblocking and drops data that does not fit; harmless for DC).
                        let accepted = sink.push(&block, start.elapsed().as_nanos() as i64);
                        if accepted < block.len() {
                            thread::sleep(Duration::from_millis(1));
                        }
                    }
                })
                .map_err(|e| Error::Backend(format!("spawn const gen thread: {e}")))?;
            self.handle = Some(handle);
            Ok(())
        }

        fn stop(&mut self) {
            self.running.store(false, Ordering::SeqCst);
            if let Some(h) = self.handle.take() {
                let _ = h.join();
            }
        }
    }

    impl Drop for ConstBackend {
        fn drop(&mut self) {
            self.stop();
        }
    }

    /// Test-only backend whose `start` always returns Err.
    struct FailingStartBackend;

    impl CaptureBackend for FailingStartBackend {
        fn native_format(&self) -> (u32, u16) {
            (48_000, 2)
        }
        fn start(&mut self, _sink: RawSink) -> Result<()> {
            Err(Error::Backend("intentional start failure".into()))
        }
        fn stop(&mut self) {}
    }

    /// Test-only backend that records start/stop call counts in shared counters.
    struct TrackingBackend {
        starts: Arc<AtomicU32>,
        stops: Arc<AtomicU32>,
    }

    impl CaptureBackend for TrackingBackend {
        fn native_format(&self) -> (u32, u16) {
            (48_000, 2)
        }
        fn start(&mut self, _sink: RawSink) -> Result<()> {
            self.starts.fetch_add(1, Ordering::SeqCst);
            Ok(())
        }
        fn stop(&mut self) {
            self.stops.fetch_add(1, Ordering::SeqCst);
        }
    }

    /// Helper that builds and starts a composite, then returns the real sink consumer.
    fn start_composite(
        mic: Box<dyn CaptureBackend>,
        system: Box<dyn CaptureBackend>,
        mic_gain: f32,
        system_gain: f32,
    ) -> (CompositeBackend, RawConsumer) {
        let mut be = CompositeBackend::new(mic, system, mic_gain, system_gain);
        assert_eq!(
            be.native_format(),
            (48_000, 2),
            "should report the internal canonical form"
        );
        let (producer, consumer) = raw_ring(RAW_RING_SAMPLES);
        let sink = RawSink::new(producer, 48_000, 2);
        be.start(sink).expect("composite start");
        (be, consumer)
    }

    /// Tolerance for value comparisons; enough to absorb f32 rounding when adding two
    /// DC values.
    const VALUE_TOL: f32 = 1e-4;

    /// Minimum sample count to consider a value "actually observed": one internal-form
    /// chunk (960 frames × 2 channels). Set the existence floor as an absolute count,
    /// not a fraction of all samples. A fraction (for example, 25%) can fail in a
    /// scheduler-dependent way when a large starvation-fill burst (up to the 48k-sample
    /// FIFO limit) inflates the denominator.
    const ONE_CHUNK_SAMPLES: usize = 1_920;

    /// Helper that collects samples from the consumer until a condition is met.
    ///
    /// `done` receives only the newly popped samples each time (the caller accumulates
    /// counts, etc. to check the condition). Return all collected samples once it returns
    /// true. A fixed wall-clock window ("N samples in 500 ms") can inherently flake if
    /// load deschedules the threads and prevents them from producing enough in that
    /// window. Instead, wait until the condition is met. `max_wait` prevents an
    /// indefinite wait under extreme load; on timeout, return what was collected (the
    /// caller's assertion detects any shortfall).
    fn collect_until(
        consumer: &mut RawConsumer,
        max_wait: Duration,
        mut done: impl FnMut(&[f32]) -> bool,
    ) -> Vec<f32> {
        let mut out = Vec::new();
        let mut scratch = vec![0.0f32; RAW_RING_SAMPLES];
        let start = Instant::now();
        loop {
            let got = consumer.pop_slice(&mut scratch);
            out.extend_from_slice(&scratch[..got]);
            if done(&scratch[..got]) || start.elapsed() >= max_wait {
                return out;
            }
            thread::sleep(Duration::from_millis(5));
        }
    }

    /// Wait limit for [`collect_until`]. Under normal conditions, it exits as soon as
    /// the condition is met; this only prevents a hang if extreme load barely lets the
    /// threads run.
    const COLLECT_MAX_WAIT: Duration = Duration::from_secs(30);

    /// Helper that counts samples matching `v` within [`VALUE_TOL`].
    fn count_near(samples: &[f32], v: f32) -> usize {
        samples
            .iter()
            .filter(|&&s| (s - v).abs() < VALUE_TOL)
            .count()
    }

    /// Helper that verifies every sample matches one of the theoretically possible
    /// values and returns the count matching the mixed value `mixed`.
    ///
    /// Why check membership in a set instead of a ratio? A ratio assertion such as "98%
    /// of steady-state samples are mixed" can fail under any threshold if load
    /// deschedules producer or mixer threads beyond the 60 ms starvation threshold:
    /// starvation-filled zeros or single-side values can appear in any interval, making
    /// the ratio inherently dependent on wall-clock timing. By contrast, for DC sources
    /// and constant gains, the mixer can produce only four values: mixed (both sides),
    /// mic alone (system starvation fill), system alone (mic starvation fill), or 0.0
    /// (priming boundary). An incorrect sum, misapplied gain, or missing clamp always
    /// produces a value outside this set. The scheduler can change the distribution but
    /// cannot create an out-of-set value, so this check is scheduler-independent.
    fn assert_only_allowed_values(
        samples: &[f32],
        mixed: f32,
        mic_only: f32,
        system_only: f32,
    ) -> usize {
        let allowed = [mixed, mic_only, system_only, 0.0];
        let mut mixed_count = 0usize;
        for (i, &s) in samples.iter().enumerate() {
            if (s - mixed).abs() < VALUE_TOL {
                mixed_count += 1;
            } else {
                assert!(
                    allowed.iter().any(|&a| (s - a).abs() < VALUE_TOL),
                    "value outside allowed set {allowed:?} (mix arithmetic error): samples[{i}] = {s}"
                );
            }
        }
        mixed_count
    }

    /// Mixing two DC sources with known amplitudes (0.2 and 0.3) at mic_gain=1.0 and
    /// system_gain=2.0 yields 0.2×1.0 + 0.3×2.0 = 0.8 (48 kHz/stereo children make
    /// every stage pass-through, so values are deterministic).
    #[test]
    fn mix_sums_two_sources_with_gains() {
        let mic = Box::new(ConstBackend::new(0.2, None));
        let system = Box::new(ConstBackend::new(0.3, None));
        let (mut be, mut consumer) = start_composite(mic, system, 1.0, 2.0);

        // Collect until there is enough data and one chunk of mixed values, waiting
        // through load-related delays.
        let (mut total, mut mixed) = (0usize, 0usize);
        let samples = collect_until(&mut consumer, COLLECT_MAX_WAIT, |new| {
            total += new.len();
            mixed += count_near(new, 0.8);
            total >= 10_000 && mixed >= ONE_CHUNK_SAMPLES
        });
        be.stop();

        assert!(
            samples.len() >= 10_000,
            "should produce enough samples: {}",
            samples.len()
        );
        // Check every sample against the allowed set: mixed value 0.2*1.0 + 0.3*2.0 =
        // 0.8, single-side values 0.2 (mic only) / 0.6 (system only) during starvation
        // fill, and 0.0 at the priming boundary. Even one out-of-set sample means the
        // mix arithmetic is wrong.
        let mixed_count = assert_only_allowed_values(&samples, 0.8, 0.2, 0.6);
        // Also require one chunk's worth to guarantee mixing actually occurred.
        assert!(
            mixed_count >= ONE_CHUNK_SAMPLES,
            "mixed value 0.8 should appear often enough: {mixed_count}/{}",
            samples.len()
        );
    }

    /// A mix that exceeds ±1.0 (0.8 + 0.8 = 1.6) is clamped to 1.0.
    #[test]
    fn mix_clamps_sum() {
        let mic = Box::new(ConstBackend::new(0.8, None));
        let system = Box::new(ConstBackend::new(0.8, None));
        let (mut be, mut consumer) = start_composite(mic, system, 1.0, 1.0);

        // Collect until one chunk of clamped mixed values is available, waiting through
        // load-related delays.
        let mut clamped = 0usize;
        let samples = collect_until(&mut consumer, COLLECT_MAX_WAIT, |new| {
            clamped += count_near(new, 1.0);
            clamped >= ONE_CHUNK_SAMPLES
        });
        be.stop();

        // The key clamp property: no sample exceeds 1.0.
        for (i, &s) in samples.iter().enumerate() {
            assert!(
                s <= 1.0,
                "should not exceed 1.0 after clamping: samples[{i}] = {s}"
            );
        }
        // Check every sample against the allowed set: clamp(0.8 + 0.8) = 1.0, single-
        // side value 0.8 during starvation fill, and 0.0 at the priming boundary. A
        // missed clamp such as 1.6 is immediately rejected as outside the set.
        let clamped_count = assert_only_allowed_values(&samples, 1.0, 0.8, 0.8);
        // Also require one chunk's worth to guarantee the clamped mix actually occurred.
        assert!(
            clamped_count >= ONE_CHUNK_SAMPLES,
            "clamped value 1.0 should appear often enough: {clamped_count}/{}",
            samples.len()
        );
    }

    /// If one side (system) stops supplying data, output continues with silence (0.0)
    /// for the starved side and only mic audio (mic-only value 0.2) begins to flow.
    #[test]
    fn mix_survives_one_side_starvation() {
        let mic = Box::new(ConstBackend::new(0.2, None));
        // System is fed for 150 ms, then stops while its thread remains alive.
        let system = Box::new(ConstBackend::new(0.3, Some(Duration::from_millis(150))));
        let (mut be, mut consumer) = start_composite(mic, system, 1.0, 1.0);

        // After system stops (150 ms wall time), the backlog drains and the starvation
        // threshold passes; mic-only value 0.2 must then start flowing. A wall-clock
        // window check such as "most values are 0.2 after 400 ms" can fail if load slows
        // backlog draining. Instead, wait until one chunk of mic-only values appears.
        let mut mic_only = 0usize;
        let samples = collect_until(&mut consumer, COLLECT_MAX_WAIT, |new| {
            mic_only += count_near(new, 0.2);
            mic_only >= ONE_CHUNK_SAMPLES
        });
        be.stop();

        // Allowed values: mixed 0.5, mic-only 0.2, system-only 0.3, and 0.0 at the
        // priming boundary.
        assert_only_allowed_values(&samples, 0.5, 0.2, 0.3);
        // Also guarantee that output continues after system stops and the mic-only
        // value actually flows with starvation fill.
        let mic_only_count = count_near(&samples, 0.2);
        assert!(
            mic_only_count >= ONE_CHUNK_SAMPLES,
            "mic-only value 0.2 should keep flowing after starvation: {mic_only_count}/{}",
            samples.len()
        );
    }

    /// If the system child's start returns Err, the already-started mic child is stopped
    /// and the composite also returns Err (it must not succeed with only one side).
    #[test]
    fn mix_start_failure_cleans_up() {
        let starts = Arc::new(AtomicU32::new(0));
        let stops = Arc::new(AtomicU32::new(0));
        let mic = Box::new(TrackingBackend {
            starts: starts.clone(),
            stops: stops.clone(),
        });
        let system = Box::new(FailingStartBackend);

        let mut be = CompositeBackend::new(mic, system, 1.0, 1.0);
        let (producer, _consumer) = raw_ring(RAW_RING_SAMPLES);
        let sink = RawSink::new(producer, 48_000, 2);

        let err = be
            .start(sink)
            .expect_err("composite should return Err when system start fails");
        assert!(
            matches!(err, Error::Backend(_)),
            "system child's Err should propagate: {err:?}"
        );
        assert_eq!(starts.load(Ordering::SeqCst), 1, "mic should start once");
        assert_eq!(
            stops.load(Ordering::SeqCst),
            1,
            "mic should stop when system fails"
        );
    }

    /// If the mic child's start returns Err, return Err immediately without touching
    /// the system child. Stop is idempotent and can be called twice.
    #[test]
    fn mix_mic_start_failure_is_immediate() {
        let starts = Arc::new(AtomicU32::new(0));
        let stops = Arc::new(AtomicU32::new(0));
        let mic = Box::new(FailingStartBackend);
        let system = Box::new(TrackingBackend {
            starts: starts.clone(),
            stops: stops.clone(),
        });

        let mut be = CompositeBackend::new(mic, system, 1.0, 1.0);
        let (producer, _consumer) = raw_ring(RAW_RING_SAMPLES);
        let sink = RawSink::new(producer, 48_000, 2);
        assert!(
            be.start(sink).is_err(),
            "mic start failure should return Err immediately"
        );
        assert_eq!(
            starts.load(Ordering::SeqCst),
            0,
            "system should not start if mic fails"
        );

        // Stop is idempotent even if not started (child stop is also idempotent by contract).
        be.stop();
        be.stop();
    }

    /// End-to-end test using the composite backend in a real [`Stream`](crate::Stream).
    /// It verifies that 20 ms/960-frame chunks flow and data has the mixed value (0.2 +
    /// 0.3 = 0.5), confirming the seq and chunk contracts work unchanged without
    /// modifying Stream itself.
    #[test]
    fn stream_delivers_mixed_chunks_end_to_end() {
        use flexaudio_core::types::StreamConfig;

        let mic = Box::new(ConstBackend::new(0.2, None));
        let system = Box::new(ConstBackend::new(0.3, None));
        let backend = Box::new(CompositeBackend::new(mic, system, 1.0, 1.0));
        let mut stream = crate::Stream::open(StreamConfig::default(), backend).expect("open");
        stream.start().expect("start");

        // Poll until one chunk of mixed samples arrives. Fixed wall-clock windows can
        // flake under load, so wait for the condition with a timeout as a hang safeguard.
        let mut chunks = Vec::new();
        let mut mixed = 0usize;
        let deadline = Instant::now() + COLLECT_MAX_WAIT;
        while Instant::now() < deadline && mixed < ONE_CHUNK_SAMPLES {
            while let Some(c) = stream.poll_chunk() {
                mixed += count_near(&c.data, 0.5);
                chunks.push(c);
            }
            thread::sleep(Duration::from_millis(5));
        }
        stream.stop();

        assert!(!chunks.is_empty(), "chunks should arrive");
        for (i, c) in chunks.iter().enumerate() {
            assert_eq!(c.frames, 960, "20ms@48k = 960 frame");
            assert_eq!(c.data.len(), 960 * 2, "stereo interleaved");
            if i > 0 {
                assert!(
                    c.seq > chunks[i - 1].seq,
                    "seq should increase monotonically"
                );
            }
        }
        // Check every sample in every chunk against the allowed set: mixed value 0.2 +
        // 0.3 = 0.5, single-side values 0.2 (mic only) / 0.3 (system only) during
        // starvation fill, and 0.0 at the priming boundary. Stream's first stage is
        // pass-through at 48 kHz/stereo (gain 1.0 leaves bytes unchanged), so mixer
        // output values arrive unchanged.
        let all: Vec<f32> = chunks.iter().flat_map(|c| c.data.iter().copied()).collect();
        let mixed_count = assert_only_allowed_values(&all, 0.5, 0.2, 0.3);
        // Also require one chunk's worth to guarantee that mixing actually occurred.
        assert!(
            mixed_count >= ONE_CHUNK_SAMPLES,
            "mixed value 0.5 should appear often enough: {mixed_count}/{}",
            all.len()
        );
    }

    // ---- Drift correction component tests (pure and deterministic) ----

    /// When system accumulates data (mic - system is negative), r moves above 1.0 to
    /// consume it faster; when mic accumulates data, r moves below 1.0.
    #[test]
    fn drift_controller_moves_toward_lagging_side() {
        let mut c = DriftController::new();
        for _ in 0..50 {
            c.update(0, 9_600);
        }
        assert!(
            c.ratio > 1.0,
            "r should be > 1.0 when system accumulates data: {}",
            c.ratio
        );

        let mut c = DriftController::new();
        for _ in 0..50 {
            c.update(9_600, 0);
        }
        assert!(
            c.ratio < 1.0,
            "r should be < 1.0 when mic accumulates data: {}",
            c.ratio
        );
    }

    /// No matter how large the level difference is, r is capped at 1.0 ± 500 ppm.
    #[test]
    fn drift_controller_clamps_at_ratio_limit() {
        let mut c = DriftController::new();
        for _ in 0..1_000 {
            c.update(0, 10_000_000);
        }
        assert!(
            (c.ratio - (1.0 + DRIFT_RATIO_LIMIT)).abs() < 1e-12,
            "should stop exactly at the upper clamp: {}",
            c.ratio
        );

        let mut c = DriftController::new();
        for _ in 0..1_000 {
            c.update(10_000_000, 0);
        }
        assert!(
            (c.ratio - (1.0 - DRIFT_RATIO_LIMIT)).abs() < 1e-12,
            "should stop exactly at the lower clamp: {}",
            c.ratio
        );
    }

    /// Even with a huge level difference, one update can move only by the slew limit.
    #[test]
    fn drift_controller_slew_limits_change_per_update() {
        let mut c = DriftController::new();
        c.update(0, 10_000_000);
        assert!(
            (c.ratio - (1.0 + DRIFT_SLEW_PER_UPDATE)).abs() < 1e-12,
            "first update should be capped exactly at the slew limit: {}",
            c.ratio
        );
        c.update(0, 10_000_000);
        assert!(
            (c.ratio - (1.0 + 2.0 * DRIFT_SLEW_PER_UPDATE)).abs() < 1e-12,
            "second update should also move by one step: {}",
            c.ratio
        );

        let mut c = DriftController::new();
        c.update(10_000_000, 0);
        assert!(
            (c.ratio - (1.0 - DRIFT_SLEW_PER_UPDATE)).abs() < 1e-12,
            "the reverse direction should also be capped at the slew limit: {}",
            c.ratio
        );
    }

    /// on_output does not update the ratio until 100 ms of mixed output (the update
    /// interval) has accumulated.
    #[test]
    fn drift_controller_updates_only_at_interval() {
        let mut c = DriftController::new();
        c.on_output(DRIFT_UPDATE_INTERVAL_SAMPLES - 1, 0, 10_000_000);
        assert!(
            (c.ratio - 1.0).abs() < 1e-15,
            "should not change before the interval: {}",
            c.ratio
        );
        c.on_output(1, 0, 10_000_000);
        assert!(
            c.ratio > 1.0,
            "should update once the interval is reached: {}",
            c.ratio
        );
    }

    /// At unity speed (r = 1.0, phase 0), the stitcher is a perfect pass-through: output
    /// matches input, the FIFO is fully consumed, and phase remains 0.
    #[test]
    fn stitcher_unity_ratio_is_passthrough() {
        let mut st = LinearStitcher::new();
        let src: Vec<f32> = (0..10).flat_map(|f| [f as f32, -(f as f32)]).collect();
        let mut fifo = src.clone();
        assert_eq!(st.producible(10, 1.0), 10);
        let mut out = Vec::new();
        st.pull(&mut fifo, 1.0, 10, &mut out);
        assert_eq!(out, src, "r=1.0 should be pass-through");
        assert!(
            fifo.is_empty(),
            "should be fully consumed: {} remaining",
            fifo.len()
        );
        assert!(st.frac.abs() < 1e-12, "phase should remain 0: {}", st.frac);
    }

    /// Reading a ramp (frame k has value k) at r = 1.25 yields the linear interpolation
    /// values at positions 0 / 1.25 / 2.5 / 3.75 (both channels, deterministically).
    #[test]
    fn stitcher_interpolates_between_frames() {
        let mut st = LinearStitcher::new();
        let mut fifo: Vec<f32> = (0..5).flat_map(|f| [f as f32, f as f32 * 10.0]).collect();
        // span = 4, floor(4 / 1.25) = 3 → 3 + 1 = 4 frames can be produced.
        assert_eq!(st.producible(5, 1.25), 4);
        let mut out = Vec::new();
        st.pull(&mut fifo, 1.25, 4, &mut out);
        let expect = [0.0f32, 1.25, 2.5, 3.75];
        for (k, &e) in expect.iter().enumerate() {
            assert!(
                (out[k * 2] - e).abs() < 1e-6,
                "interpolated value at position {e}: {}",
                out[k * 2]
            );
            assert!(
                (out[k * 2 + 1] - e * 10.0).abs() < 1e-5,
                "ch2 should interpolate at the same position: {}",
                out[k * 2 + 1]
            );
        }
        // floor(0 + 4×1.25) = 5, so all frames are consumed and phase returns to 0.
        assert!(
            fifo.is_empty(),
            "should be fully consumed: {} remaining",
            fifo.len()
        );
        assert!(st.frac.abs() < 1e-12, "phase: {}", st.frac);
    }

    // ---- Synchronous drift simulation without threads (deterministic) ----

    /// ChildLane for synchronous simulation without threads. The ring and normalizer
    /// exist only to satisfy the shape; the test supplies data directly to the FIFO via
    /// [`sim_feed`].
    fn sim_lane() -> ChildLane {
        let (_producer, consumer) = raw_ring(16);
        let normalizer = Normalizer::new(
            SAMPLE_RATE,
            CHANNELS,
            OutputFormat {
                sample_rate: SAMPLE_RATE,
                channels: CHANNELS,
            },
        )
        .expect("pass-through normalizer");
        ChildLane {
            consumer,
            normalizer,
            fifo: Vec::new(),
            last_supply: Instant::now(),
        }
    }

    /// Directly supply DC frames using the same approach as ingest (append to FIFO and
    /// enforce its safety limit).
    fn sim_feed(lane: &mut ChildLane, value: f32, frames: usize) {
        let new_len = lane.fifo.len() + frames * CHANNELS as usize;
        lane.fifo.resize(new_len, value);
        if lane.fifo.len() > FIFO_MAX_SAMPLES {
            let excess = lane.fifo.len() - FIFO_MAX_SAMPLES;
            lane.fifo.drain(..excess);
        }
        lane.last_supply = Instant::now();
    }

    /// Measurements from the synchronous drift simulation.
    struct DriftSimOutcome {
        /// Post-consumption FIFO levels (f32 samples) for each simulated second.
        system_backlog: Vec<usize>,
        mic_backlog: Vec<usize>,
        final_ratio: f64,
    }

    /// Synchronous simulation without threads. Each tick is 20 ms. Supply DC to mic at
    /// unity rate (960 frames/tick) and to system at (1 + ppm×1e-6) rate, simulating
    /// elapsed time through sample counts, then drive mix_and_push directly in a loop.
    /// It is fully deterministic and independent of threads and wall-clock time.
    ///
    /// If `fixed_unity_ratio` is set, reset the controller to 1.0 every tick, simulating
    /// "no correction." This self-check confirms the simulation reproduces the drift
    /// problem (monotonically growing FIFO levels).
    ///
    /// Check every output value on every tick: both sides are always supplying data, so
    /// only the mixed value (0.2 + 0.3 = 0.5; linear interpolation preserves DC) is
    /// allowed. Fail if even one starvation-filled 0.0 or single-side value appears;
    /// this also verifies that zero-fill does not occur in steady state.
    fn run_drift_sim(ppm: f64, seconds: usize, fixed_unity_ratio: bool) -> DriftSimOutcome {
        const TICK_FRAMES: usize = 960; // 20ms @48k
        const TICKS_PER_SEC: usize = 50;

        let mut mic = sim_lane();
        let mut system = sim_lane();
        let mut drift = DriftCorrection::new();
        let (producer, mut consumer) = raw_ring(RAW_RING_SAMPLES);
        let mut sink = RawSink::new(producer, SAMPLE_RATE, CHANNELS);
        let mut mixed: Vec<f32> = Vec::with_capacity(FIFO_MAX_SAMPLES);
        let mut scratch = vec![0.0f32; RAW_RING_SAMPLES];

        // Carry the fractional system supply frame count forward to simulate the rate
        // difference accurately in samples.
        let mut system_carry = 0.0f64;
        let mut outcome = DriftSimOutcome {
            system_backlog: Vec::new(),
            mic_backlog: Vec::new(),
            final_ratio: 1.0,
        };

        for tick in 0..seconds * TICKS_PER_SEC {
            sim_feed(&mut mic, 0.2, TICK_FRAMES);
            system_carry += TICK_FRAMES as f64 * (1.0 + ppm * 1e-6);
            let system_frames = system_carry as usize;
            system_carry -= system_frames as f64;
            sim_feed(&mut system, 0.3, system_frames);

            mix_and_push(
                &mut mic,
                &mut system,
                1.0,
                1.0,
                &mut drift,
                &mut sink,
                &mut mixed,
            );
            if fixed_unity_ratio {
                drift.controller.ratio = 1.0;
                drift.controller.ema_diff = 0.0;
            }

            let got = consumer.pop_slice(&mut scratch);
            for &s in &scratch[..got] {
                assert!(
                    (s - 0.5).abs() < VALUE_TOL,
                    "steady-state simulation should output only the mixed value 0.5 (no zero-fill or \
                     single-side values): {s}"
                );
            }

            if (tick + 1) % TICKS_PER_SEC == 0 {
                outcome.system_backlog.push(system.fifo.len());
                outcome.mic_backlog.push(mic.fifo.len());
            }
        }
        outcome.final_ratio = drift.controller.ratio;
        outcome
    }

    /// Simulate +300 ppm (system is faster) for 60 seconds: with correction, the system
    /// FIFO stays below the safety limit, grows more slowly, and remains smaller than
    /// without correction. With correction disabled (r fixed at 1.0), levels grow
    /// monotonically under the same conditions, confirming that the test reproduces the
    /// drift problem.
    #[test]
    fn mix_drift_sim_plus_300ppm_stays_bounded() {
        let corrected = run_drift_sim(300.0, 60, false);
        let uncorrected = run_drift_sim(300.0, 60, true);

        // Self-check: without correction, system levels grow monotonically each second
        // (+300 ppm ≈ +28.8 samples/second).
        for (i, w) in uncorrected.system_backlog.windows(2).enumerate() {
            assert!(
                w[1] > w[0],
                "levels should grow monotonically without correction: second {i}, {} → {}",
                w[0],
                w[1]
            );
        }

        // With correction, levels do not reach the safety limit (FIFO_MAX_SAMPLES = 500 ms).
        let max_corrected = corrected.system_backlog.iter().copied().max().unwrap();
        assert!(
            max_corrected < FIFO_MAX_SAMPLES,
            "with correction, should stay below the safety limit: max {max_corrected}"
        );

        // With correction, levels stay lower than without it, confirming that correction works.
        let last_c = *corrected.system_backlog.last().unwrap();
        let last_u = *uncorrected.system_backlog.last().unwrap();
        assert!(
            last_c < last_u,
            "corrected level {last_c} should be < uncorrected level {last_u}"
        );

        // With correction, growth slows (the increase over the first 9 seconds is > the
        // increase over the last 9 seconds).
        let sb = &corrected.system_backlog;
        let early = sb[9] - sb[0];
        let late = sb[59] - sb[50];
        assert!(
            late < early,
            "growth should slow as the ratio catches up: early +{early} vs late +{late}"
        );

        // The ratio moves toward faster system consumption and stays within the clamp.
        assert!(
            corrected.final_ratio > 1.0 + 2e-5,
            "r should move above 1.0: {}",
            corrected.final_ratio
        );
        assert!(
            corrected.final_ratio <= 1.0 + DRIFT_RATIO_LIMIT + 1e-12,
            "r should stay within the clamp: {}",
            corrected.final_ratio
        );
    }

    /// Simulate -300 ppm (system is slower) for 60 seconds: starvation zero-fill does
    /// not occur in steady state (run_drift_sim checks every output value on every
    /// tick). With correction, the mic-side level stays below the safety limit and
    /// lower than without correction.
    #[test]
    fn mix_drift_sim_minus_300ppm_no_steady_zero_fill() {
        let corrected = run_drift_sim(-300.0, 60, false);
        let uncorrected = run_drift_sim(-300.0, 60, true);

        // Self-check: without correction, mic levels grow monotonically to keep up with
        // the slower system.
        for (i, w) in uncorrected.mic_backlog.windows(2).enumerate() {
            assert!(
                w[1] > w[0],
                "mic levels should grow monotonically without correction: second {i}, {} → {}",
                w[0],
                w[1]
            );
        }

        // With correction, mic levels stay below the safety limit and lower than without it.
        let max_corrected = corrected.mic_backlog.iter().copied().max().unwrap();
        assert!(
            max_corrected < FIFO_MAX_SAMPLES,
            "with correction, should stay below the safety limit: max {max_corrected}"
        );
        let last_c = *corrected.mic_backlog.last().unwrap();
        let last_u = *uncorrected.mic_backlog.last().unwrap();
        assert!(
            last_c < last_u,
            "corrected mic level {last_c} should be < uncorrected level {last_u}"
        );

        // System does not accumulate because consumption tracks supply (at most an
        // interpolation fraction plus a recent partial chunk).
        let max_sys = corrected.system_backlog.iter().copied().max().unwrap();
        assert!(
            max_sys < ONE_CHUNK_SAMPLES,
            "system should not accumulate data: max {max_sys}"
        );

        // The ratio moves toward slower system consumption and stays within the clamp.
        assert!(
            corrected.final_ratio < 1.0 - 2e-5,
            "r should move below 1.0: {}",
            corrected.final_ratio
        );
        assert!(
            corrected.final_ratio >= 1.0 - DRIFT_RATIO_LIMIT - 1e-12,
            "r should stay within the clamp: {}",
            corrected.final_ratio
        );
    }
}