Debounce MQTT disconnects against stale_after; undo precautionary safe_state on recovery
This commit is contained in:
1 parent
7230e00e32
commit
9cc47d9a8f
1 file changed
+66
-1
+66
-1
@@ -118,6 +118,9 @@ class Engine:
|
|||||||
self.last_heartbeat = 0.0
|
self.last_heartbeat = 0.0
|
||||||
self.heartbeat_every = 300
|
self.heartbeat_every = 300
|
||||||
self.incomplete_reported = False
|
self.incomplete_reported = False
|
||||||
|
self.disconnected_since = None
|
||||||
|
self.stale_action_state = None
|
||||||
|
self.pre_stale_state = None
|
||||||
|
|
||||||
def evaluate(self, variables, now):
|
def evaluate(self, variables, now):
|
||||||
for state in self.states:
|
for state in self.states:
|
||||||
@@ -290,6 +293,34 @@ class Engine:
|
|||||||
self.last_heartbeat = now
|
self.last_heartbeat = now
|
||||||
self.report(self.summarize(variables, now))
|
self.report(self.summarize(variables, now))
|
||||||
|
|
||||||
|
async def undo_stale_action(self, session, now):
|
||||||
|
if self.stale_action_state is None:
|
||||||
|
return
|
||||||
|
|
||||||
|
previous = self.pre_stale_state
|
||||||
|
applied = self.stale_action_state
|
||||||
|
self.stale_action_state = None
|
||||||
|
self.pre_stale_state = None
|
||||||
|
|
||||||
|
if previous is None or previous == applied:
|
||||||
|
return
|
||||||
|
|
||||||
|
if self.last_commanded != applied:
|
||||||
|
return
|
||||||
|
|
||||||
|
self.report(
|
||||||
|
f"undoing the precautionary {self._word(applied)}, restoring "
|
||||||
|
f"{self._word(previous)} so the rules decide from here",
|
||||||
|
force=True,
|
||||||
|
)
|
||||||
|
await self.apply(
|
||||||
|
session,
|
||||||
|
previous,
|
||||||
|
"restored after telemetry recovered",
|
||||||
|
now,
|
||||||
|
force=True,
|
||||||
|
)
|
||||||
|
|
||||||
async def check_floor(self, session, variables, now):
|
async def check_floor(self, session, variables, now):
|
||||||
floor = self.profile.battery_floor
|
floor = self.profile.battery_floor
|
||||||
if floor is None or not floor.enabled:
|
if floor is None or not floor.enabled:
|
||||||
@@ -422,20 +453,50 @@ class Engine:
|
|||||||
return
|
return
|
||||||
|
|
||||||
if not self.source.connected():
|
if not self.source.connected():
|
||||||
|
if self.disconnected_since is None:
|
||||||
|
self.disconnected_since = now
|
||||||
|
self.report(
|
||||||
|
"MQTT session disconnected, waiting "
|
||||||
|
f"{format_duration(self.profile.stale_after)} to see if it "
|
||||||
|
"recovers before acting",
|
||||||
|
force=True,
|
||||||
|
)
|
||||||
|
|
||||||
|
waited = now - self.disconnected_since
|
||||||
|
if waited < (self.profile.stale_after or 0):
|
||||||
|
return
|
||||||
|
|
||||||
if not self.stale_reported:
|
if not self.stale_reported:
|
||||||
self.report("MQTT session disconnected", force=True)
|
self.report(
|
||||||
|
f"MQTT session still down after {format_duration(waited)}, "
|
||||||
|
f"applying on_stale={self.profile.on_stale}",
|
||||||
|
force=True,
|
||||||
|
)
|
||||||
self.stale_reported = True
|
self.stale_reported = True
|
||||||
await self.notify(
|
await self.notify(
|
||||||
None, None, {"reason": "mqtt disconnected"}, event="stale"
|
None, None, {"reason": "mqtt disconnected"}, event="stale"
|
||||||
)
|
)
|
||||||
|
|
||||||
if self.profile.on_stale == "stop":
|
if self.profile.on_stale == "stop":
|
||||||
raise RuntimeError("stopping: MQTT session disconnected")
|
raise RuntimeError("stopping: MQTT session disconnected")
|
||||||
if self.profile.on_stale == "safe_state":
|
if self.profile.on_stale == "safe_state":
|
||||||
|
if self.stale_action_state is None:
|
||||||
|
self.pre_stale_state = self.last_commanded
|
||||||
|
self.stale_action_state = self.profile.safe_state
|
||||||
await self.apply(
|
await self.apply(
|
||||||
session, self.profile.safe_state, "mqtt disconnected safe state", now
|
session, self.profile.safe_state, "mqtt disconnected safe state", now
|
||||||
)
|
)
|
||||||
return
|
return
|
||||||
|
|
||||||
|
if self.disconnected_since is not None:
|
||||||
|
self.report(
|
||||||
|
"MQTT session recovered after "
|
||||||
|
f"{format_duration(now - self.disconnected_since)}",
|
||||||
|
force=True,
|
||||||
|
)
|
||||||
|
self.disconnected_since = None
|
||||||
|
await self.undo_stale_action(session, now)
|
||||||
|
|
||||||
if age is not None and self.profile.stale_after and age > self.profile.stale_after:
|
if age is not None and self.profile.stale_after and age > self.profile.stale_after:
|
||||||
if not self.stale_reported:
|
if not self.stale_reported:
|
||||||
self.report(
|
self.report(
|
||||||
@@ -450,6 +511,9 @@ class Engine:
|
|||||||
if self.profile.on_stale == "stop":
|
if self.profile.on_stale == "stop":
|
||||||
raise RuntimeError("stopping: telemetry went stale")
|
raise RuntimeError("stopping: telemetry went stale")
|
||||||
if self.profile.on_stale == "safe_state":
|
if self.profile.on_stale == "safe_state":
|
||||||
|
if self.stale_action_state is None:
|
||||||
|
self.pre_stale_state = self.last_commanded
|
||||||
|
self.stale_action_state = self.profile.safe_state
|
||||||
await self.apply(
|
await self.apply(
|
||||||
session, self.profile.safe_state, "stale telemetry safe state", now
|
session, self.profile.safe_state, "stale telemetry safe state", now
|
||||||
)
|
)
|
||||||
@@ -458,6 +522,7 @@ class Engine:
|
|||||||
if self.stale_reported:
|
if self.stale_reported:
|
||||||
self.report("telemetry recovered", force=True)
|
self.report("telemetry recovered", force=True)
|
||||||
self.stale_reported = False
|
self.stale_reported = False
|
||||||
|
await self.undo_stale_action(session, now)
|
||||||
|
|
||||||
variables = derived_values(self.anker_profile, status)
|
variables = derived_values(self.anker_profile, status)
|
||||||
|
|
||||||
|
|||||||
Reference in new issue
Block a user