Commit e3feb7dc authored by xuwencheng's avatar xuwencheng

event

parent 0b9066fe
...@@ -5,9 +5,10 @@ edition = "2021" ...@@ -5,9 +5,10 @@ edition = "2021"
[features] [features]
default = [] default = []
socket = ["rmp", "rmp-serde", "bytes", "flate2", "qm_rust_util"] socket = ["rmp", "rmp-serde", "bytes", "flate2", "event"]
paidui = ["qm_rust_util"] paidui = []
login = [] login = []
event = []
[dependencies] [dependencies]
anyhow = { version = "1", features = ["std"], default-features = false } anyhow = { version = "1", features = ["std"], default-features = false }
...@@ -15,11 +16,11 @@ serde_json = { version = "1", default-features = false } ...@@ -15,11 +16,11 @@ serde_json = { version = "1", default-features = false }
serde_repr = { version = "0.1", default-features = false } serde_repr = { version = "0.1", default-features = false }
serde = { version = "1", features = ["derive"], default-features = false } serde = { version = "1", features = ["derive"], default-features = false }
chrono = "0.4" chrono = "0.4"
chrono-tz = "0.9.0"
tracing = { version = "0.1", default-features = false } tracing = { version = "0.1", default-features = false }
bytes = { version = "1", default-features = false, optional = true } bytes = { version = "1", default-features = false, optional = true }
rmp = { version = "0.8", default-features = false, optional = true } rmp = { version = "0.8", default-features = false, optional = true }
rmp-serde = { version = "1", default-features = false, optional = true } rmp-serde = { version = "1", default-features = false, optional = true }
flate2 = { version = "1", features = ["rust_backend"], default-features = false, optional = true } flate2 = { version = "1", features = ["rust_backend"], default-features = false, optional = true }
qm_rust_util = { git = "https://xuwencheng:password2024@git.zmcms.cn/rust-common/qm_rust_util.git", default-features = false, optional = true }
#[cfg(feature = "socket")]
pub use socket::*;
#[cfg(feature = "socket")]
pub mod socket;
\ No newline at end of file
use serde::{Deserialize, Serialize};
use serde_repr::{Deserialize_repr, Serialize_repr};
#[derive(Debug, Deserialize_repr, Serialize_repr, Eq, PartialEq, Clone, Hash)]
#[repr(u8)]
pub enum EventType {
// 业务推送事件
BusinessEvent = 0,
// 推送事件
PublishEvent = 10,
// 连接事件
ConnectEvent = 20,
}
#[derive(Debug, Clone, Deserialize, Serialize, Default)]
pub struct SocketEvent {
// 事件
pub event: String,
// 事件类型
pub event_ty: i32,
// 关键字
pub key: Option<String>,
// 消息id
pub msg_id: Option<i64>,
// 消息类型
pub msg_ty: Option<i32>,
// 事件状态
pub status: i32,
// 事件描述
pub desc: String,
pub seller_id: Option<i64>,
pub store_id: Option<i64>,
pub device_id: Option<String>,
pub client_type: Option<i32>,
pub version: Option<i32>,
pub version_desc: Option<String>,
// 扩展信息
pub extra: Option<String>,
}
use std::collections::hash_map::DefaultHasher;
use std::{
any::TypeId,
hash::{Hash, Hasher},
};
use once_cell::sync::Lazy;
use snowflake::SnowflakeIdGenerator;
use uuid::Uuid;
pub fn uuid() -> String {
Uuid::new_v4().as_simple().to_string()
}
pub fn id() -> i64 {
static mut GENERATOR: Lazy<SnowflakeIdGenerator> =
Lazy::new(|| SnowflakeIdGenerator::new(1, 1));
unsafe { GENERATOR.generate() }
}
// 生成指定日期的雪花ID函数
pub fn generate_id(time: u64) -> u64 {
// 默认: 2024-06-21 18:25:45
let last_time_millis = time;
let timestamp_part = last_time_millis << 22;
let node_id_part = (1 as u64) << 17;
let worker_id_part = (1 as u64) << 12;
let sequence_part = 0 as u64;
timestamp_part | node_id_part | worker_id_part | sequence_part
}
pub fn is_buf<T: 'static>() -> bool {
TypeId::of::<T>() == TypeId::of::<Vec<u8>>()
}
pub fn short_link_hash(input: &str) -> String {
let mut hasher = DefaultHasher::new();
input.hash(&mut hasher);
let hash = hasher.finish();
base62_encode(hash)
}
fn base62_encode(mut num: u64) -> String {
const CHARSET: &[u8] = b"0123456789abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ";
let mut encoded = Vec::new();
while num > 0 {
let rem = (num % 62) as usize;
encoded.push(CHARSET[rem]);
num /= 62;
}
encoded.reverse();
String::from_utf8(encoded).unwrap()
}
#[cfg(feature = "event")]
pub mod event;
pub mod id;
#[cfg(feature = "login")]
pub mod login;
#[cfg(feature = "paidui")] #[cfg(feature = "paidui")]
pub mod paidui; pub mod paidui;
pub mod serde;
#[cfg(feature = "socket")] #[cfg(feature = "socket")]
pub mod socket; pub mod socket;
#[cfg(feature = "login")]
pub mod login;
#[cfg(test)] #[cfg(test)]
mod tests { mod tests {
......
use chrono::{DateTime, Local}; use chrono::{DateTime, Local};
use qm_rust_util::serde::serde_date; use crate::serde::serde_date;
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use serde_json::Value; use serde_json::Value;
......
use anyhow::Result;
use serde::Deserialize;
use serde_json::Value;
pub mod string_or_i32 {
use serde::{de, Deserialize, Deserializer, Serializer};
use std::fmt::Display;
pub fn serialize<T, S>(value: &T, serializer: S) -> Result<S::Ok, S::Error>
where
T: Display,
S: Serializer,
{
serializer.collect_str(value)
}
pub fn deserialize<'de, D>(deserializer: D) -> Result<i32, D::Error>
where
D: Deserializer<'de>,
{
#[derive(Deserialize)]
#[serde(untagged)]
enum StringOrI32 {
String(String),
I32(i32),
}
match StringOrI32::deserialize(deserializer)? {
StringOrI32::String(s) => s.parse().map_err(de::Error::custom),
StringOrI32::I32(i) => Ok(i),
}
}
}
pub mod serde_date {
use chrono::{DateTime, Local, NaiveDateTime, TimeZone};
use chrono_tz::Asia::Shanghai;
use serde::{self, Deserialize, Deserializer, Serializer};
const FORMAT: &str = "%Y-%m-%d %H:%M:%S";
pub fn serialize<S>(date: &Option<DateTime<Local>>, serializer: S) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
match date {
Some(d) => serializer.serialize_str(&d.format(FORMAT).to_string()),
None => serializer.serialize_none(),
}
}
pub fn deserialize<'de, D>(deserializer: D) -> Result<Option<DateTime<Local>>, D::Error>
where
D: Deserializer<'de>,
{
let s: String = String::deserialize(deserializer)?;
let parsed_datetime =
NaiveDateTime::parse_from_str(&s, FORMAT).map_err(serde::de::Error::custom)?;
let shanghai_timezone = Shanghai;
let datetime_shanghai = shanghai_timezone
.from_local_datetime(&parsed_datetime)
.single()
.ok_or_else(|| serde::de::Error::custom("Invalid datetime in Shanghai timezone"))?;
Ok(Some(datetime_shanghai.with_timezone(&Local)))
}
}
pub fn pointer_str(value: &Value, path: &str, default: impl AsRef<str>) -> String {
pointer(value, path, default.as_ref().to_string())
}
pub fn pointer<T>(value: &Value, pointer: &str, default: T) -> T
where
T: for<'de> Deserialize<'de> + 'static,
{
if let Some(value) = value.pointer(pointer) {
serde_json::from_value(value.clone()).unwrap_or(default)
} else {
default
}
}
use crate::{event::{EventType, SocketEvent}, socket::protocol::Qos};
use super::model::TokenInfo;
pub trait BusinessEvent {
fn message(key: Option<String>, qos: Qos, msg_id: i64, msg_ty: i32) -> SocketEvent;
}
impl BusinessEvent for SocketEvent {
fn message(key: Option<String>, qos: Qos, msg_id: i64, msg_ty: i32) -> SocketEvent {
SocketEvent {
seller_id: None,
store_id: None,
device_id: None,
client_type: None,
version: None,
version_desc: None,
event: "业务推送事件".to_string(),
event_ty: EventType::BusinessEvent as i32,
key,
msg_id: Some(msg_id),
msg_ty: Some(msg_ty),
status: 1,
desc: qos.desc(),
extra: None,
}
}
}
pub trait PublishEvent {
fn push(
infos: &[TokenInfo],
key: Option<String>,
msg_id: i64,
msg_ty: i32,
status: i32,
desc: impl AsRef<str>,
extra: Option<String>,
) -> Vec<SocketEvent>;
fn ack(
info: &TokenInfo,
msg_id: i64,
desc: impl AsRef<str>,
extra: Option<String>,
) -> SocketEvent;
}
impl PublishEvent for SocketEvent {
fn push(
infos: &[TokenInfo],
key: Option<String>,
msg_id: i64,
msg_ty: i32,
status: i32,
desc: impl AsRef<str>,
extra: Option<String>,
) -> Vec<SocketEvent> {
let mut vec = Vec::new();
for info in infos {
let key = key.clone();
let extra = extra.clone();
let event = SocketEvent {
seller_id: Some(info.seller_id),
store_id: Some(info.seller_id),
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::ConnectEvent as i32,
key: key,
msg_id: Some(msg_id),
msg_ty: Some(msg_ty),
status,
desc: desc.as_ref().to_string(),
extra,
};
vec.push(event);
}
vec
}
fn ack(
info: &TokenInfo,
msg_id: i64,
desc: impl AsRef<str>,
extra: Option<String>,
) -> SocketEvent {
SocketEvent {
seller_id: Some(info.seller_id),
store_id: Some(info.seller_id),
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::ConnectEvent as i32,
key: None,
msg_id: Some(msg_id),
msg_ty: None,
status: 1,
desc: desc.as_ref().to_string(),
extra,
}
}
}
pub trait ConnectEvent {
fn connect(
token_info: &TokenInfo,
status: i32,
desc: impl AsRef<str>,
extra: Option<String>,
) -> SocketEvent;
}
impl ConnectEvent for SocketEvent {
fn connect(
info: &TokenInfo,
status: i32,
desc: impl AsRef<str>,
extra: Option<String>,
) -> SocketEvent {
SocketEvent {
seller_id: Some(info.seller_id),
store_id: Some(info.seller_id),
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: 1,
key: None,
msg_id: None,
msg_ty: None,
status,
desc: desc.as_ref().to_string(),
extra,
}
}
}
pub mod protocol; pub mod event;
\ No newline at end of file pub mod model;
pub mod protocol;
use serde::{Deserialize, Serialize};
#[derive(Clone, Debug, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct TokenInfo {
pub seller_id: i64,
pub store_id: i64,
pub device_id: String,
pub client_type: i32,
pub version: i32,
#[serde(rename = "versionName")]
pub version_desc: Option<String>,
#[serde(rename = "properties")]
pub extend: Option<String>,
}
impl TokenInfo {
pub fn generate_token(&self) -> String {
let combined_string = format!(
"{}:{}:{}:{}",
self.client_type, self.seller_id, self.store_id, self.device_id
);
crate::id::short_link_hash(&combined_string)
}
pub fn get_standard_device_id(&self, token: String) -> String {
let padded_device_id = format!("{:0>17}", self.device_id);
let sub_device_id = &padded_device_id[padded_device_id.len() - 17..];
format!("qm_cli_{}_{}", token, sub_device_id)
}
pub fn standard_device_id(&self) -> String {
// 标准设备号
if self.device_id.starts_with("qm_cli_") {
self.device_id.clone()
} else {
let token = self.generate_token();
self.get_standard_device_id(token)
}
}
pub fn standard(&self, device_id: String) -> Self {
TokenInfo {
seller_id: self.seller_id,
store_id: self.store_id,
device_id,
client_type: self.client_type,
version: self.version,
version_desc: self.version_desc.clone(),
extend: None,
}
}
}
\ No newline at end of file
...@@ -8,7 +8,6 @@ use crate::socket::protocol::{ ...@@ -8,7 +8,6 @@ use crate::socket::protocol::{
use anyhow::{anyhow, Result}; use anyhow::{anyhow, Result};
use bytes::Buf; use bytes::Buf;
use payload::{disconnect::DisConnectPayload, sub_ack::SubAckPayload}; use payload::{disconnect::DisConnectPayload, sub_ack::SubAckPayload};
use qm_rust_util::id;
use rmp::encode::RmpWrite; use rmp::encode::RmpWrite;
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use serde_repr::{Deserialize_repr, Serialize_repr}; use serde_repr::{Deserialize_repr, Serialize_repr};
...@@ -46,7 +45,7 @@ impl Protocol { ...@@ -46,7 +45,7 @@ impl Protocol {
let payload = ConnectPayload::new(device_id, username, password, client_type); let payload = ConnectPayload::new(device_id, username, password, client_type);
Ok(Self { Ok(Self {
trace_id: id::uuid(), trace_id: crate::id::uuid(),
fixed, fixed,
variable: Self::create_variable(variable)?, variable: Self::create_variable(variable)?,
payload: Self::create_payload(Some(payload))?, payload: Self::create_payload(Some(payload))?,
......
...@@ -3,7 +3,6 @@ use flate2::{ ...@@ -3,7 +3,6 @@ use flate2::{
bufread::{ZlibDecoder, ZlibEncoder}, bufread::{ZlibDecoder, ZlibEncoder},
Compression, Compression,
}; };
use qm_rust_util::id;
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use std::{io::Read, ops::Deref}; use std::{io::Read, ops::Deref};
...@@ -59,7 +58,7 @@ where ...@@ -59,7 +58,7 @@ where
type Error = anyhow::Error; type Error = anyhow::Error;
fn try_into(self) -> Result<Vec<u8>, Self::Error> { fn try_into(self) -> Result<Vec<u8>, Self::Error> {
let bytes = if id::is_buf::<T>() { let bytes = if crate::id::is_buf::<T>() {
let ptr = Box::into_raw(Box::new(self.0)) as *mut Vec<u8>; let ptr = Box::into_raw(Box::new(self.0)) as *mut Vec<u8>;
*unsafe { Box::from_raw(ptr) } *unsafe { Box::from_raw(ptr) }
} else { } else {
......
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