Commit eb4c9e5f authored by xuwencheng's avatar xuwencheng

upgrade

parent b17996b1
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use serde_repr::{Deserialize_repr, Serialize_repr}; use serde_repr::{Deserialize_repr, Serialize_repr};
use crate::event::Status;
#[derive(Debug, Deserialize_repr, Serialize_repr, Eq, PartialEq, Clone, Hash)] #[derive(Debug, Deserialize_repr, Serialize_repr, Eq, PartialEq, Clone, Hash)]
#[repr(u8)] #[repr(u8)]
pub enum EventType { pub enum EventType {
// 业务推送事件 // pos推送客户端消息入口
BusinessEvent = 0, ReceivePush = 0,
// 推送事件 // 推送事件
PublishEvent = 10, Publish = 10,
// 连接事件 // 连接事件
ConnectEvent = 20, Connect = 20,
CleanSessionEvent = 23,
// 消费事件 // 消费事件
ConsumeEvent = 30, Consume = 30,
Clean = 31,
} }
#[derive(Debug, Clone, Deserialize, Serialize, Default)] #[derive(Debug, Clone, Deserialize, Serialize, Default)]
...@@ -46,3 +48,13 @@ pub struct SocketEvent { ...@@ -46,3 +48,13 @@ pub struct SocketEvent {
// 扩展信息 // 扩展信息
pub extra: Option<String>, pub extra: Option<String>,
} }
impl SocketEvent {
pub fn new(desc: impl AsRef<str>, ty: EventType, status: Status) -> Self {
let mut event = SocketEvent::default();
event.event = desc.as_ref().to_string();
event.event_ty = ty as i32;
event.status = status as i32;
event
}
}
...@@ -25,21 +25,14 @@ impl BusinessEvent for SocketEvent { ...@@ -25,21 +25,14 @@ impl BusinessEvent for SocketEvent {
extra: impl AsRef<str>, extra: impl AsRef<str>,
) -> SocketEvent { ) -> SocketEvent {
let Header(msg_id, msg_ty, qos) = header; let Header(msg_id, msg_ty, qos) = header;
SocketEvent { let mut event = SocketEvent::new("业务推送事件", EventType::ReceivePush, status);
device_id: None, event.key = key.map(|v| v.as_ref().to_string());
client_type: None, event.qos = Some(qos as i32);
version: None, event.msg_id = Some(msg_id);
version_desc: None, event.msg_ty = Some(msg_ty);
event: "业务推送事件".to_string(), event.desc = qos.desc();
event_ty: EventType::BusinessEvent as i32, event.extra = Some(extra.as_ref().to_string());
key: key.map(|v| v.as_ref().to_string()), event
qos: Some(qos as i32),
msg_id: Some(msg_id),
msg_ty: Some(msg_ty),
status: status as i32,
desc: qos.desc(),
extra: Some(extra.as_ref().to_string()),
}
} }
} }
...@@ -58,46 +51,32 @@ impl PublishEvent for SocketEvent { ...@@ -58,46 +51,32 @@ impl PublishEvent for SocketEvent {
desc: impl AsRef<str>, desc: impl AsRef<str>,
) -> SocketEvent { ) -> SocketEvent {
let Header(msg_id, msg_ty, qos) = header; let Header(msg_id, msg_ty, qos) = header;
SocketEvent { let mut event = SocketEvent::new("推送事件", EventType::Publish, status);
device_id: Some(info.device_id.clone()), event.device_id = Some(info.device_id.clone());
client_type: Some(info.client_type), event.client_type = Some(info.client_type);
version: Some(info.version), event.version = Some(info.version);
version_desc: info.version_desc.clone(), event.version_desc = info.version_desc.clone();
event: "推送事件".to_string(), event.qos = Some(qos as i32);
event_ty: EventType::PublishEvent as i32, event.msg_id = Some(msg_id);
key: None, event.msg_ty = Some(msg_ty);
qos: Some(qos as i32), event.desc = desc.as_ref().to_string();
msg_id: Some(msg_id), event
msg_ty: Some(msg_ty),
status: status as i32,
desc: desc.as_ref().to_string(),
extra: None,
}
} }
fn ack(info: &TokenInfo, msg_id: i64, desc: impl AsRef<str>) -> SocketEvent { fn ack(info: &TokenInfo, msg_id: i64, desc: impl AsRef<str>) -> SocketEvent {
SocketEvent { let mut event = SocketEvent::new("推送事件", EventType::Publish, Status::Success);
device_id: Some(info.device_id.clone()), event.device_id = Some(info.device_id.clone());
client_type: Some(info.client_type), event.client_type = Some(info.client_type);
version: Some(info.version), event.version = Some(info.version);
version_desc: info.version_desc.clone(), event.version_desc = info.version_desc.clone();
event: "推送事件".to_string(), event.msg_id = Some(msg_id);
event_ty: EventType::PublishEvent as i32, event.desc = desc.as_ref().to_string();
key: None, event
qos: None,
msg_id: Some(msg_id),
msg_ty: None,
status: Status::Success as i32,
desc: desc.as_ref().to_string(),
extra: None,
}
} }
} }
pub trait ConnectEvent { pub trait ConnectEvent {
fn connect(info: &TokenInfo, status: Status, extra: Option<impl AsRef<str>>) -> SocketEvent; fn connect(info: &TokenInfo, status: Status, extra: Option<impl AsRef<str>>) -> SocketEvent;
fn clean_session(info: &TokenInfo, extra: impl AsRef<str>) -> SocketEvent;
} }
impl ConnectEvent for SocketEvent { impl ConnectEvent for SocketEvent {
...@@ -107,39 +86,14 @@ impl ConnectEvent for SocketEvent { ...@@ -107,39 +86,14 @@ impl ConnectEvent for SocketEvent {
} else { } else {
"设备断开连接" "设备断开连接"
}; };
SocketEvent { let mut event = SocketEvent::new("连接事件", EventType::Connect, status);
device_id: Some(info.device_id.clone()), event.device_id = Some(info.device_id.clone());
client_type: Some(info.client_type), event.client_type = Some(info.client_type);
version: Some(info.version), event.version = Some(info.version);
version_desc: info.version_desc.clone(), event.version_desc = info.version_desc.clone();
event: "连接事件".to_string(), event.desc = desc.to_string();
event_ty: EventType::ConnectEvent as i32, event.extra = extra.map(|s| s.as_ref().to_string());
key: None, event
qos: None,
msg_id: None,
msg_ty: None,
status: status as i32,
desc: desc.to_string(),
extra: extra.map(|s| s.as_ref().to_string()),
}
}
fn clean_session(info: &TokenInfo, extra: impl AsRef<str>) -> SocketEvent {
SocketEvent {
device_id: Some(info.device_id.clone()),
client_type: Some(info.client_type),
version: Some(info.version),
version_desc: info.version_desc.clone(),
event: "连接事件".to_string(),
event_ty: EventType::CleanSessionEvent as i32,
key: None,
qos: None,
msg_id: None,
msg_ty: None,
status: Status::Fail as i32,
desc: "清理历史消息".to_string(),
extra: Some(extra.as_ref().to_string()),
}
} }
} }
...@@ -150,6 +104,8 @@ pub trait ConsumeEvent { ...@@ -150,6 +104,8 @@ pub trait ConsumeEvent {
status: Status, status: Status,
extra: Option<impl AsRef<str>>, extra: Option<impl AsRef<str>>,
) -> SocketEvent; ) -> SocketEvent;
fn clean_session(info: &TokenInfo, extra: impl AsRef<str>) -> SocketEvent;
} }
impl ConsumeEvent for SocketEvent { impl ConsumeEvent for SocketEvent {
...@@ -165,20 +121,27 @@ impl ConsumeEvent for SocketEvent { ...@@ -165,20 +121,27 @@ impl ConsumeEvent for SocketEvent {
"消息消费失败" "消息消费失败"
}; };
let Header(msg_id, msg_ty, qos) = header; let Header(msg_id, msg_ty, qos) = header;
SocketEvent { let mut event = SocketEvent::new("消费事件", EventType::Consume, Status::Fail);
device_id: Some(info.device_id.clone()), event.device_id = Some(info.device_id.clone());
client_type: Some(info.client_type), event.client_type = Some(info.client_type);
version: Some(info.version), event.version = Some(info.version);
version_desc: info.version_desc.clone(), event.version_desc = info.version_desc.clone();
event: "消费事件".to_string(), event.qos = Some(qos as i32);
event_ty: EventType::ConsumeEvent as i32, event.msg_id = Some(msg_id);
key: None, event.msg_ty = Some(msg_ty);
qos: Some(qos as i32), event.desc = desc.to_string();
msg_id: Some(msg_id), event.extra = extra.map(|s| s.as_ref().to_string());
msg_ty: Some(msg_ty), event
status: status as i32, }
desc: desc.to_string(),
extra: extra.map(|s| s.as_ref().to_string()), fn clean_session(info: &TokenInfo, extra: impl AsRef<str>) -> SocketEvent {
} let mut event = SocketEvent::new("消费事件", EventType::Clean, Status::Fail);
event.device_id = Some(info.device_id.clone());
event.client_type = Some(info.client_type);
event.version = Some(info.version);
event.version_desc = info.version_desc.clone();
event.desc = "清理历史消息".to_string();
event.extra = Some(extra.as_ref().to_string());
event
} }
} }
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment