Skip to main content

sz_orm_core/cdc/
event.rs

1//! CDC 标准变更事件(v6.8.0)
2
3use serde::{Deserialize, Serialize};
4
5/// 变更事件类型
6#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
7pub enum ChangeEventType {
8    /// 插入
9    Insert,
10    /// 更新
11    Update,
12    /// 删除
13    Delete,
14}
15
16/// 变更位点
17#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
18pub enum ChangePosition {
19    /// MySQL Binlog 位点
20    MysqlBinlog {
21        /// Binlog 文件名
22        filename: String,
23        /// 文件内偏移
24        position: u64,
25    },
26    /// PostgreSQL WAL 位点
27    PostgresWal {
28        /// WAL LSN
29        lsn: u64,
30    },
31    /// SQLite update-hook 序号
32    SqliteHook {
33        /// 递增序号
34        seq: u64,
35    },
36}
37
38impl ChangePosition {
39    /// 返回位点单调递增的比较键
40    pub fn order_key(&self) -> (String, u64) {
41        match self {
42            ChangePosition::MysqlBinlog { filename, position } => (filename.clone(), *position),
43            ChangePosition::PostgresWal { lsn } => ("pg_wal".to_string(), *lsn),
44            ChangePosition::SqliteHook { seq } => ("sqlite".to_string(), *seq),
45        }
46    }
47}
48
49/// CDC 标准变更事件
50#[derive(Debug, Clone, Serialize, Deserialize)]
51pub struct ChangeEvent {
52    /// 事件唯一 ID
53    pub event_id: String,
54    /// 事件类型
55    pub event_type: ChangeEventType,
56    /// 源数据库名
57    pub source_db: String,
58    /// 源表名
59    pub source_table: String,
60    /// 行数据(JSON)
61    pub row_data: serde_json::Value,
62    /// 变更位点
63    pub position: ChangePosition,
64    /// 时间戳(毫秒)
65    pub timestamp_ms: u64,
66    /// 是否已脱敏
67    pub masked: bool,
68}
69
70impl ChangeEvent {
71    /// 创建新事件
72    pub fn new(
73        event_type: ChangeEventType,
74        source_db: &str,
75        source_table: &str,
76        row_data: serde_json::Value,
77        position: ChangePosition,
78        timestamp_ms: u64,
79    ) -> Self {
80        let event_id = format!(
81            "{}:{}:{}:{}",
82            source_db,
83            source_table,
84            timestamp_ms,
85            position.order_key().1
86        );
87        Self {
88            event_id,
89            event_type,
90            source_db: source_db.to_string(),
91            source_table: source_table.to_string(),
92            row_data,
93            position,
94            timestamp_ms,
95            masked: false,
96        }
97    }
98
99    /// 标记为已脱敏
100    pub fn with_masked(mut self, masked: bool) -> Self {
101        self.masked = masked;
102        self
103    }
104}
105
106#[cfg(test)]
107mod tests {
108    use super::*;
109
110    #[test]
111    fn event_creation_insert() {
112        let event = ChangeEvent::new(
113            ChangeEventType::Insert,
114            "my_db",
115            "users",
116            serde_json::json!({"id": 1, "name": "Alice"}),
117            ChangePosition::MysqlBinlog {
118                filename: "bin.000001".to_string(),
119                position: 100,
120            },
121            1700000000000,
122        );
123        assert_eq!(event.event_type, ChangeEventType::Insert);
124        assert_eq!(event.source_db, "my_db");
125        assert_eq!(event.source_table, "users");
126        assert!(!event.masked);
127        assert!(!event.event_id.is_empty());
128    }
129
130    #[test]
131    fn event_type_serde_roundtrip() {
132        let et = ChangeEventType::Update;
133        let json = serde_json::to_string(&et).unwrap();
134        let de: ChangeEventType = serde_json::from_str(&json).unwrap();
135        assert_eq!(et, de);
136    }
137
138    #[test]
139    fn position_order_key_monotonic() {
140        let p1 = ChangePosition::MysqlBinlog {
141            filename: "bin.000001".to_string(),
142            position: 100,
143        };
144        let p2 = ChangePosition::MysqlBinlog {
145            filename: "bin.000001".to_string(),
146            position: 200,
147        };
148        assert!(p1.order_key().1 < p2.order_key().1);
149    }
150
151    #[test]
152    fn event_id_uniqueness() {
153        let pos = ChangePosition::MysqlBinlog {
154            filename: "bin.000001".to_string(),
155            position: 100,
156        };
157        let e1 = ChangeEvent::new(
158            ChangeEventType::Insert,
159            "db",
160            "tbl",
161            serde_json::json!({}),
162            pos.clone(),
163            1000,
164        );
165        let e2 = ChangeEvent::new(
166            ChangeEventType::Insert,
167            "db",
168            "tbl",
169            serde_json::json!({}),
170            ChangePosition::MysqlBinlog {
171                filename: "bin.000001".to_string(),
172                position: 200,
173            },
174            1000,
175        );
176        assert_ne!(e1.event_id, e2.event_id);
177    }
178
179    #[test]
180    fn event_with_masked_flag() {
181        let event = ChangeEvent::new(
182            ChangeEventType::Update,
183            "db",
184            "users",
185            serde_json::json!({"phone": "13800138000"}),
186            ChangePosition::MysqlBinlog {
187                filename: "bin.000001".to_string(),
188                position: 100,
189            },
190            1000,
191        )
192        .with_masked(true);
193        assert!(event.masked);
194    }
195
196    #[test]
197    fn position_variants() {
198        let mysql_pos = ChangePosition::MysqlBinlog {
199            filename: "bin.001".to_string(),
200            position: 42,
201        };
202        let pg_pos = ChangePosition::PostgresWal { lsn: 12345 };
203        let sqlite_pos = ChangePosition::SqliteHook { seq: 99 };
204
205        assert_eq!(mysql_pos.order_key(), ("bin.001".to_string(), 42));
206        assert_eq!(pg_pos.order_key(), ("pg_wal".to_string(), 12345));
207        assert_eq!(sqlite_pos.order_key(), ("sqlite".to_string(), 99));
208    }
209}