mysql_common/binlog/events/
previous_gtids_event.rs1use std::{cmp::min, io};
24
25use crate::{
26 binlog::{
27 BinlogCtx, BinlogEvent, BinlogStruct,
28 consts::{BinlogVersion, EventType},
29 },
30 io::ParseBuf,
31 misc::{read_varlen_uint, write_varlen_uint},
32 packets::{GnoInterval, Tag},
33 proto::{MyDeserialize, MySerialize},
34};
35
36use super::BinlogEventHeader;
37
38#[derive(Debug, Clone, Eq, PartialEq, Hash)]
40pub struct PreviousGtidsSid<'a> {
41 uuid: [u8; 16],
42 tag: Option<Tag<'a>>,
43 intervals: Vec<GnoInterval>,
44}
45
46impl<'a> PreviousGtidsSid<'a> {
47 pub fn new(uuid: [u8; 16], tag: Option<Tag<'a>>, intervals: Vec<GnoInterval>) -> Self {
49 Self {
50 uuid,
51 tag,
52 intervals,
53 }
54 }
55
56 pub fn uuid(&self) -> [u8; 16] {
58 self.uuid
59 }
60
61 pub fn tag(&self) -> Option<&Tag<'a>> {
63 self.tag.as_ref()
64 }
65
66 pub fn intervals(&self) -> &[GnoInterval] {
68 &self.intervals
69 }
70
71 pub fn into_owned(self) -> PreviousGtidsSid<'static> {
73 PreviousGtidsSid {
74 uuid: self.uuid,
75 tag: self.tag.map(|t| t.into_owned()),
76 intervals: self.intervals,
77 }
78 }
79}
80
81#[derive(Debug, Clone, Eq, PartialEq, Hash)]
89pub struct PreviousGtidsEvent<'a> {
90 sids: Vec<PreviousGtidsSid<'a>>,
91}
92
93impl<'a> PreviousGtidsEvent<'a> {
94 pub fn new(sids: Vec<PreviousGtidsSid<'a>>) -> Self {
96 Self { sids }
97 }
98
99 pub fn sids(&self) -> &[PreviousGtidsSid<'a>] {
101 &self.sids
102 }
103
104 pub fn into_owned(self) -> PreviousGtidsEvent<'static> {
106 PreviousGtidsEvent {
107 sids: self.sids.into_iter().map(|s| s.into_owned()).collect(),
108 }
109 }
110
111 fn is_tagged(&self) -> bool {
113 self.sids.iter().any(|s| s.tag.is_some())
114 }
115}
116
117impl<'de> MyDeserialize<'de> for PreviousGtidsEvent<'de> {
118 const SIZE: Option<usize> = None;
119 type Ctx = BinlogCtx<'de>;
120
121 fn deserialize(_ctx: Self::Ctx, buf: &mut ParseBuf<'de>) -> io::Result<Self> {
122 if buf.len() < 8 {
123 return Err(io::Error::new(
124 io::ErrorKind::UnexpectedEof,
125 "PreviousGtidsEvent too short for header",
126 ));
127 }
128
129 let mut header_bytes = [0u8; 8];
130 header_bytes.copy_from_slice(&buf.0[..8]);
131 buf.0 = &buf.0[8..];
132 let header = u64::from_le_bytes(header_bytes);
133
134 let format = (header >> 56) as u8;
135 let tagged = format != 0;
136
137 let n_sids = if tagged {
138 (header >> 8) & ((1u64 << 48) - 1)
139 } else {
140 header & ((1u64 << 56) - 1)
141 };
142
143 let mut sids = Vec::with_capacity(n_sids as usize);
144
145 for _ in 0..n_sids {
146 if buf.len() < 16 {
148 return Err(io::Error::new(
149 io::ErrorKind::UnexpectedEof,
150 "PreviousGtidsEvent: unexpected EOF reading UUID",
151 ));
152 }
153 let mut uuid = [0u8; 16];
154 uuid.copy_from_slice(&buf.0[..16]);
155 buf.0 = &buf.0[16..];
156
157 let tag = if tagged {
159 let tag_len = read_varlen_uint(buf)? as usize;
160 if tag_len == 0 {
161 None
162 } else {
163 if buf.len() < tag_len {
164 return Err(io::Error::new(
165 io::ErrorKind::UnexpectedEof,
166 "PreviousGtidsEvent: unexpected EOF reading tag",
167 ));
168 }
169 let tag_bytes = &buf.0[..tag_len];
170 buf.0 = &buf.0[tag_len..];
171 let tag_str = std::str::from_utf8(tag_bytes).map_err(|e| {
172 io::Error::new(
173 io::ErrorKind::InvalidData,
174 format!("invalid UTF-8 tag: {}", e),
175 )
176 })?;
177 Some(Tag::new(tag_str).map_err(|e| {
178 io::Error::new(
179 io::ErrorKind::InvalidData,
180 format!("invalid GTID tag: {}", e),
181 )
182 })?)
183 }
184 } else {
185 None
186 };
187
188 if buf.len() < 8 {
190 return Err(io::Error::new(
191 io::ErrorKind::UnexpectedEof,
192 "PreviousGtidsEvent: unexpected EOF reading n_intervals",
193 ));
194 }
195 let mut n_intervals_bytes = [0u8; 8];
196 n_intervals_bytes.copy_from_slice(&buf.0[..8]);
197 buf.0 = &buf.0[8..];
198 let n_intervals = u64::from_le_bytes(n_intervals_bytes);
199
200 let mut intervals = Vec::with_capacity(n_intervals as usize);
202 for _ in 0..n_intervals {
203 let interval: GnoInterval = buf.parse(())?;
204 intervals.push(interval);
205 }
206
207 sids.push(PreviousGtidsSid {
208 uuid,
209 tag,
210 intervals,
211 });
212 }
213
214 Ok(Self { sids })
215 }
216}
217
218impl MySerialize for PreviousGtidsEvent<'_> {
219 fn serialize(&self, buf: &mut Vec<u8>) {
220 let tagged = self.is_tagged();
221 let n_sids = self.sids.len() as u64;
222
223 let header = if tagged {
224 let format: u64 = 1;
225 (format << 56) | (n_sids << 8) | format
226 } else {
227 n_sids
228 };
229
230 buf.extend_from_slice(&header.to_le_bytes());
231
232 for sid in &self.sids {
233 buf.extend_from_slice(&sid.uuid);
235
236 if tagged {
238 match &sid.tag {
239 Some(tag) => {
240 let tag_bytes = tag.as_str().as_bytes();
241 write_varlen_uint(buf, tag_bytes.len() as u64);
242 buf.extend_from_slice(tag_bytes);
243 }
244 None => {
245 write_varlen_uint(buf, 0);
246 }
247 }
248 }
249
250 buf.extend_from_slice(&(sid.intervals.len() as u64).to_le_bytes());
252
253 for interval in &sid.intervals {
255 interval.serialize(buf);
256 }
257 }
258 }
259}
260
261impl<'a> BinlogStruct<'a> for PreviousGtidsEvent<'a> {
262 fn len(&self, _version: BinlogVersion) -> usize {
263 let mut tmp = Vec::new();
264 self.serialize(&mut tmp);
265 min(tmp.len(), u32::MAX as usize - BinlogEventHeader::LEN)
266 }
267}
268
269impl<'a> BinlogEvent<'a> for PreviousGtidsEvent<'a> {
270 const EVENT_TYPE: EventType = EventType::PREVIOUS_GTIDS_EVENT;
271}