@@ -57,6 +57,26 @@ def __init__(
5757 def publisher (self ) -> Publisher :
5858 return self .vehicle_state .publisher
5959
60+ def __publish_state_echo (self , * , topic : str , value : object | None ) -> None :
61+ if value is None :
62+ return
63+ full_topic = self .vehicle_state .get_topic (topic )
64+ try :
65+ if isinstance (value , bool ):
66+ self .publisher .publish_bool (full_topic , value )
67+ elif isinstance (value , int ):
68+ self .publisher .publish_int (full_topic , value )
69+ elif isinstance (value , float ):
70+ self .publisher .publish_float (full_topic , value )
71+ else :
72+ self .publisher .publish_str (full_topic , str (value ))
73+ except Exception :
74+ LOG .warning (
75+ "Failed to publish state echo for topic %s" ,
76+ full_topic ,
77+ exc_info = True ,
78+ )
79+
6080 def __report_command_failure (
6181 self ,
6282 * ,
@@ -106,6 +126,23 @@ async def handle_mqtt_command(self, *, topic: str, payload: str) -> None:
106126 handler = handler , payload = payload , analyzed_topic = analyzed_topic
107127 )
108128
129+ async def __run_handler_and_report_success (
130+ self ,
131+ * ,
132+ handler : CommandHandlerBase ,
133+ payload : str ,
134+ analyzed_topic : _MqttCommandTopic ,
135+ ) -> None :
136+ execution_result = await handler .handle (payload )
137+ self .publisher .publish_str (analyzed_topic .response_no_global , "Success" )
138+ if execution_result .force_refresh :
139+ self .vehicle_state .set_refresh_mode (
140+ RefreshMode .FORCE ,
141+ f"after command execution on topic { analyzed_topic .command_no_vin } " ,
142+ )
143+ if execution_result .clear_command :
144+ self .publisher .clear_topic (analyzed_topic .command_no_global )
145+
109146 async def __execute_mqtt_command_handler (
110147 self ,
111148 * ,
@@ -114,65 +151,87 @@ async def __execute_mqtt_command_handler(
114151 analyzed_topic : _MqttCommandTopic ,
115152 ) -> None :
116153 topic = analyzed_topic .command_no_vin
117- topic_no_global = analyzed_topic .command_no_global
118154 result_topic = analyzed_topic .response_no_global
119155
156+ echo_topic = handler .echo_state_topic ()
157+ prior_state : object | None = None
158+ echoed = False
159+ if echo_topic is not None :
160+ prior_state = handler .capture_current_state ()
161+ echoed_value = handler .echo_payload (payload )
162+ if echoed_value is not None :
163+ self .__publish_state_echo (topic = echo_topic , value = echoed_value )
164+ echoed = True
165+
166+ def rollback_state_echo () -> None :
167+ if echoed and echo_topic is not None and prior_state is not None :
168+ self .__publish_state_echo (topic = echo_topic , value = prior_state )
169+
120170 try :
121- execution_result = await handler .handle (payload )
122- self .publisher .publish_str (result_topic , "Success" )
123- if execution_result .force_refresh :
124- self .vehicle_state .set_refresh_mode (
125- RefreshMode .FORCE , f"after command execution on topic { topic } "
126- )
127- if execution_result .clear_command :
128- self .publisher .clear_topic (topic_no_global )
171+ await self .__run_handler_and_report_success (
172+ handler = handler , payload = payload , analyzed_topic = analyzed_topic
173+ )
129174 except MqttGatewayException as e :
175+ rollback_state_echo ()
130176 self .__report_command_failure (
131177 command = topic , result_topic = result_topic , detail = e .message , exc = e
132178 )
133179 except SaicLogoutException :
134- LOG .warning (
135- "API Client was logged out, attempting immediate relogin and retry"
180+ await self .__handle_logout_and_retry (
181+ handler = handler ,
182+ payload = payload ,
183+ analyzed_topic = analyzed_topic ,
184+ rollback_state_echo = rollback_state_echo ,
136185 )
137- try :
138- await self .relogin_handler .force_login ()
139- except Exception as login_err :
140- self .__report_command_failure (
141- command = topic ,
142- result_topic = result_topic ,
143- detail = f"relogin failed ({ login_err } )" ,
144- exc = login_err ,
145- )
146- return
147- try :
148- execution_result = await handler .handle (payload )
149- self .publisher .publish_str (result_topic , "Success" )
150- if execution_result .force_refresh :
151- self .vehicle_state .set_refresh_mode (
152- RefreshMode .FORCE ,
153- f"after command execution on topic { topic } " ,
154- )
155- if execution_result .clear_command :
156- self .publisher .clear_topic (topic_no_global )
157- except Exception as retry_err :
158- self .__report_command_failure (
159- command = topic ,
160- result_topic = result_topic ,
161- detail = str (retry_err ),
162- exc = retry_err ,
163- )
164186 except SaicApiException as se :
187+ rollback_state_echo ()
165188 self .__report_command_failure (
166189 command = topic , result_topic = result_topic , detail = se .message , exc = se
167190 )
168191 except Exception as e :
192+ rollback_state_echo ()
169193 self .__report_command_failure (
170194 command = topic ,
171195 result_topic = result_topic ,
172196 detail = "unexpected error" ,
173197 exc = e ,
174198 )
175199
200+ async def __handle_logout_and_retry (
201+ self ,
202+ * ,
203+ handler : CommandHandlerBase ,
204+ payload : str ,
205+ analyzed_topic : _MqttCommandTopic ,
206+ rollback_state_echo : Callable [[], None ],
207+ ) -> None :
208+ topic = analyzed_topic .command_no_vin
209+ result_topic = analyzed_topic .response_no_global
210+ LOG .warning ("API Client was logged out, attempting immediate relogin and retry" )
211+ try :
212+ await self .relogin_handler .force_login ()
213+ except Exception as login_err :
214+ rollback_state_echo ()
215+ self .__report_command_failure (
216+ command = topic ,
217+ result_topic = result_topic ,
218+ detail = f"relogin failed ({ login_err } )" ,
219+ exc = login_err ,
220+ )
221+ return
222+ try :
223+ await self .__run_handler_and_report_success (
224+ handler = handler , payload = payload , analyzed_topic = analyzed_topic
225+ )
226+ except Exception as retry_err :
227+ rollback_state_echo ()
228+ self .__report_command_failure (
229+ command = topic ,
230+ result_topic = result_topic ,
231+ detail = str (retry_err ),
232+ exc = retry_err ,
233+ )
234+
176235 def __get_command_topics (self , topic : str ) -> _MqttCommandTopic :
177236 global_topic_removed = topic .removeprefix (self .global_mqtt_topic ).removeprefix (
178237 "/"
0 commit comments