This commit is contained in:
Mateusz Gruszczyński
2026-09-14 23:10:17 +02:00
parent 4bb9c8621a
commit 0a2fb8c6c7
47 changed files with 1337 additions and 372 deletions
+57 -3
View File
@@ -267,7 +267,7 @@ impl GreeCloudProvider {
}
pub async fn unregister_device(&self, device_id: &str) {
let cloud_ids = {
let (cloud_ids, no_registered_devices) = {
let mut registered = self.inner.registered.write().await;
let ids = registered
.iter()
@@ -275,7 +275,8 @@ impl GreeCloudProvider {
.map(|(cloud_id, _)| cloud_id.clone())
.collect::<Vec<_>>();
registered.retain(|_, item| item.device_id != device_id);
ids
let empty = registered.is_empty();
(ids, empty)
};
let mut pending = self.inner.pending.lock().await;
for cloud_id in cloud_ids {
@@ -284,6 +285,18 @@ impl GreeCloudProvider {
drop(pending);
self.inner.command_locks.lock().await.remove(device_id);
self.inner.diagnostics.write().await.remove(device_id);
// The broker can continue publishing retained/status traffic for topics from the
// previous subscription until the MQTT session is closed. Once the last registered
// Cloud device is removed there is nothing useful to receive, so close the session
// immediately instead of leaving a stale subscription alive until process restart.
if no_registered_devices {
let session = self.inner.session.write().await.take();
if let Some(session) = session {
session.mqtt.disconnect().await;
tracing::info!("GREE Cloud MQTT disconnected; no registered devices remain");
}
}
}
pub async fn poll(
@@ -1115,6 +1128,14 @@ impl GreeCloudProvider {
"topic": topic,
"bytes": payload.len(),
}));
// A broker message can race with device deletion. If the last Cloud device has
// already been unregistered, any in-flight message belongs to a stale subscription
// and must be ignored rather than reported as a payload error every few seconds.
let registered = self.inner.registered.read().await.clone();
if registered.is_empty() {
tracing::debug!(topic, "ignoring GREE Cloud MQTT payload without registered devices");
return Ok(());
}
if topic.starts_with("connect/") {
let parent = topic.split('/').nth(1).unwrap_or_default().to_string();
self.emit_cloud_mqtt(json!({
@@ -1143,7 +1164,6 @@ impl GreeCloudProvider {
"cipher": if envelope.tag.is_some() { 2 } else { 1 },
}));
let parent_from_topic = topic.split('/').nth(1).unwrap_or_default().to_string();
let registered = self.inner.registered.read().await.clone();
// Match the reference client primarily by subscribed parent topic. Some GREE
// responses use a parent/variant tcid that is not byte-for-byte equal to the
// discovery child id; restricting candidates to tcid caused valid frames to be
@@ -1752,6 +1772,40 @@ mod tests {
assert!(!LEGACY_CLOUD_PROPERTIES.contains(&"CompressorFqy"));
}
#[tokio::test]
async fn mqtt_payload_is_ignored_after_all_cloud_devices_are_removed() {
let provider = GreeCloudProvider::new(reqwest::Client::new());
let result = provider
.handle_raw_message("status/stale-parent/device", b"not-json-anymore")
.await;
assert!(result.is_ok());
}
#[tokio::test]
async fn mqtt_payload_for_unknown_parent_still_errors_when_devices_are_registered() {
let provider = GreeCloudProvider::new(reqwest::Client::new());
provider.inner.registered.write().await.insert(
"AABBCCDDEEFF".into(),
RegisteredDevice {
device_id: "cloud-test".into(),
key: "0123456789abcdef".into(),
parent_mac: "AABBCCDDEE".into(),
cipher_version: 1,
},
);
let err = provider
.handle_raw_message(
"status/1122334455/device",
br#"{"pack":"x","tcid":"112233445566"}"#,
)
.await
.unwrap_err();
assert!(err
.to_string()
.contains("does not match a registered device"));
}
#[tokio::test]
async fn provider_dispatcher_selects_transport_from_connection_type() {
let local = GreeClient::new(