1#[cfg(not(feature = "gen_blocks"))]
2use async_stream::stream;
3use futures_util::{FutureExt, Stream};
4
5#[cfg(feature = "gen_blocks")]
6use crate::stream::stream;
7use crate::{BoxComponent, Child, Component, ComponentMessage, ComponentSender};
8
9#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
11#[non_exhaustive]
12pub enum RunEvent<T, E> {
13 Event(T),
15 UpdateErr(E),
17 RenderErr(E),
19}
20
21impl<T, E> RunEvent<T, E> {
22 pub fn flatten(self) -> Result<T, E> {
24 match self {
25 RunEvent::Event(t) => Ok(t),
26 RunEvent::UpdateErr(e) | RunEvent::RenderErr(e) => Err(e),
27 }
28 }
29}
30
31pub struct Root<T: Component> {
33 model: T,
34 sender: ComponentSender<T>,
35}
36
37impl<T: Component> Root<T> {
38 pub async fn init<'a>(init: impl Into<T::Init<'a>>) -> Result<Self, T::Error> {
40 let sender = ComponentSender::new();
41 let model = T::init(init.into(), &sender).await?;
42 Ok(Self { model, sender })
43 }
44
45 pub(crate) fn new(model: T, sender: ComponentSender<T>) -> Self {
46 Self { model, sender }
47 }
48
49 pub fn post(&mut self, message: T::Message) {
51 self.sender.post(message);
52 }
53
54 pub async fn emit(&mut self, message: T::Message) -> Result<bool, T::Error> {
56 self.model.update(message, &self.sender).await
57 }
58
59 pub fn sender(&self) -> &ComponentSender<T> {
61 &self.sender
62 }
63
64 pub fn into_child(self) -> Child<T> {
66 Child::new(self.model, self.sender)
67 }
68
69 pub fn run(&mut self) -> impl Stream<Item = RunEvent<T::Event, T::Error>> + use<'_, T> {
71 run_events_impl(&mut self.model, &self.sender)
72 }
73}
74
75fn run_events_impl<'a, T: Component>(
76 model: &'a mut T,
77 sender: &'a ComponentSender<T>,
78) -> impl Stream<Item = RunEvent<T::Event, T::Error>> + 'a {
79 stream! {
80 if let Err(e) = model.render(sender) {
81 yield RunEvent::RenderErr(e);
82 }
83 if let Err(e) = model.render_children() {
84 yield RunEvent::RenderErr(e);
85 }
86 loop {
87 let fut_start = model.start(sender);
88 let fut_recv = sender.wait();
89 futures_util::select! {
90 x = fut_start.fuse() => match x {},
91 _ = fut_recv.fuse() => {}
92 }
93 let mut need_render = false;
94 let mut children_need_render = match model.update_children().await {
95 Ok(v) => v,
96 Err(e) => {
97 yield RunEvent::UpdateErr(e);
98 false
99 }
100 };
101 for msg in sender.fetch_all() {
102 match msg {
103 ComponentMessage::Message(msg) => {
104 need_render |= match model.update(msg, sender).await {
105 Ok(v) => v,
106 Err(e) => {
107 yield RunEvent::UpdateErr(e);
108 false
109 }
110 };
111 }
112 ComponentMessage::Event(e) => yield RunEvent::Event(e),
113 };
114 }
115 children_need_render |= need_render;
116 if need_render && let Err(e) = model.render(sender) {
117 yield RunEvent::RenderErr(e);
118 }
119 if children_need_render && let Err(e) = model.render_children() {
120 yield RunEvent::RenderErr(e);
121 }
122 }
123 }
124}
125
126impl<T: Component + 'static> Root<T> {
127 pub fn into_boxed(self) -> Root<BoxComponent<T::Message, T::Event, T::Error>> {
129 let sender = ComponentSender(self.sender.0);
130 Root::new(BoxComponent::new(self.model), sender)
131 }
132}
133
134#[cfg(test)]
135mod test {
136 use async_stream::stream;
137 use futures_util::{Stream, StreamExt};
138
139 use crate::*;
140
141 struct TestComponent;
142
143 #[derive(Debug, PartialEq, Eq)]
144 enum TestEvent {
145 Event1,
146 Event2,
147 }
148
149 enum TestMessage {
150 Msg1,
151 Msg2,
152 }
153
154 impl Component for TestComponent {
155 type Error = ();
156 type Event = TestEvent;
157 type Init<'a> = Vec<TestMessage>;
158 type Message = TestMessage;
159
160 async fn init(init: Self::Init<'_>, sender: &ComponentSender<Self>) -> Result<Self, ()> {
161 for m in init {
162 sender.post(m);
163 }
164 Ok(Self)
165 }
166
167 async fn update(
168 &mut self,
169 message: Self::Message,
170 sender: &ComponentSender<Self>,
171 ) -> Result<bool, ()> {
172 match message {
173 TestMessage::Msg1 => {
174 sender.output(TestEvent::Event1);
175 Ok(false)
176 }
177 TestMessage::Msg2 => {
178 sender.output(TestEvent::Event2);
179 Ok(false)
180 }
181 }
182 }
183 }
184
185 async fn run_events<'a, T: Component>(
186 init: impl Into<T::Init<'a>>,
187 ) -> impl Stream<Item = RunEvent<T::Event, T::Error>> {
188 let mut root = Root::<T>::init(init)
190 .await
191 .expect("failed to init component");
192 stream! {
193 for await event in root.run() {
194 yield event;
195 }
196 }
197 }
198
199 async fn run_once<'a, T: Component>(
200 init: impl Into<T::Init<'a>>,
201 ) -> RunEvent<T::Event, T::Error> {
202 let stream = run_events::<T>(init).await;
203 let mut stream = std::pin::pin!(stream);
204 stream.next().await.expect("component exits without event")
205 }
206
207 #[compio::test]
208 async fn test_run() {
209 let event = run_once::<TestComponent>(vec![TestMessage::Msg1]).await;
210 assert_eq!(event, RunEvent::Event(TestEvent::Event1));
211
212 let event = run_once::<TestComponent>(vec![TestMessage::Msg2, TestMessage::Msg1]).await;
213 assert_eq!(event, RunEvent::Event(TestEvent::Event2));
214 }
215
216 #[compio::test]
217 async fn test_run_component() {
218 let events = run_events::<TestComponent>(vec![
219 TestMessage::Msg1,
220 TestMessage::Msg2,
221 TestMessage::Msg1,
222 ])
223 .await;
224 assert_send_sync(&events);
225 let expects = [TestEvent::Event1, TestEvent::Event2, TestEvent::Event1];
226 let zip = events.zip(futures_util::stream::iter(expects.into_iter()));
227 let mut zip = std::pin::pin!(zip);
228 while let Some((e, ex)) = zip.next().await {
229 assert_eq!(e, RunEvent::Event(ex));
230 }
231 }
232
233 fn assert_send_sync<T: Send + Sync>(_: &T) {}
234
235 #[compio::test]
236 async fn test_boxed_run() {
237 let mut boxed_child = Root::<TestComponent>::init(vec![
238 TestMessage::Msg1,
239 TestMessage::Msg2,
240 TestMessage::Msg1,
241 ])
242 .await
243 .expect("failed to init component")
244 .into_boxed();
245 let events = boxed_child.run();
246 let expects = [TestEvent::Event1, TestEvent::Event2, TestEvent::Event1];
247 let zip = events.zip(futures_util::stream::iter(expects.into_iter()));
248 let mut zip = std::pin::pin!(zip);
249 while let Some((e, ex)) = zip.next().await {
250 assert_eq!(e, RunEvent::Event(ex));
251 }
252 }
253}