1use sva_samples::Extent;
4
5use super::RenderConfig;
6use super::table::spill::Spill;
7use super::table::{Pulled, Table};
8use super::until::{Known, Until};
9use crate::cache::Recording;
10use crate::error::EngineError;
11use crate::flops::Work;
12use crate::query::{DEFAULT_FRAME_SECS, Representation};
13
14pub(super) struct Driver {
15 pub(super) table: Table,
16 pub(super) start: i64,
17 pub(super) at: i64,
18 last: i64,
19 block: usize,
20 until: Option<Until>,
21 frame: usize,
22 stop: Option<i64>,
23 end: Option<i64>,
24 keep: i64,
26 heard: i64,
28 most_bytes: usize,
29 pub(super) work: Work,
30 pub(super) recording: Recording,
31 pub(super) spill: Option<Spill>,
32}
33
34pub struct Block {
35 planes: Vec<Vec<f64>>,
36 start: i64,
37}
38
39impl Block {
40 pub fn start(&self) -> i64 {
41 self.start
42 }
43
44 pub fn len(&self) -> usize {
45 self.planes[0].len()
46 }
47
48 pub fn is_empty(&self) -> bool {
49 self.len() == 0
50 }
51
52 pub fn width(&self) -> usize {
53 self.planes.len()
54 }
55
56 pub fn plane(&self, c: usize) -> &[f64] {
57 &self.planes[c]
58 }
59
60 pub(super) fn widened(mut self, width: usize) -> Block {
62 if self.planes.len() == 1 {
63 let copies = vec![self.planes[0].clone(); width.saturating_sub(1)];
64 self.planes.extend(copies);
65 }
66 self
67 }
68}
69
70pub(super) fn frame(config: &RenderConfig, rate: u32) -> usize {
72 let secs = config
73 .asks
74 .iter()
75 .find_map(|ask| match ask.representation {
76 Representation::Envelope { frame_secs } => Some(frame_secs),
77 _ => None,
78 })
79 .flatten()
80 .unwrap_or(DEFAULT_FRAME_SECS);
81 ((secs * f64::from(rate)).round() as usize).max(1)
82}
83
84impl Driver {
85 pub(super) fn new(
86 table: Table,
87 range: Extent,
88 block: usize,
89 config: &RenderConfig,
90 recording: Recording,
91 ) -> Driver {
92 let frame = frame(config, config.rate);
93 Driver {
94 table,
95 start: range.start,
96 at: range.start,
97 heard: range.start,
98 last: range.end,
99 block,
100 keep: config.until.as_ref().map_or(0, |_| frame as i64),
101 until: config.until.clone(),
102 frame,
103 stop: None,
104 end: None,
105 most_bytes: 0,
106 work: Work {
107 waves: Some(0),
108 ..Work::default()
109 },
110 recording,
111 spill: None,
112 }
113 }
114
115 pub(super) fn replace(&mut self, table: Table, last: i64) {
116 self.table = table;
117 self.bound(last);
118 }
119
120 pub(super) fn bound(&mut self, last: i64) {
122 self.last = last;
123 if self.stop.is_none() {
124 self.end = None;
125 }
126 }
127
128 pub(super) fn last(&self) -> i64 {
129 self.last
130 }
131
132 fn next_to(&self, n: usize) -> i64 {
133 self.last.min(self.at.saturating_add(n as i64)).max(self.at)
134 }
135
136 pub(super) fn read(&mut self, n: usize) -> Result<Option<Block>, EngineError> {
138 let from = self.at;
139 if !self.pulled(n)? {
140 return Ok(None);
141 }
142 let end = self.end.map_or(self.at, |end| end.clamp(from, self.at));
143 let held = self.table.samples(self.table.root, Extent::new(from, end));
144 Ok(Some(Block {
145 planes: held.planes,
146 start: from,
147 }))
148 }
149
150 pub(super) fn skip(&mut self, to: i64) -> Result<Vec<usize>, EngineError> {
154 if self.end.is_some_and(|end| self.at >= end) {
155 return Ok(Vec::new());
156 }
157 let to = to.min(self.last);
158 let silenced = self.table.skipped(Extent::new(to, self.last))?;
159 self.recording.reach(to);
160 (self.at, self.heard) = (to, to);
161 let future = (to < self.last).then(|| Extent::new(to, self.last));
162 self.table.release(future, Extent::NOWHERE, self.start);
163 if self.last <= to {
164 self.end = Some(to);
165 }
166 Ok(silenced)
167 }
168
169 pub(super) fn pull(&mut self) -> Result<bool, EngineError> {
170 self.pulled(self.block)
171 }
172
173 pub(super) fn pulled(&mut self, n: usize) -> Result<bool, EngineError> {
175 let from = self.at;
176 if self.end.is_some_and(|end| from >= end) {
177 return Ok(false);
178 }
179 let to = self.next_to(n);
180 self.recording.reach(from);
181 let window = Extent::new(from, to);
182 if from == self.start {
183 let history = self
184 .table
185 .history(window, self.block as i64, &mut self.recording)?;
186 self.priced(&history);
187 }
188 let asked = self.spill.as_ref().map(|_| self.table.demand(window));
189 let pulled = self.table.pull(window, &mut self.recording)?;
190 self.priced(&pulled);
191 if let (Some(spill), Some(asked)) = (&mut self.spill, asked) {
192 spill.take(&self.table, &asked);
193 }
194 self.work.samples += (to - from) as u64;
195 self.at = to;
196 self.settle(from, to);
197 let future = (to < self.last).then(|| Extent::new(to, self.last));
198 let keep = Extent::new(from.min(to.saturating_sub(self.keep)), to);
199 self.table.release(future, keep, self.start);
200 Ok(true)
201 }
202
203 fn priced(&mut self, pulled: &Pulled) {
204 self.work.priced_flops += pulled.priced;
205 self.work.waves = self.work.waves.map(|held| held + pulled.waves);
206 self.most_bytes = self.most_bytes.max(pulled.most_bytes);
207 }
208
209 pub(super) fn most_bytes(&self) -> usize {
210 self.most_bytes
211 }
212
213 fn settle(&mut self, from: i64, to: i64) {
216 if self.end.is_some() {
217 return;
218 }
219 if self.last <= to {
220 self.end = Some(to);
221 }
222 let Some(until) = &self.until else {
223 return;
224 };
225 let frame = self.frame as i64;
226 let open = self.start + (from - 1 - self.start).div_euclid(frame) * frame;
227 let base = open.max(self.heard);
228 let heard = self.table.samples(self.table.root, Extent::new(base, to));
229 let planes: Vec<&[f64]> = heard.planes.iter().map(Vec::as_slice).collect();
230 let rate = self.table.values[self.table.root].grid.rate;
231 let known = Known::new(planes, base, self.start, self.frame, rate, to == self.last);
232 if let Some(at) = until.first(&known, base, to) {
233 self.stop = Some(at);
234 self.end = Some(at.max(from));
235 }
236 }
237
238 pub(super) fn end(&self) -> Option<i64> {
239 self.end
240 }
241
242 pub(super) fn stop(&self) -> Option<i64> {
243 self.stop
244 }
245}