1use rill_core::interpolate::Interpolate;
2use rill_core::time::{ClockTick, RenderContext};
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#[derive(Debug, Clone, Copy, PartialEq)]
12pub enum InterpMode {
13 Nearest,
15 Linear,
17 Cubic,
19}
20
21#[derive(Debug, Clone)]
25pub struct TimeSeriesChannel<T> {
26 pub name: String,
28 pub timestamps: Vec<f64>,
30 pub values: Vec<T>,
32}
33
34impl<T> TimeSeriesChannel<T> {
35 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 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 pub fn len(&self) -> usize {
55 self.timestamps.len()
56 }
57
58 pub fn is_empty(&self) -> bool {
60 self.timestamps.is_empty()
61 }
62
63 pub fn push(&mut self, t: f64, value: T) {
65 self.timestamps.push(t);
66 self.values.push(value);
67 }
68}
69
70pub struct TimeSeriesReader<T> {
80 channels: Vec<TimeSeriesChannel<T>>,
81 interp: InterpMode,
82}
83
84impl<T: Transcendental + Copy> TimeSeriesReader<T> {
85 pub fn new() -> Self {
87 Self {
88 channels: Vec::new(),
89 interp: InterpMode::Nearest,
90 }
91 }
92
93 pub fn with_interp(mut self, mode: InterpMode) -> Self {
95 self.interp = mode;
96 self
97 }
98
99 pub fn set_interp(&mut self, mode: InterpMode) {
101 self.interp = mode;
102 }
103
104 pub fn interp_mode(&self) -> InterpMode {
106 self.interp
107 }
108
109 pub fn add_channel(&mut self, channel: TimeSeriesChannel<T>) {
111 self.channels.push(channel);
112 }
113
114 pub fn num_channels(&self) -> usize {
116 self.channels.len()
117 }
118
119 pub fn channel(&self, index: usize) -> Option<&TimeSeriesChannel<T>> {
121 self.channels.get(index)
122 }
123
124 pub fn channel_mut(&mut self, index: usize) -> Option<&mut TimeSeriesChannel<T>> {
126 self.channels.get_mut(index)
127 }
128
129 pub fn duration(&self) -> f64 {
131 self.channels
132 .iter()
133 .map(|c| c.duration())
134 .fold(0.0, f64::max)
135 }
136
137 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 let idx = match ch
148 .timestamps
149 .binary_search_by(|&ts| ts.partial_cmp(&t).unwrap())
150 {
151 Ok(i) => {
152 return ch.values[i];
154 }
155 Err(i) => {
156 if i == 0 {
158 return ch.values[0]; }
160 if i >= ch.len() {
161 return ch.values[ch.len() - 1]; }
163 i - 1 }
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 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 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
217pub 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 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 pub fn reader(&self) -> &TimeSeriesReader<T> {
253 &self.reader
254 }
255
256 pub fn reader_mut(&mut self) -> &mut TimeSeriesReader<T> {
258 &mut self.reader
259 }
260
261 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 _ctx: &RenderContext,
447 _control_inputs: &[T],
448 _clock_inputs: &[RenderContext],
449 _tick: &ClockTick,
450 ) -> ProcessResult<()> {
451 if !self.playing || self.reader.num_channels() == 0 {
452 for port in self.outputs.iter_mut() {
453 port.write().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.write();
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 ctx = RenderContext::new(0, 4, 44100.0);
674 let tick = ClockTick::new(0, 4, 44100.0, String::new());
675 node.generate(&ctx, &[], &[], &tick).unwrap();
676
677 let port = node.output_port(0).unwrap();
678 let buf = port.read();
679 assert!(ae(buf[0], 1.0));
680 assert!(ae(buf[1], 1.5));
681 assert!(ae(buf[2], 2.0));
682 assert!(ae(buf[3], 2.0));
683 }
684}