1use core::fmt::Display;
19use core::future::Future;
20use core::pin::pin;
21
22use embassy_futures::select::{select, select_slice};
23
24use crate::crypto::Crypto;
25use crate::dm::clusters::net_comm;
26use crate::dm::networks::wireless::NoopWirelessNetCtl;
27use crate::dm::DataModel;
28use crate::error::Error;
29use crate::im::busy::BusyInteractionModel;
30use crate::im::events::DEFAULT_MAX_EVENTS_BUF_SIZE;
31use crate::im::subscriptions::DEFAULT_MAX_SUBSCRIPTIONS;
32use crate::im::{IMBuffer, InteractionModel, PROTO_ID_INTERACTION_MODEL};
33use crate::persist::KvBlobStoreAccess;
34use crate::sc::busy::BusySecureChannel;
35use crate::sc::SecureChannel;
36use crate::transport::exchange::Exchange;
37use crate::utils::select::Coalesce;
38use crate::utils::storage::pooled::Buffers;
39use crate::Matter;
40
41const RESPOND_BUSY_MS: u32 = 500;
44
45pub trait ExchangeHandler {
50 async fn handle(&self, exchange: Exchange<'_>) -> Result<(), Error>;
51}
52
53impl<T> ExchangeHandler for &T
54where
55 T: ExchangeHandler,
56{
57 fn handle(&self, exchange: Exchange<'_>) -> impl Future<Output = Result<(), Error>> {
58 (*self).handle(exchange)
59 }
60}
61
62pub struct ChainedExchangeHandler<H, T> {
66 pub handler_proto: u16,
67 pub handler: H,
68 pub next: T,
69}
70
71impl<H, T> ChainedExchangeHandler<H, T> {
72 pub const fn new(handler_proto: u16, handler: H, next: T) -> Self {
76 Self {
77 handler_proto,
78 handler,
79 next,
80 }
81 }
82
83 pub const fn chain<H2>(
89 self,
90 handler_proto: u16,
91 handler: H2,
92 ) -> ChainedExchangeHandler<H2, Self> {
93 ChainedExchangeHandler::new(handler_proto, handler, self)
94 }
95}
96
97impl<H, T> ExchangeHandler for ChainedExchangeHandler<H, T>
98where
99 H: ExchangeHandler,
100 T: ExchangeHandler,
101{
102 async fn handle(&self, mut exchange: Exchange<'_>) -> Result<(), Error> {
103 exchange.recv_fetch().await?;
106 let proto_id = exchange.rx()?.meta().proto_id;
107
108 if proto_id == self.handler_proto {
109 self.handler.handle(exchange).await
110 } else {
111 self.next.handle(exchange).await
112 }
113 }
114}
115
116pub struct EmptyExchangeHandler;
123
124impl ExchangeHandler for EmptyExchangeHandler {
125 async fn handle(&self, _exchange: Exchange<'_>) -> Result<(), Error> {
126 Ok(())
127 }
128}
129
130pub struct Responder<'a, T> {
135 name: &'a str,
136 handler: T,
137 matter: &'a Matter<'a>,
138 respond_after_ms: u32,
139}
140
141impl<'a, T> Responder<'a, T>
142where
143 T: ExchangeHandler,
144{
145 #[inline(always)]
153 pub const fn new(
154 name: &'a str,
155 handler: T,
156 matter: &'a Matter<'a>,
157 respond_after_ms: u32,
158 ) -> Self {
159 Self {
160 name,
161 handler,
162 matter,
163 respond_after_ms,
164 }
165 }
166
167 pub const fn name(&self) -> &str {
169 self.name
170 }
171
172 pub fn handler(&self) -> &T {
174 &self.handler
175 }
176
177 pub async fn run<const N: usize>(&self) -> Result<(), Error> {
179 info!("{}: Creating {} handlers", self.name, N);
180
181 let mut handlers = heapless::Vec::<_, N>::new();
182 debug!(
183 "{}: Handlers size: {}B",
184 self.name,
185 core::mem::size_of_val(&handlers)
186 );
187
188 for handler_id in 0..N {
189 unwrap!(handlers.push(self.handle(handler_id)).map_err(|_| ())); }
191
192 let handlers = pin!(handlers);
193 let handlers = unsafe { handlers.map_unchecked_mut(|handlers| handlers.as_mut_slice()) };
194
195 select_slice(handlers).await.0
196 }
197
198 #[inline(always)]
200 pub async fn handle(&self, handler_id: impl Display) -> Result<(), Error> {
201 loop {
202 let _ = self.respond_once(&handler_id).await;
204 }
205 }
206
207 #[inline(always)]
210 pub async fn respond_once(&self, handler_id: impl Display) -> Result<(), Error> {
211 let exchange = Exchange::accept_after(self.matter, self.respond_after_ms).await?;
212 let exchange_id = exchange.id();
214
215 if self.log_warn() {
216 warn!(
217 "{}: Handler {} / exchange {}: Starting",
218 self.name,
219 display2format!(&handler_id),
220 exchange_id
221 );
222 } else {
223 debug!(
224 "{}: Handler {} / exchange {}: Starting",
225 self.name,
226 display2format!(&handler_id),
227 exchange_id
228 );
229 }
230
231 let result = self.handler.handle(exchange).await;
232
233 if let Err(err) = &result {
234 error!(
235 "{}: Handler {} / exchange {}: Abandoned because of error {:?}",
236 self.name,
237 display2format!(&handler_id),
238 exchange_id,
239 err
240 );
241 } else if self.log_warn() {
242 warn!(
243 "{}: Handler {} / exchange {}: Completed",
244 self.name,
245 display2format!(&handler_id),
246 exchange_id
247 );
248 } else {
249 debug!(
250 "{}: Handler {} / exchange {}: Completed",
251 self.name,
252 display2format!(&handler_id),
253 exchange_id
254 );
255 }
256
257 result
258 }
259
260 fn log_warn(&self) -> bool {
261 self.respond_after_ms > 0
262 }
263}
264
265pub type DefaultExchangeHandler<
267 'd,
268 'a,
269 C,
270 B,
271 T,
272 K,
273 N,
274 NC = NoopWirelessNetCtl,
275 const NS: usize = DEFAULT_MAX_SUBSCRIPTIONS,
276 const NE: usize = DEFAULT_MAX_EVENTS_BUF_SIZE,
277> = ChainedExchangeHandler<
278 &'d InteractionModel<'a, C, B, T, K, N, NC, NS, NE>,
279 SecureChannel<'d, &'d C>,
280>;
281
282impl<'d, 'a, C, B, T, K, N, NC, const NS: usize, const NE: usize>
283 Responder<'a, DefaultExchangeHandler<'d, 'a, C, B, T, K, N, NC, NS, NE>>
284where
285 B: Buffers<IMBuffer>,
286{
287 #[inline(always)]
290 pub const fn new_default(
291 data_model: &'d InteractionModel<'a, C, B, T, K, N, NC, NS, NE>,
292 ) -> Self
293 where
294 C: Crypto,
295 T: DataModel,
296 K: KvBlobStoreAccess,
297 N: net_comm::Networks,
298 {
299 Self::new(
300 "Responder",
301 ChainedExchangeHandler::new(
302 PROTO_ID_INTERACTION_MODEL,
303 data_model,
304 SecureChannel::new(data_model.crypto(), data_model),
305 ),
306 data_model.matter(),
307 0,
308 )
309 }
310}
311
312pub type BusyExchangeHandler = ChainedExchangeHandler<BusyInteractionModel, BusySecureChannel>;
314
315impl<'a> Responder<'a, BusyExchangeHandler> {
316 #[inline(always)]
323 pub const fn new_busy(matter: &'a Matter<'a>, respond_after_ms: u32) -> Self {
324 Self::new(
325 "Busy Responder",
326 ChainedExchangeHandler::new(
327 PROTO_ID_INTERACTION_MODEL,
328 BusyInteractionModel::new(),
329 BusySecureChannel::new(),
330 ),
331 matter,
332 respond_after_ms,
333 )
334 }
335}
336
337pub struct DefaultResponder<
339 'd,
340 'a,
341 C,
342 B,
343 T,
344 K,
345 N,
346 NC = NoopWirelessNetCtl,
347 const NS: usize = DEFAULT_MAX_SUBSCRIPTIONS,
348 const NE: usize = DEFAULT_MAX_EVENTS_BUF_SIZE,
349> where
350 B: Buffers<IMBuffer>,
351{
352 responder: Responder<'a, DefaultExchangeHandler<'d, 'a, C, B, T, K, N, NC, NS, NE>>,
353 busy_responder: Responder<'a, BusyExchangeHandler>,
354}
355
356impl<'d, 'a, C, B, T, K, N, NC, const NS: usize, const NE: usize>
357 DefaultResponder<'d, 'a, C, B, T, K, N, NC, NS, NE>
358where
359 C: Crypto,
360 B: Buffers<IMBuffer>,
361 T: DataModel,
362 K: KvBlobStoreAccess,
363 N: net_comm::Networks,
364{
365 #[inline(always)]
367 pub const fn new(data_model: &'d InteractionModel<'a, C, B, T, K, N, NC, NS, NE>) -> Self {
368 Self {
369 responder: Responder::new_default(data_model),
370 busy_responder: Responder::new_busy(data_model.matter(), RESPOND_BUSY_MS),
371 }
372 }
373
374 pub async fn run<const A: usize, const O: usize>(&self) -> Result<(), Error> {
376 let mut actual = pin!(self.responder.run::<A>());
377 let mut busy = pin!(self.busy_responder.run::<O>());
378
379 select(&mut actual, &mut busy).coalesce().await
380 }
381
382 #[allow(clippy::type_complexity)]
386 pub const fn responder(
387 &self,
388 ) -> &Responder<
389 'a,
390 ChainedExchangeHandler<
391 &'d InteractionModel<'a, C, B, T, K, N, NC, NS, NE>,
392 SecureChannel<'d, &'d C>,
393 >,
394 > {
395 &self.responder
396 }
397
398 pub const fn busy_responder(
402 &self,
403 ) -> &Responder<'a, ChainedExchangeHandler<BusyInteractionModel, BusySecureChannel>> {
404 &self.busy_responder
405 }
406}