1use log::debug;
3
4use crate::block::{Block, BlockRet};
5use crate::stream::{ReadStream, WriteStream};
6use crate::{Result, Sample};
7
8#[derive(rustradio_macros::Block)]
10#[rustradio(crate, noeof)]
11pub struct Delay<T: Sample> {
12 delay: usize,
13 current_delay: usize,
14
15 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 #[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 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 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 let mut o = self.dst.write_buf()?;
80 if o.is_empty() {
81 return Ok(BlockRet::WaitForStream(&self.dst, 1));
82 }
83
84 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 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 #[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 Ok(())
262 }
263}