Skip to main content

rill_sampler/
timeseries.rs

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