Skip to main content

prns_runtime_embassy/manifold/driver/egress/
mod.rs

1use heapless::Vec as HeaplessVec;
2
3use crate::engine::{EngineReaction, FanTarget, InstantMillis, Journaled};
4use crate::interfaces::InterfaceIfac;
5use crate::interfaces::{InterfaceDescriptor, InterfaceId, InterfaceKind};
6use crate::manifold::announce_pacer::{AnnouncePacer, FixedPacerQueue};
7use crate::manifold::grant::{FrameTarget, LaneWriteOutcome, ManifoldLaneWriter};
8use crate::manifold::interface_seam::EMBEDDED_MAX_WIRE_FRAME_LEN;
9use crate::manifold::kernel::{
10    route_reaction as route_engine_reaction, AnnounceDirective, DirectiveEgress,
11};
12
13fn lane_serves(lane_key: InterfaceId, target: InterfaceId) -> bool {
14    if lane_key == target {
15        return true;
16    }
17    match (lane_key.kind(), target.kind()) {
18        (Some(supervisor), Some(child)) => supervisor.member_kind() == Some(child),
19        _ => false,
20    }
21}
22
23#[derive(Debug, Clone, Copy, PartialEq, Eq)]
24#[must_use]
25pub enum EgressOutcome {
26    Enqueued,
27    LaneFull {
28        lane: InterfaceId,
29    },
30    FrameTooLarge {
31        lane: InterfaceId,
32        frame_len: usize,
33        capacity: usize,
34    },
35    NoLane,
36}
37
38fn egress_outcome(lane: InterfaceId, outcome: LaneWriteOutcome) -> EgressOutcome {
39    match outcome {
40        LaneWriteOutcome::Written => EgressOutcome::Enqueued,
41        LaneWriteOutcome::Full => EgressOutcome::LaneFull { lane },
42        LaneWriteOutcome::FrameTooLarge {
43            frame_len,
44            capacity,
45        } => EgressOutcome::FrameTooLarge {
46            lane,
47            frame_len,
48            capacity,
49        },
50    }
51}
52
53/// Nonblocking direct and fleet egress.
54pub trait ManifoldEgress {
55    fn enqueue(&mut self, target: InterfaceId, bytes: &[u8]) -> EgressOutcome;
56    fn enqueue_broadcast(
57        &mut self,
58        supervisor: InterfaceKind,
59        fan: FanTarget,
60        bytes: &[u8],
61    ) -> EgressOutcome;
62    fn lane_for(&self, target: InterfaceId) -> Option<InterfaceId> {
63        Some(target)
64    }
65    fn fleet_lane(&self, _supervisor: InterfaceKind) -> Option<InterfaceId> {
66        None
67    }
68}
69
70/// Fixed-set egress with erased slot sizes, allowing heterogeneous lanes in one borrowed slice without allocation.
71pub struct EmbassyEgress<'a> {
72    lanes: &'a mut [(InterfaceId, &'a mut dyn ManifoldLaneWriter)],
73}
74
75impl<'a> EmbassyEgress<'a> {
76    #[must_use]
77    pub fn new(lanes: &'a mut [(InterfaceId, &'a mut dyn ManifoldLaneWriter)]) -> Self {
78        Self { lanes }
79    }
80}
81
82impl ManifoldEgress for EmbassyEgress<'_> {
83    fn enqueue(&mut self, target: InterfaceId, bytes: &[u8]) -> EgressOutcome {
84        for (id, producer) in self.lanes.iter_mut() {
85            if lane_serves(*id, target) {
86                return egress_outcome(*id, producer.try_write(FrameTarget::Direct(target), bytes));
87            }
88        }
89        EgressOutcome::NoLane
90    }
91
92    fn enqueue_broadcast(
93        &mut self,
94        supervisor: InterfaceKind,
95        fan: FanTarget,
96        bytes: &[u8],
97    ) -> EgressOutcome {
98        for (id, producer) in self.lanes.iter_mut() {
99            if id.kind() == Some(supervisor) {
100                return egress_outcome(*id, producer.try_write(FrameTarget::Fan(fan), bytes));
101            }
102        }
103        EgressOutcome::NoLane
104    }
105
106    fn lane_for(&self, target: InterfaceId) -> Option<InterfaceId> {
107        self.lanes
108            .iter()
109            .map(|(id, _)| *id)
110            .find(|id| lane_serves(*id, target))
111    }
112
113    fn fleet_lane(&self, supervisor: InterfaceKind) -> Option<InterfaceId> {
114        self.lanes
115            .iter()
116            .map(|(id, _)| *id)
117            .find(|id| id.kind() == Some(supervisor))
118    }
119}
120
121pub(super) const MAX_PACED_INTERFACES: usize = 2;
122const PACER_DEPTH: usize = 2;
123
124pub(super) struct InterfacePacer {
125    pub(super) id: InterfaceId,
126    pacer: AnnouncePacer<FixedPacerQueue<PACER_DEPTH, FrameTarget>, FrameTarget>,
127}
128
129impl InterfacePacer {
130    pub(super) fn from_descriptor(id: InterfaceId, descriptor: &InterfaceDescriptor) -> Self {
131        Self {
132            id,
133            pacer: AnnouncePacer::new(descriptor.announce_bandwidth_cap, descriptor.bitrate),
134        }
135    }
136}
137
138pub(super) fn route_reaction(
139    reaction: EngineReaction<'_>,
140    egress: &mut impl ManifoldEgress,
141    ifacs: &[InterfaceIfac],
142    pacers: &mut [InterfacePacer],
143    now: InstantMillis,
144    app: &mut impl FnMut(Journaled<'_>),
145) {
146    let mut directive_egress = EmbassyDirectiveEgress {
147        egress,
148        ifacs,
149        pacers,
150        now,
151    };
152    route_engine_reaction(reaction, &mut directive_egress, app);
153}
154
155struct EmbassyDirectiveEgress<'a, E> {
156    egress: &'a mut E,
157    ifacs: &'a [InterfaceIfac],
158    pacers: &'a mut [InterfacePacer],
159    now: InstantMillis,
160}
161
162impl<E: ManifoldEgress> EmbassyDirectiveEgress<'_, E> {
163    fn offer_to_fleet_pacer(
164        &mut self,
165        supervisor: InterfaceKind,
166        fan: FanTarget,
167        bytes: &[u8],
168        hops: u8,
169    ) {
170        let Some(lane) = self.egress.fleet_lane(supervisor) else {
171            enqueue_broadcast_for_wire(self.egress, self.ifacs, supervisor, fan, bytes);
172            return;
173        };
174        match self.pacers.iter_mut().find(|entry| entry.id == lane) {
175            Some(entry) => {
176                let _ = entry.pacer.offer_tagged(
177                    bytes,
178                    hops,
179                    self.now,
180                    FrameTarget::Fan(fan),
181                    |frame, target| {
182                        enqueue_paced_for_wire(self.egress, self.ifacs, lane, target, frame)
183                    },
184                );
185            }
186            None => {
187                enqueue_broadcast_for_wire(self.egress, self.ifacs, supervisor, fan, bytes);
188            }
189        }
190    }
191}
192
193impl<E: ManifoldEgress> DirectiveEgress for EmbassyDirectiveEgress<'_, E> {
194    fn send(&mut self, target: InterfaceId, bytes: &[u8]) {
195        enqueue_for_wire(self.egress, self.ifacs, target, bytes);
196    }
197
198    fn send_announce(&mut self, target: InterfaceId, announce: AnnounceDirective<'_>) {
199        offer_to_pacer(
200            self.pacers,
201            target,
202            announce.bytes(),
203            announce.hops(),
204            self.now,
205            self.egress,
206            self.ifacs,
207        );
208    }
209
210    fn send_to_fleet(&mut self, supervisor: InterfaceKind, fan: FanTarget, bytes: &[u8]) {
211        enqueue_broadcast_for_wire(self.egress, self.ifacs, supervisor, fan, bytes);
212    }
213
214    fn send_announce_to_fleet(
215        &mut self,
216        supervisor: InterfaceKind,
217        fan: FanTarget,
218        announce: AnnounceDirective<'_>,
219    ) {
220        self.offer_to_fleet_pacer(supervisor, fan, announce.bytes(), announce.hops());
221    }
222
223    fn emit_frame(
224        &mut self,
225        target: InterfaceId,
226        _size_hint: usize,
227        fill: &mut dyn FnMut(&mut [u8]) -> Option<usize>,
228    ) {
229        emit_for_wire(self.egress, self.ifacs, target, fill);
230    }
231}
232
233/// Erased slot sizes require one bounded stack buffer before the frame enters its lane. `fill` runs exactly once even when the lane is full.
234fn emit_for_wire(
235    egress: &mut impl ManifoldEgress,
236    ifacs: &[InterfaceIfac],
237    target: InterfaceId,
238    fill: &mut dyn FnMut(&mut [u8]) -> Option<usize>,
239) {
240    let mut frame = [0u8; EMBEDDED_MAX_WIRE_FRAME_LEN];
241    if let Some(len) = fill(&mut frame) {
242        enqueue_for_wire(egress, ifacs, target, &frame[..len]);
243    }
244}
245
246pub(super) fn ifac_for(ifacs: &[InterfaceIfac], id: InterfaceId) -> Option<&InterfaceIfac> {
247    if ifacs.is_empty() {
248        return None;
249    }
250    ifacs.iter().find(|entry| entry.id == id)
251}
252
253pub(super) fn enqueue_for_wire(
254    egress: &mut impl ManifoldEgress,
255    ifacs: &[InterfaceIfac],
256    target: InterfaceId,
257    bytes: &[u8],
258) {
259    let lane = egress.lane_for(target).unwrap_or(target);
260    match ifac_for(ifacs, lane) {
261        Some(entry) => {
262            let mut wire = [0u8; EMBEDDED_MAX_WIRE_FRAME_LEN];
263            if let Some(masked_len) = entry.context.mask_outbound(bytes, &mut wire) {
264                let _ = egress.enqueue(target, &wire[..masked_len]);
265            }
266        }
267        None => {
268            let _ = egress.enqueue(target, bytes);
269        }
270    }
271}
272
273pub(super) fn enqueue_broadcast_for_wire(
274    egress: &mut impl ManifoldEgress,
275    ifacs: &[InterfaceIfac],
276    supervisor: InterfaceKind,
277    fan: FanTarget,
278    bytes: &[u8],
279) {
280    match egress
281        .fleet_lane(supervisor)
282        .and_then(|lane| ifac_for(ifacs, lane))
283    {
284        Some(entry) => {
285            let mut wire = [0u8; EMBEDDED_MAX_WIRE_FRAME_LEN];
286            if let Some(masked_len) = entry.context.mask_outbound(bytes, &mut wire) {
287                let _ = egress.enqueue_broadcast(supervisor, fan, &wire[..masked_len]);
288            }
289        }
290        None => {
291            let _ = egress.enqueue_broadcast(supervisor, fan, bytes);
292        }
293    }
294}
295
296fn offer_to_pacer(
297    pacers: &mut [InterfacePacer],
298    target: InterfaceId,
299    bytes: &[u8],
300    hops: u8,
301    now: InstantMillis,
302    egress: &mut impl ManifoldEgress,
303    ifacs: &[InterfaceIfac],
304) {
305    let lane = egress.lane_for(target).unwrap_or(target);
306    match pacers.iter_mut().find(|entry| entry.id == lane) {
307        Some(entry) => {
308            let _ = entry.pacer.offer_tagged(
309                bytes,
310                hops,
311                now,
312                FrameTarget::Direct(target),
313                |frame, target| enqueue_paced_for_wire(egress, ifacs, lane, target, frame),
314            );
315        }
316        None => enqueue_for_wire(egress, ifacs, target, bytes),
317    }
318}
319
320fn enqueue_paced_for_wire(
321    egress: &mut impl ManifoldEgress,
322    ifacs: &[InterfaceIfac],
323    lane: InterfaceId,
324    target: FrameTarget,
325    bytes: &[u8],
326) {
327    match target {
328        FrameTarget::Direct(target) => enqueue_for_wire(egress, ifacs, target, bytes),
329        FrameTarget::Fan(fan) => {
330            if let Some(supervisor) = lane.kind() {
331                enqueue_broadcast_for_wire(egress, ifacs, supervisor, fan, bytes);
332            }
333        }
334    }
335}
336
337pub(super) fn flush_due_pacers(
338    pacers: &mut [InterfacePacer],
339    now: InstantMillis,
340    egress: &mut impl ManifoldEgress,
341    ifacs: &[InterfaceIfac],
342) {
343    for entry in pacers.iter_mut() {
344        let lane = entry.id;
345        let _ = entry.pacer.release_due_tagged(now, |frame, target| {
346            enqueue_paced_for_wire(egress, ifacs, lane, target, frame)
347        });
348    }
349}
350
351pub(super) fn soonest_pacer_release(pacers: &[InterfacePacer]) -> Option<InstantMillis> {
352    pacers
353        .iter()
354        .filter_map(|entry| entry.pacer.next_release())
355        .min_by_key(|deadline| deadline.0)
356}
357
358pub struct PooledEgress<const LANE_COUNT: usize> {
359    pub(crate) lanes: HeaplessVec<(InterfaceId, &'static mut dyn ManifoldLaneWriter), LANE_COUNT>,
360}
361
362impl<const LANE_COUNT: usize> PooledEgress<LANE_COUNT> {
363    #[must_use]
364    pub fn new() -> Self {
365        Self {
366            lanes: HeaplessVec::new(),
367        }
368    }
369
370    pub(crate) fn push(
371        &mut self,
372        id: InterfaceId,
373        producer: &'static mut dyn ManifoldLaneWriter,
374    ) -> Result<(), &'static mut dyn ManifoldLaneWriter> {
375        self.lanes
376            .push((id, producer))
377            .map_err(|(_, producer)| producer)
378    }
379
380    pub(crate) fn retag(&mut self, old_id: InterfaceId, new_id: InterfaceId) {
381        for (id, _) in self.lanes.iter_mut() {
382            if *id == old_id {
383                *id = new_id;
384            }
385        }
386    }
387}
388
389impl<const LANE_COUNT: usize> ManifoldEgress for PooledEgress<LANE_COUNT> {
390    fn enqueue(&mut self, target: InterfaceId, bytes: &[u8]) -> EgressOutcome {
391        for (id, producer) in self.lanes.iter_mut() {
392            if lane_serves(*id, target) {
393                return egress_outcome(*id, producer.try_write(FrameTarget::Direct(target), bytes));
394            }
395        }
396        EgressOutcome::NoLane
397    }
398
399    fn enqueue_broadcast(
400        &mut self,
401        supervisor: InterfaceKind,
402        fan: FanTarget,
403        bytes: &[u8],
404    ) -> EgressOutcome {
405        for (id, producer) in self.lanes.iter_mut() {
406            if id.kind() == Some(supervisor) {
407                return egress_outcome(*id, producer.try_write(FrameTarget::Fan(fan), bytes));
408            }
409        }
410        EgressOutcome::NoLane
411    }
412
413    fn lane_for(&self, target: InterfaceId) -> Option<InterfaceId> {
414        self.lanes
415            .iter()
416            .map(|(id, _)| *id)
417            .find(|id| lane_serves(*id, target))
418    }
419
420    fn fleet_lane(&self, supervisor: InterfaceKind) -> Option<InterfaceId> {
421        self.lanes
422            .iter()
423            .map(|(id, _)| *id)
424            .find(|id| id.kind() == Some(supervisor))
425    }
426}
427
428impl<const LANE_COUNT: usize> Default for PooledEgress<LANE_COUNT> {
429    fn default() -> Self {
430        Self::new()
431    }
432}
433
434#[cfg(test)]
435mod tests;