rustdv_methodology/
analysis.rs1use std::cell::RefCell;
52use std::rc::Rc;
53
54use crate::component::{Component, ComponentNode};
55use crate::port::{bind_or_panic, sink_of, PortName, PortOwner, PublishIf, SinkHandle};
56
57struct HubInner<T: 'static> {
64 subs: RefCell<Vec<Rc<dyn SinkHandle<T>>>>,
65}
66
67impl<T: 'static> HubInner<T> {
68 fn broadcast(&self, item: &T) {
73 let subs: Vec<Rc<dyn SinkHandle<T>>> = self.subs.borrow().clone();
76 for sub in subs {
77 sub.deliver(item);
78 }
79 }
80}
81
82impl<T: 'static> PublishIf<T> for HubInner<T> {
83 fn write(&self, item: &T) {
84 self.broadcast(item);
85 }
86}
87
88pub struct PublishExport<T: 'static> {
91 inner: Rc<HubInner<T>>,
92}
93
94impl<T: 'static> PublishExport<T> {
95 pub fn connect(&self, owner: &dyn PortOwner, name: PortName<dyn PublishIf<T>>) {
96 bind_or_panic(owner, name, self.inner.clone() as Rc<dyn PublishIf<T>>);
97 }
98}
99
100pub struct SubscribeExport<T: 'static> {
103 inner: Rc<HubInner<T>>,
104}
105
106impl<T: 'static> SubscribeExport<T> {
107 pub fn connect(&self, owner: &dyn PortOwner, name: PortName<dyn SinkHandle<T>>) {
113 match sink_of(owner, name) {
114 Ok(sink) => self.inner.subs.borrow_mut().push(sink),
115 Err(e) => panic!("{e}"),
116 }
117 }
118}
119
120pub struct AnalysisBus<T: 'static> {
131 inner: Rc<HubInner<T>>,
132}
133
134impl<T: 'static> Clone for AnalysisBus<T> {
135 fn clone(&self) -> Self {
137 AnalysisBus { inner: self.inner.clone() }
138 }
139}
140
141impl<T: 'static> Default for AnalysisBus<T> {
142 fn default() -> Self {
143 AnalysisBus::new()
144 }
145}
146
147impl<T: 'static> AnalysisBus<T> {
148 pub fn new() -> AnalysisBus<T> {
149 AnalysisBus { inner: Rc::new(HubInner { subs: RefCell::new(Vec::new()) }) }
150 }
151
152 pub fn pub_export(&self) -> PublishExport<T> {
154 PublishExport { inner: self.inner.clone() }
155 }
156
157 pub fn sub_export(&self) -> SubscribeExport<T> {
160 SubscribeExport { inner: self.inner.clone() }
161 }
162
163 pub fn write(&self, item: &T) {
165 self.inner.broadcast(item);
166 }
167
168 pub fn subscriber_count(&self) -> usize {
170 self.inner.subs.borrow().len()
171 }
172}
173
174impl<T: 'static> Component for AnalysisBus<T> {}
176
177impl<T: 'static> ComponentNode for AnalysisBus<T> {
178 fn node_name(&self) -> &'static str {
179 "AnalysisBus"
180 }
181 fn children_mut(&mut self) -> Vec<(String, &mut (dyn ComponentNode + 'static))> {
182 Vec::new()
183 }
184}
185
186#[cfg(test)]
191mod tests {
192 use super::*;
193 use crate::port::{
194 PortField, PortName, PortOwner, PublishPort, SinkHandle, SubscribePort, Subscriber,
195 };
196 use crate::shared::RustdvShared;
197 use std::any::Any;
198
199 #[derive(Default)]
200 struct Tally {
201 seen: Vec<u8>,
202 }
203 impl Subscriber<u8> for Tally {
204 fn write(&mut self, item: &u8) {
205 self.seen.push(*item);
206 }
207 }
208
209 struct Source {
210 ap: PublishPort<u8>,
211 }
212 impl Source {
213 const AP: PortName<dyn PublishIf<u8>> = PortName::new("ap");
214 }
215 impl PortOwner for Source {
216 fn owner_port_slot(&self, name: &str) -> Option<Rc<dyn Any>> {
217 (name == "ap").then(|| self.ap.slot_any())
218 }
219 fn owner_label(&self) -> &'static str {
220 "Source"
221 }
222 }
223
224 struct Listener {
225 input: SubscribePort<u8>,
226 tally: RustdvShared<Tally>,
227 }
228 impl Listener {
229 const INPUT: PortName<dyn SinkHandle<u8>> = PortName::new("input");
230 fn new() -> Listener {
231 let l = Listener { input: SubscribePort::default(), tally: RustdvShared::default() };
232 l.input.subscribe(l.tally.clone());
233 l
234 }
235 }
236 impl PortOwner for Listener {
237 fn owner_port_slot(&self, name: &str) -> Option<Rc<dyn Any>> {
238 (name == "input").then(|| self.input.slot_any())
239 }
240 fn owner_label(&self) -> &'static str {
241 "Listener"
242 }
243 }
244
245 #[test]
246 fn one_write_reaches_every_subscriber() {
247 let bus: AnalysisBus<u8> = AnalysisBus::new();
248 let src = Source { ap: PublishPort::default() };
249 let a = Listener::new();
250 let b = Listener::new();
251
252 bus.pub_export().connect(&src, Source::AP);
253 bus.sub_export().connect(&a, Listener::INPUT);
254 bus.sub_export().connect(&b, Listener::INPUT);
255 assert_eq!(bus.subscriber_count(), 2);
256
257 src.ap.write(&7);
258 assert_eq!(a.tally.get().seen, vec![7]);
259 assert_eq!(b.tally.get().seen, vec![7], "several subscribers is what makes it a broadcast");
260 }
261
262 #[test]
263 fn subscribers_are_called_in_connection_order() {
264 let bus: AnalysisBus<u8> = AnalysisBus::new();
265 let src = Source { ap: PublishPort::default() };
266 let first = Listener::new();
267 let second = Listener::new();
268 bus.pub_export().connect(&src, Source::AP);
269 bus.sub_export().connect(&first, Listener::INPUT);
270 bus.sub_export().connect(&second, Listener::INPUT);
271
272 for n in 1..=3u8 {
273 src.ap.write(&n);
274 }
275 assert_eq!(first.tally.get().seen, vec![1, 2, 3]);
276 assert_eq!(second.tally.get().seen, vec![1, 2, 3]);
277 }
278
279 #[test]
282 fn the_bus_stores_nothing() {
283 let bus: AnalysisBus<u8> = AnalysisBus::new();
284 let src = Source { ap: PublishPort::default() };
285 bus.pub_export().connect(&src, Source::AP);
286
287 src.ap.write(&1); src.ap.write(&2);
289
290 let late = Listener::new();
291 bus.sub_export().connect(&late, Listener::INPUT);
292 assert!(late.tally.get().seen.is_empty(), "nothing was buffered for a late subscriber");
293
294 src.ap.write(&3);
295 assert_eq!(late.tally.get().seen, vec![3], "only what arrives after it connects");
296 }
297
298 #[test]
301 fn writing_with_no_subscribers_is_legal() {
302 let bus: AnalysisBus<u8> = AnalysisBus::new();
303 let src = Source { ap: PublishPort::default() };
304 bus.pub_export().connect(&src, Source::AP);
305 assert_eq!(bus.subscriber_count(), 0);
306 src.ap.write(&1); }
308
309 #[test]
310 fn an_unconnected_publish_port_does_not_panic() {
311 let src = Source { ap: PublishPort::default() };
312 assert!(!src.ap.has_subscribers());
313 src.ap.write(&1); }
315
316 #[test]
319 fn delivery_happens_before_write_returns() {
320 let bus: AnalysisBus<u8> = AnalysisBus::new();
321 let src = Source { ap: PublishPort::default() };
322 let sub = Listener::new();
323 bus.pub_export().connect(&src, Source::AP);
324 bus.sub_export().connect(&sub, Listener::INPUT);
325
326 src.ap.write(&5);
327 assert_eq!(sub.tally.get().seen, vec![5], "already delivered, no scheduling in between");
328 }
329
330 #[test]
333 #[should_panic(expected = "has no subscriber")]
334 fn connecting_a_subscriber_with_no_sink_is_a_named_error() {
335 let bus: AnalysisBus<u8> = AnalysisBus::new();
336 let bare = Listener { input: SubscribePort::default(), tally: RustdvShared::default() };
337 bus.sub_export().connect(&bare, Listener::INPUT);
338 }
339
340 #[test]
341 fn a_clone_is_the_same_bus() {
342 let bus: AnalysisBus<u8> = AnalysisBus::new();
343 let other = bus.clone();
344 let sub = Listener::new();
345 other.sub_export().connect(&sub, Listener::INPUT);
346 assert_eq!(bus.subscriber_count(), 1, "two handles, one subscriber list");
347 bus.write(&4);
348 assert_eq!(sub.tally.get().seen, vec![4]);
349 }
350}