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
// SPDX-License-Identifier: Apache-2.0
// Copyright (c) 2026 Fábio Henrique de Lima Silva (fhl.bsb@gmail.com) All rights reserved.
use super::*;
use std::sync::atomic::Ordering;
use std::time::{Duration, Instant};
#[test]
fn test_dsp_bridge_concurrent_access() {
// 1. DspBridge setup:
// We use Box::leak to obtain a 'static reference, simulating the runtime
// behavior where the object lives for the entire duration of the PipeWire host.
// This allows safely converting the reference into raw pointers (*const/*mut).
let bridge: &'static DspBridge = Box::leak(Box::new(DspBridge {
buffers: [
BridgeBuffer {
buf_l: [0.0; MAX_BRIDGE_BUF],
buf_r: [0.0; MAX_BRIDGE_BUF],
n_samples: 0,
},
BridgeBuffer {
buf_l: [0.0; MAX_BRIDGE_BUF],
buf_r: [0.0; MAX_BRIDGE_BUF],
n_samples: 0,
},
],
active_read_idx: std::sync::atomic::AtomicUsize::new(0),
generation: std::sync::atomic::AtomicU64::new(0),
consumed_gen: std::sync::atomic::AtomicU64::new(0),
dropped_frames: std::sync::atomic::AtomicU32::new(0),
}));
// Raw pointers for the threads (writer/reader)
let bridge_ptr_writer = bridge as *const DspBridge as *mut DspBridge as usize;
let bridge_ptr_reader = bridge as *const DspBridge;
// 2. Writer Thread (Simulates RT Capture Callback):
// This thread fills the "back buffer" (the buffer not being read)
// and then toggles the atomic index to make it the new "front buffer".
let writer_handle = std::thread::spawn(move || {
let mut counter = 0.0f32;
let bridge_ptr_writer = bridge_ptr_writer as *mut DspBridge;
for _ in 0..1000 {
let bridge_ref = unsafe { &mut *bridge_ptr_writer };
// Prevents the writer from overwriting the buffer being read
while bridge_ref.generation.load(Ordering::Acquire)
> bridge_ref.consumed_gen.load(Ordering::Acquire)
{
std::thread::yield_now();
}
// Locates the inactive buffer (back-buffer) for writing.
let back_idx = 1 - bridge_ref.active_read_idx.load(Ordering::Relaxed);
let back_buf = &mut bridge_ref.buffers[back_idx];
// Fills with sequential data to verify integrity in the reader.
for i in 0..64 {
back_buf.buf_l[i] = counter;
back_buf.buf_r[i] = counter;
counter += 1.0;
}
back_buf.n_samples = 64;
bridge_ref
.active_read_idx
.store(back_idx, Ordering::Release);
bridge_ref.generation.fetch_add(1, Ordering::Release);
// Small delay to simulate DSP processing time and allow interleaving.
std::thread::sleep(Duration::from_micros(10));
}
});
// 3. Reader Thread (Simulates RT Playback Callback):
// This thread monitors the 'generation' counter. When it changes, it
// consumes the new "front buffer".
let start = Instant::now();
let mut last_gen = 0;
let mut reads = 0;
let mut last_val_read = -1.0;
while reads < 1000 && start.elapsed() < Duration::from_millis(500) {
let bridge_ref = unsafe { &*bridge_ptr_reader };
let current_gen = bridge_ref.generation.load(Ordering::Acquire);
if current_gen != last_gen {
let read_idx = bridge_ref.active_read_idx.load(Ordering::Acquire);
let front_buf = &bridge_ref.buffers[read_idx];
assert_eq!(front_buf.n_samples, 64);
// Integrity Check: data in a buffer must be contiguous.
let first_val = front_buf.buf_l[0];
for i in 0..64 {
assert_eq!(
front_buf.buf_l[i],
first_val + i as f32,
"Buffer mixing detected in channel L"
);
assert_eq!(
front_buf.buf_r[i],
first_val + i as f32,
"Buffer mixing detected in channel R"
);
}
// Monotonicity Check:
// Even if we skip frames (which may occur in tests under load),
// we must never read data older than previously read.
assert!(
first_val > last_val_read,
"Read older data than previously seen! (Stale read)"
);
last_val_read = front_buf.buf_l[63];
last_gen = current_gen;
bridge_ref
.consumed_gen
.store(current_gen, Ordering::Release);
reads += 1;
}
if last_gen == 1000 {
break;
}
}
bridge.consumed_gen.store(u64::MAX, Ordering::Release);
writer_handle.join().unwrap();
// 4. Performance Check:
// The test should complete quickly. If it takes too long, it indicates deadlocks
// or severe starvation (even though the design is lock-free).
assert!(
start.elapsed() < Duration::from_millis(500),
"Test took too long to execute"
);
}