|
|
|
|
@ -189,7 +189,7 @@ impl MQTTState {
|
|
|
|
|
self.protocol_version = conn.protocol_version;
|
|
|
|
|
if self.connected {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::DoubleConnect);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::DoubleConnect);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
} else {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
@ -200,7 +200,7 @@ impl MQTTState {
|
|
|
|
|
MQTTOperation::PUBLISH(ref publish) => {
|
|
|
|
|
if !self.connected {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::UnintroducedMessage);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::UnintroducedMessage);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
@ -219,13 +219,13 @@ impl MQTTState {
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
} else {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::MissingMsgId);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::MissingMsgId);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
}
|
|
|
|
|
},
|
|
|
|
|
_ => {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::InvalidQosLevel);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::InvalidQosLevel);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
@ -233,7 +233,7 @@ impl MQTTState {
|
|
|
|
|
MQTTOperation::SUBSCRIBE(ref subscribe) => {
|
|
|
|
|
if !self.connected {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::UnintroducedMessage);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::UnintroducedMessage);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
@ -253,7 +253,7 @@ impl MQTTState {
|
|
|
|
|
},
|
|
|
|
|
_ => {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::InvalidQosLevel);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::InvalidQosLevel);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
@ -261,7 +261,7 @@ impl MQTTState {
|
|
|
|
|
MQTTOperation::UNSUBSCRIBE(ref unsubscribe) => {
|
|
|
|
|
if !self.connected {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::UnintroducedMessage);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::UnintroducedMessage);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
@ -281,7 +281,7 @@ impl MQTTState {
|
|
|
|
|
},
|
|
|
|
|
_ => {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::InvalidQosLevel);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::InvalidQosLevel);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
@ -294,7 +294,7 @@ impl MQTTState {
|
|
|
|
|
self.connected = true;
|
|
|
|
|
} else {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::MissingConnect);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::MissingConnect);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
}
|
|
|
|
|
},
|
|
|
|
|
@ -302,7 +302,7 @@ impl MQTTState {
|
|
|
|
|
| MQTTOperation::PUBREL(ref v) => {
|
|
|
|
|
if !self.connected {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::UnintroducedMessage);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::UnintroducedMessage);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
@ -310,7 +310,7 @@ impl MQTTState {
|
|
|
|
|
(*tx).msg.push(msg);
|
|
|
|
|
} else {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::MissingPublish);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::MissingPublish);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
}
|
|
|
|
|
},
|
|
|
|
|
@ -318,7 +318,7 @@ impl MQTTState {
|
|
|
|
|
| MQTTOperation::PUBCOMP(ref v) => {
|
|
|
|
|
if !self.connected {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::UnintroducedMessage);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::UnintroducedMessage);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
@ -328,14 +328,14 @@ impl MQTTState {
|
|
|
|
|
(*tx).pkt_id = None;
|
|
|
|
|
} else {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::MissingPublish);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::MissingPublish);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
}
|
|
|
|
|
},
|
|
|
|
|
MQTTOperation::SUBACK(ref suback) => {
|
|
|
|
|
if !self.connected {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::UnintroducedMessage);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::UnintroducedMessage);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
@ -345,14 +345,14 @@ impl MQTTState {
|
|
|
|
|
(*tx).pkt_id = None;
|
|
|
|
|
} else {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::MissingSubscribe);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::MissingSubscribe);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
}
|
|
|
|
|
},
|
|
|
|
|
MQTTOperation::UNSUBACK(ref unsuback) => {
|
|
|
|
|
if !self.connected {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::UnintroducedMessage);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::UnintroducedMessage);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
@ -362,14 +362,14 @@ impl MQTTState {
|
|
|
|
|
(*tx).pkt_id = None;
|
|
|
|
|
} else {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::MissingUnsubscribe);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::MissingUnsubscribe);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
}
|
|
|
|
|
},
|
|
|
|
|
MQTTOperation::UNASSIGNED => {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
tx.complete = true;
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::UnassignedMsgtype);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::UnassignedMsgtype);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
},
|
|
|
|
|
MQTTOperation::TRUNCATED(_) => {
|
|
|
|
|
@ -381,7 +381,7 @@ impl MQTTState {
|
|
|
|
|
| MQTTOperation::DISCONNECT(_) => {
|
|
|
|
|
if !self.connected {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::UnintroducedMessage);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::UnintroducedMessage);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
@ -393,7 +393,7 @@ impl MQTTState {
|
|
|
|
|
| MQTTOperation::PINGRESP => {
|
|
|
|
|
if !self.connected {
|
|
|
|
|
let mut tx = self.new_tx(msg, toclient);
|
|
|
|
|
&mut MQTTState::set_event(&mut tx, MQTTEvent::UnintroducedMessage);
|
|
|
|
|
MQTTState::set_event(&mut tx, MQTTEvent::UnintroducedMessage);
|
|
|
|
|
self.transactions.push(tx);
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
|