1use std::cell::RefCell;
14use std::rc::Rc;
15
16use async_trait::async_trait;
17use gwr_engine::engine::Engine;
18use gwr_engine::port::{OutPort, PortStateResult};
19use gwr_engine::traits::{Runnable, SimObject};
20use gwr_engine::types::SimResult;
21use gwr_model_builder::{EntityDisplay, EntityGet};
22use gwr_track::entity::Entity;
23use gwr_track::tracker::aka::Aka;
24
25#[macro_export]
26macro_rules! option_box_repeat {
27 ($value:expr ; $repeat:expr) => {
28 Some(Box::new(std::iter::repeat($value).take($repeat)))
29 };
30}
31use crate::types::DataGenerator;
32use crate::{connect_tx, take_option};
33
34#[macro_export]
35macro_rules! option_box_chain {
36 ($value1:expr , $value2:expr) => {
37 Some(Box::new((*($value1.unwrap())).chain(*($value2.unwrap()))))
38 };
39}
40
41#[derive(EntityGet, EntityDisplay)]
42pub struct Source<T>
43where
44 T: SimObject,
45{
46 entity: Rc<Entity>,
47 data_generator: RefCell<Option<DataGenerator<T>>>,
48 tx: RefCell<Option<OutPort<T>>>,
49}
50
51impl<T> Source<T>
52where
53 T: SimObject,
54{
55 pub fn new_and_register_with_renames(
56 engine: &Engine,
57 parent: &Rc<Entity>,
58 name: &str,
59 aka: Option<&Aka>,
60 data_generator: Option<DataGenerator<T>>,
61 ) -> Rc<Self> {
62 let entity = Rc::new(Entity::new(parent, name));
63 let tx = OutPort::new_with_renames(&entity, "tx", aka);
64 let rc_self = Rc::new(Self {
65 entity,
66 data_generator: RefCell::new(data_generator),
67 tx: RefCell::new(Some(tx)),
68 });
69 engine.register(rc_self.clone());
70 rc_self
71 }
72
73 pub fn new_and_register(
74 engine: &Engine,
75 parent: &Rc<Entity>,
76 name: &str,
77 data_generator: Option<DataGenerator<T>>,
78 ) -> Rc<Self> {
79 Self::new_and_register_with_renames(engine, parent, name, None, data_generator)
80 }
81
82 pub fn set_generator(&self, data_generator: Option<DataGenerator<T>>) {
83 *self.data_generator.borrow_mut() = data_generator;
84 }
85
86 pub fn connect_port_tx(&self, port_state: PortStateResult<T>) -> SimResult {
87 connect_tx!(self.tx, connect ; port_state)
88 }
89}
90
91#[async_trait(?Send)]
92impl<T> Runnable for Source<T>
93where
94 T: SimObject,
95{
96 async fn run(&self) -> SimResult {
97 let mut data_generator = match self.data_generator.borrow_mut().take() {
98 Some(data_generator) => data_generator,
99 None => return Ok(()),
100 };
101
102 let mut tx = take_option!(self.tx);
103 loop {
104 let value = data_generator.next();
105 if let Some(value) = value {
106 self.entity.track_exit(value.id());
107 tx.put(value)?.await;
108 } else {
109 break;
110 }
111 }
112 Ok(())
113 }
114}