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
53pub 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
70pub 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
233fn 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;