sim_lib_stream_host/placement/
lan.rs1use sim_kernel::{CapabilityName, Error, Expr, Result, Symbol};
4use sim_lib_stream_core::{
5 BridgeLatency, ClockDomain, DomainBridgeDescriptor, LatencyClass, PlacedFragment,
6 StreamCapability, StreamEnvelope, TransportProfile,
7};
8
9const DEFAULT_BEATS_PER_BAR: u32 = 4;
10
11pub fn lan_peer_site_symbol() -> Symbol {
13 Symbol::qualified("stream/site", "lan-peer")
14}
15
16pub fn lan_jitter_buffered_mode_symbol() -> Symbol {
18 Symbol::qualified("stream/lan-mode", "jitter-buffered")
19}
20
21pub fn lan_bar_delay_mode_symbol() -> Symbol {
23 Symbol::qualified("stream/lan-mode", "collab-bardelay")
24}
25
26pub fn lan_pinned_sample_refusal_diagnostic() -> Symbol {
28 Symbol::qualified("stream/lan-diagnostic", "pinned-sample-remote-refused")
29}
30
31pub fn lan_pinned_sample_experimental_diagnostic() -> Symbol {
33 Symbol::qualified("stream/lan-diagnostic", "pinned-sample-remote-experimental")
34}
35
36pub fn lan_experimental_remote_sample_capability() -> CapabilityName {
38 CapabilityName::new("stream.lan.experimental-remote-sample")
39}
40
41#[derive(Clone, Copy, Debug, PartialEq, Eq)]
47pub enum LanPlacementMode {
48 JitterBuffered {
50 jitter_packets: u32,
52 latency_comp_frames: u64,
54 },
55 BarDelay {
57 bars: u32,
59 beats_per_bar: u32,
61 tempo_bpm: u32,
63 jitter_packets: u32,
65 latency_comp_frames: u64,
67 },
68}
69
70impl LanPlacementMode {
71 pub fn jitter_buffered(jitter_packets: u32, latency_comp_frames: u64) -> Result<Self> {
86 if jitter_packets == 0 {
87 return Err(Error::Eval(
88 "LAN jitter buffer must retain at least one packet".to_owned(),
89 ));
90 }
91 Ok(Self::JitterBuffered {
92 jitter_packets,
93 latency_comp_frames,
94 })
95 }
96
97 pub fn bar_delay(
102 bars: u32,
103 tempo_bpm: u32,
104 jitter_packets: u32,
105 latency_comp_frames: u64,
106 ) -> Result<Self> {
107 if bars == 0 {
108 return Err(Error::Eval(
109 "LAN bar-delay mode must delay at least one bar".to_owned(),
110 ));
111 }
112 if tempo_bpm == 0 {
113 return Err(Error::Eval(
114 "LAN bar-delay mode tempo must be greater than zero".to_owned(),
115 ));
116 }
117 if jitter_packets == 0 {
118 return Err(Error::Eval(
119 "LAN jitter buffer must retain at least one packet".to_owned(),
120 ));
121 }
122 Ok(Self::BarDelay {
123 bars,
124 beats_per_bar: DEFAULT_BEATS_PER_BAR,
125 tempo_bpm,
126 jitter_packets,
127 latency_comp_frames,
128 })
129 }
130
131 pub fn symbol(self) -> Symbol {
133 match self {
134 Self::JitterBuffered { .. } => lan_jitter_buffered_mode_symbol(),
135 Self::BarDelay { .. } => lan_bar_delay_mode_symbol(),
136 }
137 }
138
139 pub fn latency_class(self) -> LatencyClass {
141 match self {
142 Self::JitterBuffered { .. } => LatencyClass::BufferedPreview,
143 Self::BarDelay { .. } => LatencyClass::CollabBarDelay,
144 }
145 }
146
147 pub fn jitter_packets(self) -> u32 {
149 match self {
150 Self::JitterBuffered { jitter_packets, .. } | Self::BarDelay { jitter_packets, .. } => {
151 jitter_packets
152 }
153 }
154 }
155
156 pub fn latency_comp_frames(self) -> u64 {
158 match self {
159 Self::JitterBuffered {
160 latency_comp_frames,
161 ..
162 }
163 | Self::BarDelay {
164 latency_comp_frames,
165 ..
166 } => latency_comp_frames,
167 }
168 }
169
170 pub fn bar_delay_millis(self) -> Option<u64> {
173 match self {
174 Self::JitterBuffered { .. } => None,
175 Self::BarDelay {
176 bars,
177 beats_per_bar,
178 tempo_bpm,
179 ..
180 } => Some(
181 u64::from(bars)
182 .saturating_mul(u64::from(beats_per_bar))
183 .saturating_mul(60_000)
184 / u64::from(tempo_bpm),
185 ),
186 }
187 }
188
189 pub fn transport_profile(self) -> Result<TransportProfile> {
191 match self {
192 Self::JitterBuffered { .. } => Ok(TransportProfile::lan_buffered_audio_preview()),
193 Self::BarDelay { .. } => TransportProfile::new(
194 Symbol::qualified("stream/profile", "lan-collab-bardelay"),
195 LatencyClass::CollabBarDelay,
196 vec![
197 StreamCapability::Remote,
198 StreamCapability::Bounded,
199 StreamCapability::Preview,
200 StreamCapability::Lossy,
201 ],
202 ),
203 }
204 }
205
206 fn bridges(self) -> Vec<DomainBridgeDescriptor> {
207 vec![
208 DomainBridgeDescriptor::jitter_buffer(self.jitter_packets()),
209 DomainBridgeDescriptor::latency_comp_delay(self.latency_comp_frames()),
210 ]
211 }
212}
213
214#[derive(Clone, Debug, PartialEq, Eq)]
216pub struct LanPlacementRequest {
217 fragment: PlacedFragment,
218 mode: LanPlacementMode,
219 realtime_pinned: bool,
220 capabilities: Vec<CapabilityName>,
221}
222
223impl LanPlacementRequest {
224 pub fn new(fragment: PlacedFragment, mode: LanPlacementMode) -> Self {
227 Self {
228 fragment,
229 mode,
230 realtime_pinned: false,
231 capabilities: Vec::new(),
232 }
233 }
234
235 pub fn with_realtime_pin(mut self, realtime_pinned: bool) -> Self {
237 self.realtime_pinned = realtime_pinned;
238 self
239 }
240
241 pub fn with_capability(mut self, capability: CapabilityName) -> Self {
243 self.capabilities.push(capability);
244 self
245 }
246
247 pub fn plan(&self) -> Result<LanPlacementReport> {
253 let experimental = self
254 .capabilities
255 .contains(&lan_experimental_remote_sample_capability());
256 if self.realtime_pinned && self.fragment_has_sample_domain() && !experimental {
257 let diagnostic = lan_pinned_sample_refusal_diagnostic();
258 return Err(Error::Eval(format!(
259 "{}: pinned sample-domain nodes cannot be sample-locked across LAN",
260 diagnostic.as_qualified_str()
261 )));
262 }
263
264 let mut diagnostics = Vec::new();
265 if self.realtime_pinned && self.fragment_has_sample_domain() {
266 diagnostics.push(lan_pinned_sample_experimental_diagnostic());
267 }
268
269 let profile = self.mode.transport_profile()?;
270 let output_envelopes =
271 remote_output_envelopes(&self.fragment.output_envelopes(), &profile, &diagnostics)?;
272 Ok(LanPlacementReport {
273 fragment_id: self.fragment.id().clone(),
274 site: lan_peer_site_symbol(),
275 mode: self.mode,
276 bridges: self.mode.bridges(),
277 output_envelopes,
278 diagnostics,
279 })
280 }
281
282 fn fragment_has_sample_domain(&self) -> bool {
283 self.fragment
284 .input_edges()
285 .iter()
286 .chain(self.fragment.output_edges())
287 .any(|edge| edge.rate_contract().clock_domain() == ClockDomain::Sample)
288 }
289}
290
291#[derive(Clone, Debug, PartialEq, Eq)]
296pub struct LanPlacementReport {
297 fragment_id: Symbol,
298 site: Symbol,
299 mode: LanPlacementMode,
300 bridges: Vec<DomainBridgeDescriptor>,
301 output_envelopes: Vec<StreamEnvelope>,
302 diagnostics: Vec<Symbol>,
303}
304
305impl LanPlacementReport {
306 pub fn fragment_id(&self) -> &Symbol {
308 &self.fragment_id
309 }
310
311 pub fn site(&self) -> &Symbol {
313 &self.site
314 }
315
316 pub fn mode(&self) -> LanPlacementMode {
318 self.mode
319 }
320
321 pub fn latency_class(&self) -> LatencyClass {
323 self.mode.latency_class()
324 }
325
326 pub fn bridges(&self) -> &[DomainBridgeDescriptor] {
328 &self.bridges
329 }
330
331 pub fn output_envelopes(&self) -> &[StreamEnvelope] {
333 &self.output_envelopes
334 }
335
336 pub fn diagnostics(&self) -> &[Symbol] {
338 &self.diagnostics
339 }
340
341 pub fn added_bridge_latency(&self) -> BridgeLatency {
343 self.bridges
344 .iter()
345 .fold(BridgeLatency::zero(), |latency, bridge| {
346 latency.plus(bridge.latency())
347 })
348 }
349
350 pub fn bar_delay_millis(&self) -> Option<u64> {
352 self.mode.bar_delay_millis()
353 }
354
355 pub fn to_expr(&self) -> Expr {
357 let latency = self.added_bridge_latency();
358 Expr::Map(vec![
359 (
360 Expr::Symbol(Symbol::new("fragment")),
361 Expr::Symbol(self.fragment_id.clone()),
362 ),
363 (
364 Expr::Symbol(Symbol::new("site")),
365 Expr::Symbol(self.site.clone()),
366 ),
367 (
368 Expr::Symbol(Symbol::new("mode")),
369 Expr::Symbol(self.mode.symbol()),
370 ),
371 (
372 Expr::Symbol(Symbol::new("latency-class")),
373 Expr::Symbol(self.latency_class().symbol()),
374 ),
375 (
376 Expr::Symbol(Symbol::new("bar-delay-ms")),
377 Expr::String(self.bar_delay_millis().unwrap_or(0).to_string()),
378 ),
379 (
380 Expr::Symbol(Symbol::new("bridge-latency-frames")),
381 Expr::String(latency.frame_count().to_string()),
382 ),
383 (
384 Expr::Symbol(Symbol::new("bridge-latency-packets")),
385 Expr::String(latency.packet_count().to_string()),
386 ),
387 (
388 Expr::Symbol(Symbol::new("bridges")),
389 Expr::List(
390 self.bridges
391 .iter()
392 .map(|bridge| Expr::Symbol(bridge.kind().symbol()))
393 .collect(),
394 ),
395 ),
396 (
397 Expr::Symbol(Symbol::new("diagnostics")),
398 Expr::List(self.diagnostics.iter().cloned().map(Expr::Symbol).collect()),
399 ),
400 (
401 Expr::Symbol(Symbol::new("output-profiles")),
402 Expr::List(
403 self.output_envelopes
404 .iter()
405 .map(|envelope| Expr::Symbol(envelope.profile().name().clone()))
406 .collect(),
407 ),
408 ),
409 ])
410 }
411}
412
413fn remote_output_envelopes(
414 envelopes: &[StreamEnvelope],
415 profile: &TransportProfile,
416 diagnostics: &[Symbol],
417) -> Result<Vec<StreamEnvelope>> {
418 envelopes
419 .iter()
420 .map(|envelope| {
421 let mut envelope_diagnostics = envelope.diagnostics().to_vec();
422 envelope_diagnostics.extend(diagnostics.iter().cloned());
423 StreamEnvelope::new_with_clock_domains(
424 envelope.stream_id().clone(),
425 envelope.packet_id().clone(),
426 envelope.media(),
427 envelope.direction(),
428 envelope.sequence(),
429 envelope.ticks().to_vec(),
430 envelope.clock_domain(),
431 envelope.clock_domains().to_vec(),
432 profile.clone(),
433 envelope_diagnostics,
434 envelope.packet().clone(),
435 )
436 })
437 .collect()
438}