Skip to main content

sim_lib_stream_host/placement/
lan.rs

1//! LAN peer placement policy for stream fragments.
2
3use 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
11/// Stable site name for a non-real-time node hosted by a LAN peer.
12pub fn lan_peer_site_symbol() -> Symbol {
13    Symbol::qualified("stream/site", "lan-peer")
14}
15
16/// Stable mode name for jitter-buffered LAN placement.
17pub fn lan_jitter_buffered_mode_symbol() -> Symbol {
18    Symbol::qualified("stream/lan-mode", "jitter-buffered")
19}
20
21/// Stable mode name for musically aligned bar-delay collaboration.
22pub fn lan_bar_delay_mode_symbol() -> Symbol {
23    Symbol::qualified("stream/lan-mode", "collab-bardelay")
24}
25
26/// Diagnostic emitted when a pinned sample-domain node is refused across LAN.
27pub fn lan_pinned_sample_refusal_diagnostic() -> Symbol {
28    Symbol::qualified("stream/lan-diagnostic", "pinned-sample-remote-refused")
29}
30
31/// Diagnostic emitted when experimental pinned sample-domain LAN placement is used.
32pub fn lan_pinned_sample_experimental_diagnostic() -> Symbol {
33    Symbol::qualified("stream/lan-diagnostic", "pinned-sample-remote-experimental")
34}
35
36/// Capability required to try pinned sample-domain placement across LAN.
37pub fn lan_experimental_remote_sample_capability() -> CapabilityName {
38    CapabilityName::new("stream.lan.experimental-remote-sample")
39}
40
41/// Placement mode for a stream fragment hosted on a LAN peer.
42///
43/// Selects how packets crossing the LAN are buffered and time-aligned: a plain
44/// jitter buffer for buffered preview, or a musically aligned bar delay for
45/// collaborative play.
46#[derive(Clone, Copy, Debug, PartialEq, Eq)]
47pub enum LanPlacementMode {
48    /// Jitter-buffered preview placement.
49    JitterBuffered {
50        /// Packets retained in the jitter buffer.
51        jitter_packets: u32,
52        /// Latency-compensation delay applied, in frames.
53        latency_comp_frames: u64,
54    },
55    /// Musically aligned collaborative placement delayed by whole bars.
56    BarDelay {
57        /// Number of bars of alignment delay.
58        bars: u32,
59        /// Beats per bar used to size the bar delay.
60        beats_per_bar: u32,
61        /// Tempo in beats per minute used to size the bar delay.
62        tempo_bpm: u32,
63        /// Packets retained in the jitter buffer.
64        jitter_packets: u32,
65        /// Latency-compensation delay applied, in frames.
66        latency_comp_frames: u64,
67    },
68}
69
70impl LanPlacementMode {
71    /// Builds a jitter-buffered mode retaining `jitter_packets` packets.
72    ///
73    /// Returns an evaluation error when `jitter_packets` is zero.
74    ///
75    /// # Examples
76    ///
77    /// ```
78    /// use sim_lib_stream_host::LanPlacementMode;
79    ///
80    /// let mode = LanPlacementMode::jitter_buffered(4, 128).unwrap();
81    /// assert_eq!(mode.jitter_packets(), 4);
82    /// assert_eq!(mode.latency_comp_frames(), 128);
83    /// assert!(mode.bar_delay_millis().is_none());
84    /// ```
85    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    /// Builds a bar-delay mode delaying `bars` bars at `tempo_bpm`.
98    ///
99    /// Uses a default of four beats per bar. Returns an evaluation error when
100    /// `bars`, `tempo_bpm`, or `jitter_packets` is zero.
101    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    /// Returns the stable mode symbol.
132    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    /// Returns the latency class this mode places the fragment into.
140    pub fn latency_class(self) -> LatencyClass {
141        match self {
142            Self::JitterBuffered { .. } => LatencyClass::BufferedPreview,
143            Self::BarDelay { .. } => LatencyClass::CollabBarDelay,
144        }
145    }
146
147    /// Returns the number of packets retained in the jitter buffer.
148    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    /// Returns the latency-compensation delay in frames.
157    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    /// Returns the bar-delay length in milliseconds, or `None` for
171    /// jitter-buffered placement.
172    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    /// Returns the transport profile advertised for this mode.
190    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/// Request to place a stream fragment on a LAN peer under a chosen mode.
215#[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    /// Builds a request to place `fragment` using `mode`, unpinned and with no
225    /// extra capabilities.
226    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    /// Marks whether the fragment is pinned to realtime (sample-locked) play.
236    pub fn with_realtime_pin(mut self, realtime_pinned: bool) -> Self {
237        self.realtime_pinned = realtime_pinned;
238        self
239    }
240
241    /// Grants an additional capability to the request.
242    pub fn with_capability(mut self, capability: CapabilityName) -> Self {
243        self.capabilities.push(capability);
244        self
245    }
246
247    /// Plans the placement, returning a report or an evaluation error.
248    ///
249    /// Refuses a realtime-pinned sample-domain fragment across the LAN unless
250    /// the experimental remote-sample capability is granted, in which case it
251    /// proceeds and records an experimental diagnostic.
252    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/// Outcome of planning a LAN fragment placement.
292///
293/// Records the placement site, mode, the domain bridges inserted, the rewritten
294/// output envelopes carrying the remote transport profile, and any diagnostics.
295#[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    /// Returns the placed fragment identifier.
307    pub fn fragment_id(&self) -> &Symbol {
308        &self.fragment_id
309    }
310
311    /// Returns the placement site symbol.
312    pub fn site(&self) -> &Symbol {
313        &self.site
314    }
315
316    /// Returns the placement mode.
317    pub fn mode(&self) -> LanPlacementMode {
318        self.mode
319    }
320
321    /// Returns the latency class of the placement.
322    pub fn latency_class(&self) -> LatencyClass {
323        self.mode.latency_class()
324    }
325
326    /// Returns the domain bridges inserted by the placement.
327    pub fn bridges(&self) -> &[DomainBridgeDescriptor] {
328        &self.bridges
329    }
330
331    /// Returns the rewritten output envelopes.
332    pub fn output_envelopes(&self) -> &[StreamEnvelope] {
333        &self.output_envelopes
334    }
335
336    /// Returns the diagnostics recorded during planning.
337    pub fn diagnostics(&self) -> &[Symbol] {
338        &self.diagnostics
339    }
340
341    /// Returns the total latency added by the inserted bridges.
342    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    /// Returns the bar-delay length in milliseconds when the mode uses one.
351    pub fn bar_delay_millis(&self) -> Option<u64> {
352        self.mode.bar_delay_millis()
353    }
354
355    /// Builds a browse/inspection expression summarizing the placement.
356    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}