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#[derive(Debug, Clone, Copy, PartialEq)]
13pub enum InterpMode {
14 Nearest,
16 Linear,
18 Cubic,
20}
21
22#[derive(Debug, Clone)]
26pub struct TimeSeriesChannel<T> {
27 pub name: String,
29 pub timestamps: Vec<f64>,
31 pub values: Vec<T>,
33}
34
35impl<T> TimeSeriesChannel<T> {
36 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 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 pub fn len(&self) -> usize {
56 self.timestamps.len()
57 }
58
59 pub fn is_empty(&self) -> bool {
61 self.timestamps.is_empty()
62 }
63
64 pub fn push(&mut self, t: f64, value: T) {
66 self.timestamps.push(t);
67 self.values.push(value);
68 }
69}
70
71pub struct TimeSeriesReader<T> {
81 channels: Vec<TimeSeriesChannel<T>>,
82 interp: InterpMode,
83}
84
85impl<T: Transcendental + Copy> TimeSeriesReader<T> {
86 pub fn new() -> Self {
88 Self {
89 channels: Vec::new(),
90 interp: InterpMode::Nearest,
91 }
92 }
93
94 pub fn with_interp(mut self, mode: InterpMode) -> Self {
96 self.interp = mode;
97 self
98 }
99
100 pub fn set_interp(&mut self, mode: InterpMode) {
102 self.interp = mode;
103 }
104
105 pub fn interp_mode(&self) -> InterpMode {
107 self.interp
108 }
109
110 pub fn add_channel(&mut self, channel: TimeSeriesChannel<T>) {
112 self.channels.push(channel);
113 }
114
115 pub fn num_channels(&self) -> usize {
117 self.channels.len()
118 }
119
120 pub fn channel(&self, index: usize) -> Option<&TimeSeriesChannel<T>> {
122 self.channels.get(index)
123 }
124
125 pub fn channel_mut(&mut self, index: usize) -> Option<&mut TimeSeriesChannel<T>> {
127 self.channels.get_mut(index)
128 }
129
130 pub fn duration(&self) -> f64 {
132 self.channels
133 .iter()
134 .map(|c| c.duration())
135 .fold(0.0, f64::max)
136 }
137
138 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 let idx = match ch
149 .timestamps
150 .binary_search_by(|&ts| ts.partial_cmp(&t).unwrap())
151 {
152 Ok(i) => {
153 return ch.values[i];
155 }
156 Err(i) => {
157 if i == 0 {
159 return ch.values[0]; }
161 if i >= ch.len() {
162 return ch.values[ch.len() - 1]; }
164 i - 1 }
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 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 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
218pub 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 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 pub fn reader(&self) -> &TimeSeriesReader<T> {
254 &self.reader
255 }
256
257 pub fn reader_mut(&mut self) -> &mut TimeSeriesReader<T> {
259 &mut self.reader
260 }
261
262 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 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 if self.time >= dur && dur > 0.0 {
475 if self.speed > 0.0 {
476 self.time = dur; 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
488pub 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 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}