Skip to main content

rustradio/
delay.rs

1//! Delay stream. Good for syncing up streams.
2use log::debug;
3
4use crate::block::{Block, BlockRet};
5use crate::stream::{ReadStream, WriteStream};
6use crate::{Result, Sample};
7
8/// Delay stream. Good for syncing up streams.
9#[derive(rustradio_macros::Block)]
10#[rustradio(crate, noeof)]
11pub struct Delay<T: Sample> {
12    delay: usize,
13    current_delay: usize,
14
15    // Skip is the number of samples we're needlessly ahead. This can happen
16    // when the delay changes mid stream.
17    skip: usize,
18
19    #[rustradio(in)]
20    src: ReadStream<T>,
21    #[rustradio(out)]
22    dst: WriteStream<T>,
23}
24
25impl<T: Sample> Delay<T> {
26    /// Create new Delay block.
27    #[must_use]
28    pub fn new(src: ReadStream<T>, delay: usize) -> (Self, ReadStream<T>) {
29        let (dst, dr) = crate::stream::new_stream();
30        (
31            Self {
32                src,
33                dst,
34                delay,
35                current_delay: delay,
36                skip: 0,
37            },
38            dr,
39        )
40    }
41
42    /// Change the delay.
43    pub fn set_delay(&mut self, delay: usize) {
44        if delay > self.delay {
45            self.current_delay += delay - self.delay;
46        } else {
47            let reduce = self.delay - delay;
48            let cdskip = std::cmp::min(self.current_delay, reduce);
49            self.current_delay -= cdskip;
50            self.skip += reduce - cdskip;
51        }
52        self.delay = delay;
53    }
54}
55
56impl<T: Sample> crate::block::BlockEOF for Delay<T> {
57    fn eof(&mut self) -> bool {
58        self.current_delay == 0 && self.src.eof()
59    }
60}
61
62impl<T: Sample> Block for Delay<T> {
63    fn work(&mut self) -> Result<BlockRet<'_>> {
64        loop {
65            // Check if we need to catch up.
66            let (input, tags) = self.src.read_buf()?;
67            if self.skip > 0 {
68                let n = std::cmp::min(input.len(), self.skip);
69                if n == 0 {
70                    return Ok(BlockRet::WaitForStream(&self.src, 1));
71                }
72                input.consume(n);
73                debug!("Delay: skipped {n}");
74                self.skip -= n;
75                continue;
76            }
77
78            // Everything except catch-up requires output space.
79            let mut o = self.dst.write_buf()?;
80            if o.is_empty() {
81                return Ok(BlockRet::WaitForStream(&self.dst, 1));
82            }
83
84            // Check if we're still delaying, thus filling with default.
85            if self.current_delay > 0 {
86                let n = std::cmp::min(self.current_delay, o.len());
87                o.slice()[..n].fill(T::default());
88                o.produce(n, &[]);
89                self.current_delay -= n;
90                continue;
91            }
92
93            // Neither skipping nor delaying. just plain copy.
94            if input.is_empty() {
95                return Ok(BlockRet::WaitForStream(&self.src, 1));
96            }
97
98            let n = std::cmp::min(input.len(), o.len());
99            assert_ne!(
100                n, 0,
101                "can't happen: we already checked both input and output"
102            );
103            o.fill_from_slice(&input.slice()[..n]);
104            let tags = tags
105                .into_iter()
106                .filter(|tag| tag.pos() < n)
107                .collect::<Vec<_>>();
108            o.produce(n, &tags);
109            input.consume(n);
110        }
111    }
112}
113
114#[cfg(test)]
115mod tests {
116    use super::*;
117
118    // TODO: test tag propagation.
119
120    #[test]
121    fn delay_zero() -> Result<()> {
122        let s = ReadStream::from_slice(&[1.0f32, 2.0, 3.0]);
123        let (mut delay, o) = Delay::new(s, 0);
124
125        delay.work()?;
126        let (res, _) = o.read_buf()?;
127        assert_eq!(res.slice(), vec![1.0f32, 2.0, 3.0]);
128        Ok(())
129    }
130
131    #[test]
132    fn delay_one() -> Result<()> {
133        let s = ReadStream::from_slice(&[1.0f32, 2.0, 3.0]);
134        let (mut delay, o) = Delay::new(s, 1);
135
136        delay.work()?;
137        let (res, _) = o.read_buf()?;
138        assert_eq!(res.slice(), vec![0.0f32, 1.0, 2.0, 3.0]);
139        Ok(())
140    }
141
142    #[test]
143    fn delay_increase_before_work_extends_remaining_delay() -> Result<()> {
144        let s = ReadStream::from_slice(&[1u32, 2]);
145        let (mut delay, o) = Delay::new(s, 1);
146
147        delay.set_delay(2);
148        delay.work()?;
149        let (res, _) = o.read_buf()?;
150        assert_eq!(res.slice(), &[0, 0, 1, 2]);
151        Ok(())
152    }
153
154    #[test]
155    fn delay_decrease_before_work_reduces_remaining_delay() -> Result<()> {
156        let s = ReadStream::from_slice(&[1u32, 2]);
157        let (mut delay, o) = Delay::new(s, 3);
158
159        delay.set_delay(1);
160        delay.work()?;
161        let (res, _) = o.read_buf()?;
162        assert_eq!(res.slice(), &[0, 1, 2]);
163        Ok(())
164    }
165
166    #[test]
167    fn delay_reduced_twice_accumulates_pending_skip() -> Result<()> {
168        let cap = crate::stream::DEFAULT_STREAM_SIZE / std::mem::size_of::<u32>();
169        let input = (0..cap as u32).collect::<Vec<_>>();
170        let s = ReadStream::from_slice(&input);
171        let (mut delay, o) = Delay::new(s, cap + 10);
172
173        delay.work()?;
174        {
175            let (res, _) = o.read_buf()?;
176            let len = res.len();
177            assert_eq!(len, cap);
178            assert!(res.iter().all(|v| *v == 0));
179            res.consume(len);
180        }
181
182        delay.set_delay(cap - 1);
183        delay.set_delay(cap - 2);
184        delay.work()?;
185        let (res, _) = o.read_buf()?;
186        assert_eq!(res.slice(), &input[2..]);
187        Ok(())
188    }
189
190    #[test]
191    fn eof_waits_for_pending_delay() -> Result<()> {
192        let cap = crate::stream::DEFAULT_STREAM_SIZE / std::mem::size_of::<u32>();
193        let s = ReadStream::<u32>::from_slice(&[]);
194        let (mut delay, o) = Delay::new(s, cap + 1);
195
196        assert!(!crate::block::BlockEOF::eof(&mut delay));
197        assert!(matches![delay.work()?, BlockRet::WaitForStream(_, 1)]);
198        assert!(!crate::block::BlockEOF::eof(&mut delay));
199
200        let (res, _) = o.read_buf()?;
201        assert_eq!(res.len(), cap);
202        res.consume(cap);
203
204        assert!(matches![delay.work()?, BlockRet::WaitForStream(_, 1)]);
205        assert!(crate::block::BlockEOF::eof(&mut delay));
206        Ok(())
207    }
208
209    #[test]
210    fn delay_change() -> Result<()> {
211        let s = ReadStream::from_slice(&[1u32, 2]);
212        let (mut delay, o) = Delay::new(s, 1);
213
214        delay.work()?;
215        {
216            let (res, _) = o.read_buf()?;
217            assert_eq!(res.slice(), vec![0, 1, 2]);
218        }
219
220        // TODO: fix
221        /*
222        // 3,4 => 0,3,4
223        {
224            let mut b = s.write_buf()?;
225            b.fill_from_slice(&[3, 4]);
226            b.produce(2, &[]);
227        }
228        delay.set_delay(2);
229        delay.work()?;
230        {
231            let (res, _) = o.read_buf()?;
232            assert_eq!(res.slice(), vec![0, 1, 2, 0, 3, 4]);
233        }
234
235        // 5,6 => 0,3,4
236        {
237            let mut b = s.write_buf()?;
238            b.fill_from_slice(&[5, 6]);
239            b.produce(2, &[]);
240        }
241        delay.set_delay(0);
242        delay.work()?;
243        {
244            let (res, _) = o.read_buf()?;
245            assert_eq!(res.slice(), vec![0, 1, 2, 0, 3, 4]);
246        }
247
248        // 7 => 7
249        {
250            let mut b = s.write_buf()?;
251            b.slice()[0] = 7;
252            b.produce(1, &[]);
253        }
254        delay.set_delay(0);
255        delay.work()?;
256        {
257            let (res, _) = o.read_buf()?;
258            assert_eq!(res.slice(), vec![0, 1, 2, 0, 3, 4, 7]);
259        }
260        */
261        Ok(())
262    }
263}