Skip to main content

rill_sampler/
timeseries.rs

1use rill_core::interpolate::Interpolate;
2use rill_core::time::ClockTick;
3use rill_core::traits::{
4    NodeCategory, NodeId, NodeMetadata, NodeState, ParamValue, ParameterId, Port, SignalNode,
5    Source,
6};
7use rill_core::Transcendental;
8use rill_core::{ProcessError, ProcessResult};
9use std::marker::PhantomData;
10
11/// Interpolation strategy for reading between samples.
12#[derive(Debug, Clone, Copy, PartialEq)]
13pub enum InterpMode {
14    /// Nearest-neighbour (no interpolation). Works for any `T`.
15    Nearest,
16    /// Linear interpolation. Requires `T: Transcendental`.
17    Linear,
18    /// Cubic Hermite interpolation. Requires `T: Transcendental`.
19    Cubic,
20}
21
22/// One channel of an unevenly-sampled time series.
23///
24/// `timestamps` must be monotonically non-decreasing and aligned with `values`.
25#[derive(Debug, Clone)]
26pub struct TimeSeriesChannel<T> {
27    /// Channel display name (e.g. `"engine_speed"`).
28    pub name: String,
29    /// Timestamps in seconds from start (monotonic).
30    pub timestamps: Vec<f64>,
31    /// Sample values aligned with `timestamps`.
32    pub values: Vec<T>,
33}
34
35impl<T> TimeSeriesChannel<T> {
36    /// Create a new empty channel with the given name.
37    pub fn new(name: impl Into<String>) -> Self {
38        Self {
39            name: name.into(),
40            timestamps: Vec::new(),
41            values: Vec::new(),
42        }
43    }
44
45    /// Total duration covered by this channel (seconds).
46    pub fn duration(&self) -> f64 {
47        if self.timestamps.len() < 2 {
48            0.0
49        } else {
50            self.timestamps[self.timestamps.len() - 1] - self.timestamps[0]
51        }
52    }
53
54    /// Number of samples in this channel.
55    pub fn len(&self) -> usize {
56        self.timestamps.len()
57    }
58
59    /// Returns `true` if the channel has no samples.
60    pub fn is_empty(&self) -> bool {
61        self.timestamps.is_empty()
62    }
63
64    /// Push one sample (caller must ensure timestamps stay monotonic).
65    pub fn push(&mut self, t: f64, value: T) {
66        self.timestamps.push(t);
67        self.values.push(value);
68    }
69}
70
71/// Unevenly-sampled time series reader.
72///
73/// Reads from multiple independent channels at a virtual uniform rate,
74/// using the [`Interpolate`] trait for fractional-index interpolation.
75///
76/// # Type parameter
77///
78/// `T` must implement `Transcendental` for `Linear` / `Cubic` modes.
79/// `Nearest` mode only requires `Copy`.
80pub struct TimeSeriesReader<T> {
81    channels: Vec<TimeSeriesChannel<T>>,
82    interp: InterpMode,
83}
84
85impl<T: Transcendental + Copy> TimeSeriesReader<T> {
86    /// Create an empty reader with nearest-neighbour interpolation.
87    pub fn new() -> Self {
88        Self {
89            channels: Vec::new(),
90            interp: InterpMode::Nearest,
91        }
92    }
93
94    /// Set the interpolation mode (builder pattern).
95    pub fn with_interp(mut self, mode: InterpMode) -> Self {
96        self.interp = mode;
97        self
98    }
99
100    /// Set the interpolation mode.
101    pub fn set_interp(&mut self, mode: InterpMode) {
102        self.interp = mode;
103    }
104
105    /// Return the current interpolation mode.
106    pub fn interp_mode(&self) -> InterpMode {
107        self.interp
108    }
109
110    /// Add a channel to the reader.
111    pub fn add_channel(&mut self, channel: TimeSeriesChannel<T>) {
112        self.channels.push(channel);
113    }
114
115    /// Number of registered channels.
116    pub fn num_channels(&self) -> usize {
117        self.channels.len()
118    }
119
120    /// Immutable access to a channel by index.
121    pub fn channel(&self, index: usize) -> Option<&TimeSeriesChannel<T>> {
122        self.channels.get(index)
123    }
124
125    /// Mutable access to a channel by index.
126    pub fn channel_mut(&mut self, index: usize) -> Option<&mut TimeSeriesChannel<T>> {
127        self.channels.get_mut(index)
128    }
129
130    /// Total time span across all channels (union of ranges).
131    pub fn duration(&self) -> f64 {
132        self.channels
133            .iter()
134            .map(|c| c.duration())
135            .fold(0.0, f64::max)
136    }
137
138    /// Read value from a single channel at an arbitrary timestamp.
139    pub fn at_time(&self, channel: usize, t: f64) -> T {
140        let Some(ch) = self.channels.get(channel) else {
141            return T::ZERO;
142        };
143        if ch.len() < 2 {
144            return ch.values.first().copied().unwrap_or(T::ZERO);
145        }
146
147        // Binary search for the segment containing t
148        let idx = match ch
149            .timestamps
150            .binary_search_by(|&ts| ts.partial_cmp(&t).unwrap())
151        {
152            Ok(i) => {
153                // Exact match: return the value directly
154                return ch.values[i];
155            }
156            Err(i) => {
157                // i is where t would be inserted
158                if i == 0 {
159                    return ch.values[0]; // before start → clamp
160                }
161                if i >= ch.len() {
162                    return ch.values[ch.len() - 1]; // past end → clamp
163                }
164                i - 1 // segment index
165            }
166        };
167
168        let t0 = ch.timestamps[idx];
169        let t1 = ch.timestamps[idx + 1];
170        let span = t1 - t0;
171        if span <= 0.0 {
172            return ch.values[idx];
173        }
174
175        let frac = (t - t0) / span;
176        let index = idx as f64 + frac;
177
178        match self.interp {
179            InterpMode::Nearest => ch.values.interpolate_nearest(index),
180            InterpMode::Linear => ch.values.interpolate_linear(index),
181            InterpMode::Cubic => ch.values.interpolate_cubic(index),
182        }
183    }
184
185    /// Fill a planar output buffer.
186    ///
187    /// Layout: `[ch0_s0, ch0_s1, ..., ch0_s{BUF-1}, ch1_s0, ...]`
188    /// i.e. `output[ch * buf_size + i]`.
189    pub fn read_block(&self, time: f64, sample_rate: f64, output: &mut [T]) {
190        let nch = self.channels.len();
191        if nch == 0 {
192            for s in output.iter_mut() {
193                *s = T::ZERO;
194            }
195            return;
196        }
197        let buf_size = output.len() / nch;
198        let dt = 1.0 / sample_rate;
199        for (ch, s) in output.chunks_mut(buf_size).enumerate() {
200            for (i, v) in s.iter_mut().enumerate() {
201                *v = self.at_time(ch, time + i as f64 * dt);
202            }
203        }
204    }
205
206    /// All registered channels as a slice.
207    pub fn channels(&self) -> &[TimeSeriesChannel<T>] {
208        &self.channels
209    }
210}
211
212impl<T: Transcendental + Copy> Default for TimeSeriesReader<T> {
213    fn default() -> Self {
214        Self::new()
215    }
216}
217
218// ---------------------------------------------------------------------------
219// Graph node
220// ---------------------------------------------------------------------------
221
222/// Source node wrapping [`TimeSeriesReader`].
223///
224/// Produces one output port per channel, each filled at the configured
225/// virtual `sample_rate`. All automatable via patchbay.
226pub struct TimeSeriesNode<T: Transcendental, const BUF_SIZE: usize> {
227    reader: TimeSeriesReader<T>,
228    sample_rate: f64,
229    playing: bool,
230    time: f64,
231    speed: f64,
232    outputs: Vec<Port<T, BUF_SIZE>>,
233    state: Option<NodeState<T, BUF_SIZE>>,
234    _phantom: PhantomData<[T; BUF_SIZE]>,
235}
236
237impl<T: Transcendental + Copy, const BUF_SIZE: usize> TimeSeriesNode<T, BUF_SIZE> {
238    /// Create a new node with linear interpolation at 100 Hz virtual rate.
239    pub fn new() -> Self {
240        Self {
241            reader: TimeSeriesReader::new().with_interp(InterpMode::Linear),
242            sample_rate: 100.0,
243            playing: true,
244            time: 0.0,
245            speed: 1.0,
246            outputs: Vec::new(),
247            state: None,
248            _phantom: PhantomData,
249        }
250    }
251
252    /// Immutable reference to the inner reader.
253    pub fn reader(&self) -> &TimeSeriesReader<T> {
254        &self.reader
255    }
256
257    /// Mutable reference to the inner reader.
258    pub fn reader_mut(&mut self) -> &mut TimeSeriesReader<T> {
259        &mut self.reader
260    }
261
262    /// Replace all channels and rebuild output ports.
263    pub fn set_channels(&mut self, channels: Vec<TimeSeriesChannel<T>>) {
264        self.outputs.clear();
265        for (i, ch) in channels.iter().enumerate() {
266            self.outputs
267                .push(Port::output(NodeId(0), i as u16, &ch.name));
268        }
269        self.reader = TimeSeriesReader {
270            channels,
271            interp: self.reader.interp,
272        };
273        self.time = 0.0;
274    }
275
276    fn param_to_t(value: ParamValue) -> Option<T> {
277        match value {
278            ParamValue::Float(f) => Some(T::from_f32(f)),
279            ParamValue::Int(i) => Some(T::from_f32(i as f32)),
280            _ => None,
281        }
282    }
283}
284
285impl<T: Transcendental + Copy, const BUF_SIZE: usize> Default for TimeSeriesNode<T, BUF_SIZE> {
286    fn default() -> Self {
287        Self::new()
288    }
289}
290
291impl<T: Transcendental + Copy, const BUF_SIZE: usize> SignalNode<T, BUF_SIZE>
292    for TimeSeriesNode<T, BUF_SIZE>
293{
294    fn metadata(&self) -> NodeMetadata {
295        NodeMetadata {
296            name: "TimeSeries".to_string(),
297            type_name: None,
298            category: NodeCategory::Source,
299            description: "Unevenly-sampled time series reader with multiple output channels".into(),
300            author: "Rill".to_string(),
301            version: env!("CARGO_PKG_VERSION").to_string(),
302            signal_inputs: 0,
303            signal_outputs: self.outputs.len(),
304            control_inputs: 0,
305            control_outputs: 0,
306            clock_inputs: 0,
307            clock_outputs: 0,
308            feedback_ports: 0,
309            parameters: vec![],
310        }
311    }
312
313    fn init(&mut self, sample_rate: f32) {
314        self.state = Some(NodeState::new(sample_rate));
315    }
316
317    fn reset(&mut self) {
318        self.time = 0.0;
319        self.playing = true;
320        if let Some(state) = &mut self.state {
321            state.reset();
322        }
323    }
324
325    fn get_parameter(&self, id: &ParameterId) -> Option<ParamValue> {
326        match id.as_str() {
327            "sample_rate" => Some(ParamValue::Float(self.sample_rate as f32)),
328            "interpolation" => {
329                let s = match self.reader.interp_mode() {
330                    InterpMode::Nearest => "nearest",
331                    InterpMode::Linear => "linear",
332                    InterpMode::Cubic => "cubic",
333                };
334                Some(ParamValue::Choice(s.into()))
335            }
336            "play" => Some(ParamValue::Bool(self.playing)),
337            "position" => {
338                let dur = self.reader.duration();
339                if dur > 0.0 {
340                    Some(ParamValue::Float((self.time / dur) as f32))
341                } else {
342                    Some(ParamValue::Float(0.0))
343                }
344            }
345            "speed" => Some(ParamValue::Float(self.speed as f32)),
346            _ => None,
347        }
348    }
349
350    fn set_parameter(&mut self, id: &ParameterId, value: ParamValue) -> ProcessResult<()> {
351        match id.as_str() {
352            "sample_rate" => {
353                if let Some(r) = Self::param_to_t(value) {
354                    self.sample_rate = r.to_f64().clamp(0.1, 1_000_000.0);
355                    Ok(())
356                } else {
357                    Err(ProcessError::Parameter("Expected float".into()))
358                }
359            }
360            "interpolation" => {
361                if let ParamValue::Choice(s) = &value {
362                    self.reader.set_interp(match s.as_str() {
363                        "linear" => InterpMode::Linear,
364                        "cubic" => InterpMode::Cubic,
365                        _ => InterpMode::Nearest,
366                    });
367                    Ok(())
368                } else {
369                    Err(ProcessError::Parameter("Expected choice".into()))
370                }
371            }
372            "play" => {
373                if let ParamValue::Bool(b) = value {
374                    self.playing = b;
375                    Ok(())
376                } else {
377                    Err(ProcessError::Parameter("Expected bool".into()))
378                }
379            }
380            "speed" => {
381                if let Some(s) = Self::param_to_t(value) {
382                    self.speed = s.to_f64().clamp(0.0, 100.0);
383                    Ok(())
384                } else {
385                    Err(ProcessError::Parameter("Expected float".into()))
386                }
387            }
388            _ => Err(ProcessError::Parameter(format!(
389                "Unknown parameter: {}",
390                id
391            ))),
392        }
393    }
394
395    fn id(&self) -> NodeId {
396        NodeId(0)
397    }
398
399    fn set_id(&mut self, _id: NodeId) {}
400
401    fn input_port(&self, _index: usize) -> Option<&Port<T, BUF_SIZE>> {
402        None
403    }
404
405    fn input_port_mut(&mut self, _index: usize) -> Option<&mut Port<T, BUF_SIZE>> {
406        None
407    }
408
409    fn output_port(&self, index: usize) -> Option<&Port<T, BUF_SIZE>> {
410        self.outputs.get(index)
411    }
412
413    fn output_port_mut(&mut self, index: usize) -> Option<&mut Port<T, BUF_SIZE>> {
414        self.outputs.get_mut(index)
415    }
416
417    fn control_port(&self, _index: usize) -> Option<&Port<T, BUF_SIZE>> {
418        None
419    }
420
421    fn control_port_mut(&mut self, _index: usize) -> Option<&mut Port<T, BUF_SIZE>> {
422        None
423    }
424
425    fn state(&self) -> &NodeState<T, BUF_SIZE> {
426        self.state.as_ref().unwrap()
427    }
428
429    fn state_mut(&mut self) -> &mut NodeState<T, BUF_SIZE> {
430        self.state.as_mut().unwrap()
431    }
432
433    fn num_signal_inputs(&self) -> usize {
434        0
435    }
436
437    fn num_signal_outputs(&self) -> usize {
438        self.outputs.len()
439    }
440}
441
442impl<T: Transcendental + Copy, const BUF_SIZE: usize> Source<T, BUF_SIZE>
443    for TimeSeriesNode<T, BUF_SIZE>
444{
445    fn generate(
446        &mut self,
447        _clock: &ClockTick,
448        _control_inputs: &[T],
449        _clock_inputs: &[ClockTick],
450    ) -> ProcessResult<()> {
451        if !self.playing || self.reader.num_channels() == 0 {
452            for port in self.outputs.iter_mut() {
453                port.buffer.as_mut_array().fill(T::ZERO);
454            }
455            return Ok(());
456        }
457
458        let nch = self.reader.num_channels();
459        let dur = self.reader.duration();
460        let dt = 1.0 / self.sample_rate;
461
462        // Planar per-channel writes
463        for (ch_idx, port) in self.outputs.iter_mut().enumerate().take(nch) {
464            let buf = port.buffer.as_mut_array();
465            for (i, v) in buf.iter_mut().enumerate() {
466                let t = self.time + i as f64 * dt;
467                *v = self.reader.at_time(ch_idx, t);
468            }
469        }
470
471        self.time += BUF_SIZE as f64 * dt * self.speed;
472
473        // Clamp and optionally pause at end
474        if self.time >= dur && dur > 0.0 {
475            if self.speed > 0.0 {
476                self.time = dur; // hold last values
477                self.playing = false;
478            }
479        } else if self.time < 0.0 {
480            self.time = 0.0;
481            self.playing = false;
482        }
483
484        Ok(())
485    }
486}
487
488// ---------------------------------------------------------------------------
489// CSV loader
490// ---------------------------------------------------------------------------
491
492/// Load a time-series reader from a CSV string.
493///
494/// Expected format (header optional):
495/// ```csv
496/// t,channel,value
497/// 0.001,engine_speed,1500
498/// 0.001,oil_temp,85
499/// ```
500///
501/// Lines that cannot be parsed are silently skipped.
502pub fn from_csv<T: Transcendental + Copy>(input: &str) -> TimeSeriesReader<T> {
503    use std::collections::BTreeMap;
504
505    let mut raw: BTreeMap<String, Vec<(f64, T)>> = BTreeMap::new();
506
507    for line in input.lines() {
508        let line = line.trim();
509        if line.is_empty() || line.starts_with("t,") || line.starts_with("timestamp,") {
510            continue;
511        }
512        let mut parts = line.splitn(3, ',');
513        let t: f64 = match parts.next().and_then(|s| s.trim().parse().ok()) {
514            Some(v) => v,
515            None => continue,
516        };
517        let name = match parts.next() {
518            Some(s) => s.trim().to_string(),
519            None => continue,
520        };
521        let value: f64 = match parts.next().and_then(|s| s.trim().parse().ok()) {
522            Some(v) => v,
523            None => continue,
524        };
525
526        raw.entry(name).or_default().push((t, T::from_f64(value)));
527    }
528
529    let mut reader = TimeSeriesReader::new();
530    for (name, mut samples) in raw {
531        // Sort by timestamp (BTreeMap iteration is key-ordered, but values
532        // within a channel may arrive out of order in CSV)
533        samples.sort_by(|a, b| a.0.partial_cmp(&b.0).unwrap());
534        let mut ch = TimeSeriesChannel::new(&name);
535        for (t, v) in samples {
536            ch.push(t, v);
537        }
538        reader.add_channel(ch);
539    }
540
541    reader
542}
543
544#[cfg(test)]
545mod tests {
546    use super::*;
547
548    fn ae(a: f64, b: f64) -> bool {
549        (a - b).abs() < 1e-6
550    }
551
552    #[test]
553    fn test_at_time_exact() {
554        let mut ch = TimeSeriesChannel::new("test");
555        ch.push(0.0, 10.0);
556        ch.push(1.0, 20.0);
557        ch.push(2.0, 30.0);
558        let mut reader = TimeSeriesReader::new();
559        reader.add_channel(ch);
560
561        assert!(ae(reader.at_time(0, 0.0), 10.0));
562        assert!(ae(reader.at_time(0, 1.0), 20.0));
563        assert!(ae(reader.at_time(0, 2.0), 30.0));
564    }
565
566    #[test]
567    fn test_at_time_interpolated() {
568        let mut ch = TimeSeriesChannel::new("test");
569        ch.push(0.0, 0.0);
570        ch.push(1.0, 2.0);
571        let mut reader = TimeSeriesReader::new().with_interp(InterpMode::Linear);
572        reader.add_channel(ch);
573
574        assert!(ae(reader.at_time(0, 0.5), 1.0));
575        assert!(ae(reader.at_time(0, 0.25), 0.5));
576    }
577
578    #[test]
579    fn test_at_time_clamp() {
580        let mut ch = TimeSeriesChannel::new("test");
581        ch.push(1.0, 100.0);
582        ch.push(2.0, 200.0);
583        let mut reader = TimeSeriesReader::new();
584        reader.add_channel(ch);
585
586        assert!(ae(reader.at_time(0, 0.0), 100.0));
587        assert!(ae(reader.at_time(0, 5.0), 200.0));
588    }
589
590    #[test]
591    fn test_nearest_mode() {
592        let mut ch = TimeSeriesChannel::new("test");
593        ch.push(0.0, 10.0);
594        ch.push(1.0, 20.0);
595        let mut reader = TimeSeriesReader::new().with_interp(InterpMode::Nearest);
596        reader.add_channel(ch);
597
598        assert!(ae(reader.at_time(0, 0.49), 10.0));
599        assert!(ae(reader.at_time(0, 0.5), 20.0));
600    }
601
602    #[test]
603    fn test_empty_channel() {
604        let ch = TimeSeriesChannel::new("empty");
605        let mut reader = TimeSeriesReader::new();
606        reader.add_channel(ch);
607        assert!(ae(reader.at_time(0, 0.5), 0.0));
608    }
609
610    #[test]
611    fn test_read_block() {
612        let mut ch = TimeSeriesChannel::new("ch");
613        ch.push(0.0, 1.0);
614        ch.push(1.0, 3.0);
615        let mut reader = TimeSeriesReader::new().with_interp(InterpMode::Linear);
616        reader.add_channel(ch);
617
618        let mut out = [0.0_f64; 4];
619        reader.read_block(0.0, 2.0, &mut out);
620        assert!(ae(out[0], 1.0));
621        assert!(ae(out[1], 2.0));
622        assert!(ae(out[2], 3.0));
623        assert!(ae(out[3], 3.0));
624    }
625
626    #[test]
627    fn test_read_multichannel() {
628        let mut ch1 = TimeSeriesChannel::new("a");
629        ch1.push(0.0, 10.0);
630        ch1.push(1.0, 20.0);
631        let mut ch2 = TimeSeriesChannel::new("b");
632        ch2.push(0.0, 100.0);
633        ch2.push(1.0, 200.0);
634        let mut reader = TimeSeriesReader::new().with_interp(InterpMode::Linear);
635        reader.add_channel(ch1);
636        reader.add_channel(ch2);
637
638        let mut out = [0.0_f64; 4];
639        reader.read_block(0.5, 2.0, &mut out);
640        assert!(ae(out[0], 15.0));
641        assert!(ae(out[1], 20.0));
642        assert!(ae(out[2], 150.0));
643        assert!(ae(out[3], 200.0));
644    }
645
646    #[test]
647    fn test_csv_loading() {
648        let csv = "\
649t,channel,value
6500.0,speed,100
6510.5,speed,200
6520.0,temp,25
6530.5,temp,30
654";
655        let reader: TimeSeriesReader<f64> = from_csv(csv);
656        assert_eq!(reader.num_channels(), 2);
657        let speed = reader.channel(0).unwrap();
658        assert_eq!(speed.name, "speed");
659        assert!(ae(speed.values[0], 100.0));
660        assert!(ae(speed.values[1], 200.0));
661    }
662
663    #[test]
664    fn test_timeseries_node_basic() {
665        let mut ch = TimeSeriesChannel::new("test");
666        ch.push(0.0, 1.0);
667        ch.push(1.0, 2.0);
668        let mut node = TimeSeriesNode::<f64, 4>::new();
669        node.set_channels(vec![ch]);
670        node.init(44100.0);
671        node.sample_rate = 2.0;
672
673        let clock = ClockTick::new(0, 4, 44100.0);
674        node.generate(&clock, &[], &[]).unwrap();
675
676        let port = node.output_port(0).unwrap();
677        let buf = port.buffer.as_array();
678        assert!(ae(buf[0], 1.0));
679        assert!(ae(buf[1], 1.5));
680        assert!(ae(buf[2], 2.0));
681        assert!(ae(buf[3], 2.0));
682    }
683}