1use crate::dm::{
28 AsyncHandler, AttrDetails, CmdDetails, HandlerContext, InvokeContextInstance,
29 InvokeReplyInstance, ReadContextInstance, ReadReplyInstance, WriteContextInstance,
30};
31use crate::error::{Error, ErrorCode};
32use crate::im::encoding::{AttrResp, AttrStatus, CmdResp, CmdStatus, IMStatusCode};
33use crate::tlv::{TLVElement, TLVWrite, TagType, ToTLV};
34use crate::transport::exchange::Exchange;
35use crate::utils::storage::WriteBuf;
36
37pub struct HandlerInvoker<'a, 'b, C> {
38 exchange: &'b mut Exchange<'a>,
39 context: C,
40}
41
42impl<'a, 'b, C> HandlerInvoker<'a, 'b, C>
43where
44 C: HandlerContext,
45{
46 pub const fn new(exchange: &'b mut Exchange<'a>, context: C) -> Self {
47 Self { exchange, context }
48 }
49
50 pub fn exchange(&mut self) -> &mut Exchange<'a> {
51 self.exchange
52 }
53
54 pub async fn process_read(
55 &mut self,
56 item: &Result<AttrDetails, AttrStatus>,
57 tw: &mut WriteBuf<'_>,
58 ) -> Result<(), Error> {
59 let tail = tw.get_tail();
60
61 let result = self.do_process_read(item, &mut *tw).await;
62
63 if result.is_err() {
64 tw.rewind_to(tail);
66 }
67
68 result
69 }
70
71 async fn do_process_read(
72 &mut self,
73 item: &Result<AttrDetails, AttrStatus>,
74 tw: &mut WriteBuf<'_>,
75 ) -> Result<(), Error> {
76 let result = match item {
77 Ok(attr) => {
78 let pos = tw.get_tail();
79
80 let result = self.read(attr, &mut *tw).await;
81
82 match result {
83 Ok(()) => Ok(None),
84 Err(e) if e.code() != ErrorCode::NoSpace => {
85 error!("Error reading attribute: {}", e);
86
87 tw.rewind_to(pos);
88
89 Ok(attr.status(e.into()))
90 }
91 Err(e) => Err(e),
92 }
93 }
94 Err(status) => {
95 error!("Error processing attribute read: {:?}", status);
96 Ok(Some(status.clone()))
97 }
98 };
99
100 match result {
101 Ok(Some(status)) => AttrResp::Status(status).to_tlv(&TagType::Anonymous, tw),
102 Ok(None) => Ok(()),
103 Err(err) => Err(err),
104 }
105 }
106
107 pub async fn read(&mut self, attr: &AttrDetails, tw: &mut WriteBuf<'_>) -> Result<(), Error> {
108 self.context
109 .handler()
110 .read(
111 ReadContextInstance::new(self.exchange, &self.context, attr),
112 ReadReplyInstance::new(attr, tw),
113 )
114 .await
115 }
116
117 pub async fn process_write(
118 &mut self,
119 item: &Result<(AttrDetails, TLVElement<'_>), AttrStatus>,
120 tw: &mut WriteBuf<'_>,
121 ) -> Result<(), Error> {
122 let tail = tw.get_tail();
123
124 let result = self.do_process_write(item, &mut *tw).await;
125
126 if result.is_err() {
127 tw.rewind_to(tail);
129 }
130
131 result
132 }
133
134 async fn do_process_write(
135 &mut self,
136 item: &Result<(AttrDetails, TLVElement<'_>), AttrStatus>,
137 tw: &mut WriteBuf<'_>,
138 ) -> Result<(), Error> {
139 let result = match item {
140 Ok((attr, data)) => {
141 let pos = tw.get_tail();
142
143 let result = self.write(attr, data).await;
144
145 match result {
146 Ok(()) => {
147 self.context.notify_attr_changed(
153 attr.endpoint_id,
154 attr.cluster_id,
155 attr.attr_id,
156 );
157 Ok(attr.status(IMStatusCode::Success))
158 }
159 Err(err) if err.code() != ErrorCode::NoSpace => {
160 error!("Error writing attribute: {}", err);
161
162 tw.rewind_to(pos);
163
164 Ok(attr.status(err.into()))
165 }
166 Err(err) => Err(err),
167 }
168 }
169 Err(status) => {
170 error!("Error processing attribute write: {:?}", status);
171 Ok(Some(status.clone()))
172 }
173 };
174
175 match result {
176 Ok(Some(status)) => status.to_tlv(&TagType::Anonymous, tw),
177 Ok(None) => Ok(()),
178 Err(err) => Err(err),
179 }
180 }
181
182 pub async fn write(&mut self, attr: &AttrDetails, data: &TLVElement<'_>) -> Result<(), Error> {
183 self.context
184 .handler()
185 .write(WriteContextInstance::new(
186 self.exchange,
187 &self.context,
188 attr,
189 data,
190 ))
191 .await
192 }
193
194 pub async fn process_invoke(
195 &mut self,
196 item: &Result<(CmdDetails, TLVElement<'_>), CmdStatus>,
197 tw: &mut WriteBuf<'_>,
198 ) -> Result<(), Error> {
199 let tail = tw.get_tail();
200
201 let result = self.do_process_invoke(item, &mut *tw).await;
202
203 if result.is_err() {
204 tw.rewind_to(tail);
206 }
207
208 result
209 }
210
211 async fn do_process_invoke(
212 &mut self,
213 item: &Result<(CmdDetails, TLVElement<'_>), CmdStatus>,
214 tw: &mut WriteBuf<'_>,
215 ) -> Result<(), Error> {
216 let result = match item {
217 Ok((cmd, data)) => {
218 let pos = tw.get_tail();
219
220 let result = self.invoke(cmd, data, &mut *tw).await;
221
222 match result {
223 Ok(()) => {
224 if pos == tw.get_tail() {
225 Ok(cmd.status(IMStatusCode::Success))
226 } else {
227 Ok(None)
228 }
229 }
230 Err(err) if err.code() != ErrorCode::NoSpace => {
231 error!("Error invoking command: {}", err);
232
233 tw.rewind_to(pos);
234
235 Ok(cmd.status(err.into()))
236 }
237 Err(err) => Err(err),
238 }
239 }
240 Err(status) => {
241 error!("Error processing command: {:?}", status);
242 Ok(Some(status.clone()))
243 }
244 };
245
246 match result {
247 Ok(Some(status)) => CmdResp::Status(status).to_tlv(&TagType::Anonymous, tw),
248 Ok(None) => Ok(()),
249 Err(err) => Err(err),
250 }
251 }
252
253 pub async fn invoke(
254 &mut self,
255 cmd: &CmdDetails,
256 data: &TLVElement<'_>,
257 tw: &mut WriteBuf<'_>,
258 ) -> Result<(), Error> {
259 self.context
260 .handler()
261 .invoke(
262 InvokeContextInstance::new(self.exchange, &self.context, cmd, data),
263 InvokeReplyInstance::new(cmd, tw),
264 )
265 .await
266 }
267}