1use std::rc::Rc;
27
28use async_trait::async_trait;
29use gwr_components::arbiter::{Arbiter, Arbitrate};
30use gwr_components::connect_port;
31use gwr_components::flow_controls::limiter::Limiter;
32use gwr_components::flow_controls::rate_limiter::RateLimiter;
33use gwr_components::router::{Route, Router};
34use gwr_components::store::{ByteStore, Store};
35use gwr_engine::engine::Engine;
36use gwr_engine::port::PortStateResult;
37use gwr_engine::time::clock::Clock;
38use gwr_engine::traits::{Routable, SimObject};
39use gwr_engine::types::{SimError, SimResult};
40use gwr_model_builder::{EntityDisplay, EntityGet, Runnable};
41use gwr_track::build_aka;
42use gwr_track::entity::Entity;
43use gwr_track::tracker::aka::Aka;
44
45pub const RING_INDEX: usize = 0;
47pub const IO_INDEX: usize = 1;
49
50pub struct RingConfig<T>
51where
52 T: SimObject,
53{
54 rx_buffer_bytes: usize,
55 tx_buffer_bytes: usize,
56 write_limiter: Rc<RateLimiter<T>>,
57}
58
59impl<T> RingConfig<T>
60where
61 T: SimObject,
62{
63 #[must_use]
64 pub fn new(
65 rx_buffer_bytes: usize,
66 tx_buffer_bytes: usize,
67 write_limiter: Rc<RateLimiter<T>>,
68 ) -> Self {
69 Self {
70 rx_buffer_bytes,
71 tx_buffer_bytes,
72 write_limiter,
73 }
74 }
75}
76
77#[derive(EntityGet, EntityDisplay, Runnable)]
78pub struct RingNode<T>
79where
80 T: SimObject + Routable,
81{
82 entity: Rc<Entity>,
83 rx_buffer_limiter: Rc<Limiter<T>>,
84 tx_buffer: Rc<Store<T>>,
85 arbiter: Rc<Arbiter<T>>,
86 router: Rc<Router<T>>,
87}
88
89impl<T> RingNode<T>
90where
91 T: SimObject + Routable,
92{
93 #[expect(clippy::too_many_arguments)]
94 pub fn new_and_register_with_renames(
95 engine: &Engine,
96 clock: &Clock,
97 parent: &Rc<Entity>,
98 name: &str,
99 aka: Option<&Aka>,
100 config: &RingConfig<T>,
101 routing_algorithm: Box<dyn Route<T>>,
102 policy: Box<dyn Arbitrate<T>>,
103 ) -> Result<Rc<Self>, SimError> {
104 let entity = Rc::new(Entity::new(parent, name));
105
106 let rx_buffer_limiter_aka = build_aka!(aka, &entity, &[("ring_rx", "rx")]);
107 let rx_buffer_limiter = Limiter::new_and_register_with_renames(
108 engine,
109 clock,
110 &entity,
111 "limit_rx",
112 Some(&rx_buffer_limiter_aka),
113 config.write_limiter.clone(),
114 );
115 let rx_buffer =
116 ByteStore::new_and_register(engine, clock, &entity, "rx_buf", config.rx_buffer_bytes)?;
117 connect_port!(rx_buffer_limiter, tx => rx_buffer, rx)
118 .expect("Internal ports should connect without error");
119
120 let tx_buffer_limiter = Limiter::new_and_register(
121 engine,
122 clock,
123 &entity,
124 "limit_tx",
125 config.write_limiter.clone(),
126 );
127 let tx_buffer_aka = build_aka!(aka, &entity, &[("ring_tx", "tx")]);
128 let tx_buffer = ByteStore::new_and_register_with_renames(
129 engine,
130 clock,
131 &entity,
132 "tx_buf",
133 Some(&tx_buffer_aka),
134 config.tx_buffer_bytes,
135 )?;
136 connect_port!(tx_buffer_limiter, tx => tx_buffer, rx)
137 .expect("Internal ports should connect without error");
138
139 let router_aka = build_aka!(aka, &entity, &[("io_tx", &format!("tx_{IO_INDEX}"))]);
140 let router = Router::new_and_register_with_renames(
141 engine,
142 clock,
143 &entity,
144 "router",
145 Some(&router_aka),
146 2,
147 routing_algorithm,
148 );
149 connect_port!(rx_buffer, tx => router, rx)
150 .expect("Internal ports should connect without error");
151
152 let arbiter_aka = build_aka!(aka, &entity, &[("io_rx", &format!("rx_{IO_INDEX}"))]);
153 let arbiter = Arbiter::new_and_register_with_renames(
154 engine,
155 clock,
156 &entity,
157 "arb",
158 Some(&arbiter_aka),
159 2,
160 policy,
161 );
162 connect_port!(router, tx, RING_INDEX => arbiter, rx, RING_INDEX)
163 .expect("Internal ports should connect without error");
164 connect_port!(arbiter, tx => tx_buffer_limiter, rx)
165 .expect("Internal ports should connect without error");
166
167 let rc_self = Rc::new(Self {
168 entity,
169 rx_buffer_limiter,
170 tx_buffer,
171 arbiter,
172 router,
173 });
174 engine.register(rc_self.clone());
175 Ok(rc_self)
176 }
177
178 pub fn new_and_register(
179 engine: &Engine,
180 clock: &Clock,
181 parent: &Rc<Entity>,
182 name: &str,
183 config: &RingConfig<T>,
184 routing_algorithm: Box<dyn Route<T>>,
185 policy: Box<dyn Arbitrate<T>>,
186 ) -> Result<Rc<Self>, SimError> {
187 Self::new_and_register_with_renames(
188 engine,
189 clock,
190 parent,
191 name,
192 None,
193 config,
194 routing_algorithm,
195 policy,
196 )
197 }
198
199 pub fn connect_port_ring_tx(&self, port_state: PortStateResult<T>) -> SimResult {
200 self.tx_buffer.connect_port_tx(port_state)
201 }
202
203 pub fn connect_port_io_tx(&self, port_state: PortStateResult<T>) -> SimResult {
204 self.router.connect_port_tx_i(IO_INDEX, port_state)
205 }
206
207 pub fn port_ring_rx(&self) -> PortStateResult<T> {
208 self.rx_buffer_limiter.port_rx()
209 }
210
211 pub fn port_io_rx(&self) -> PortStateResult<T> {
212 self.arbiter.port_rx_i(IO_INDEX)
213 }
214}