1use bitvec::prelude::*;
10use byteorder::{LittleEndian, WriteBytesExt};
11use bytes::BufMut;
12use saturating::Saturating as S;
13
14use crate::{
15 io::ParseBuf,
16 misc::raw::{RawConst, RawFlags, int::*},
17 proto::{MyDeserialize, MySerialize},
18};
19
20pub use self::{
21 anonymous_gtid_event::AnonymousGtidEvent,
22 begin_load_query_event::BeginLoadQueryEvent,
23 delete_rows_event::DeleteRowsEvent,
24 delete_rows_event_v1::DeleteRowsEventV1,
25 execute_load_query_event::ExecuteLoadQueryEvent,
26 format_description_event::FormatDescriptionEvent,
27 gtid_event::GtidEvent,
28 incident_event::IncidentEvent,
29 intvar_event::IntvarEvent,
30 partial_update_rows_event::PartialUpdateRowsEvent,
31 previous_gtids_event::{PreviousGtidsEvent, PreviousGtidsSid},
32 query_event::{QueryEvent, StatusVar, StatusVarVal, StatusVars, StatusVarsIterator},
33 rand_event::RandEvent,
34 rotate_event::RotateEvent,
35 rows_event::{RowsEvent, RowsEventRows},
36 rows_query_event::RowsQueryEvent,
37 table_map_event::*,
38 transaction_payload_event::{TransactionPayloadEvent, TransactionPayloadReader},
39 update_rows_event::UpdateRowsEvent,
40 update_rows_event_v1::UpdateRowsEventV1,
41 user_var_event::UserVarEvent,
42 write_rows_event::WriteRowsEvent,
43 write_rows_event_v1::WriteRowsEventV1,
44 xid_event::XidEvent,
45};
46
47use std::{
48 any::type_name,
49 borrow::Cow,
50 cmp::min,
51 io::{self, Read, Write},
52};
53
54use super::{
55 BinlogCtx, BinlogEvent,
56 consts::{
57 BinlogChecksumAlg, BinlogVersion, EventFlags, EventType, RowsEventFlags,
58 UnknownChecksumAlg, UnknownEventType,
59 },
60 misc::LimitWrite,
61};
62
63mod anonymous_gtid_event;
64mod begin_load_query_event;
65mod delete_rows_event;
66mod delete_rows_event_v1;
67mod execute_load_query_event;
68mod format_description_event;
69mod gtid_event;
70
71mod incident_event;
72mod intvar_event;
73mod partial_update_rows_event;
74mod previous_gtids_event;
75mod query_event;
76mod rand_event;
77mod rotate_event;
78mod rows_event;
79mod rows_query_event;
80mod table_map_event;
81mod transaction_payload_event;
82mod update_rows_event;
83mod update_rows_event_v1;
84mod user_var_event;
85mod write_rows_event;
86mod write_rows_event_v1;
87mod xid_event;
88
89#[derive(Debug, Clone, Eq, PartialEq)]
94pub struct Event {
95 fde: FormatDescriptionEvent<'static>,
97 header: BinlogEventHeader,
99 data: Vec<u8>,
106 footer: BinlogEventFooter,
108 checksum: [u8; BinlogEventFooter::BINLOG_CHECKSUM_LEN],
112}
113
114impl Event {
115 pub fn read<'a, T: Read>(
117 fde: &'a FormatDescriptionEvent<'a>,
118 mut input: T,
119 ) -> io::Result<Self> {
120 let binlog_header_len = BinlogEventHeader::LEN;
121 let mut fde = fde.clone().into_owned();
122
123 let mut header_buf = [0u8; BinlogEventHeader::LEN];
124 input.read_exact(&mut header_buf)?;
125 let header = BinlogEventHeader::deserialize((), &mut ParseBuf(&header_buf))?;
126
127 let mut data = vec![0_u8; (S(header.event_size() as usize) - S(binlog_header_len)).0];
128 input.read_exact(&mut data).unwrap();
129
130 let is_fde = header.event_type.0 == EventType::FORMAT_DESCRIPTION_EVENT as u8;
131 let mut bytes_to_truncate = 0;
132 let mut checksum = [0_u8; BinlogEventFooter::BINLOG_CHECKSUM_LEN];
133
134 let footer = if is_fde {
135 let footer = BinlogEventFooter::read(&data)?;
136 if footer.checksum_alg.is_some() {
137 bytes_to_truncate += BinlogEventFooter::BINLOG_CHECKSUM_ALG_DESC_LEN;
139 }
140 fde = fde.with_footer(footer);
142 footer
143 } else {
144 fde.footer()
145 };
146
147 let contains_checksum = footer.checksum_alg.is_some()
150 && (is_fde || footer.checksum_alg != Some(RawConst::new(0)))
151 && footer.checksum_enabled;
152
153 if contains_checksum {
154 bytes_to_truncate += BinlogEventFooter::BINLOG_CHECKSUM_LEN;
156 checksum.copy_from_slice(&data[data.len() - BinlogEventFooter::BINLOG_CHECKSUM_LEN..]);
157 }
158
159 data.truncate(data.len() - bytes_to_truncate);
160
161 Ok(Self {
162 fde,
163 header,
164 data,
165 footer,
166 checksum,
167 })
168 }
169
170 pub fn write<T: Write>(&self, version: BinlogVersion, mut output: T) -> io::Result<()> {
172 let is_fde = self.header.event_type.0 == EventType::FORMAT_DESCRIPTION_EVENT as u8;
173 let mut output = output.limit(S(self.len(version)));
174
175 let mut header_buf = Vec::with_capacity(BinlogEventHeader::LEN);
176 self.header.serialize(&mut header_buf);
177 output.write_all(&header_buf)?;
178 output.write_all(&self.data)?;
179
180 if let Ok(Some(alg)) = self.footer.get_checksum_alg() {
181 if is_fde {
182 output.write_u8(alg as u8)?;
183 }
184 if alg == BinlogChecksumAlg::BINLOG_CHECKSUM_ALG_CRC32 || is_fde {
185 output.write_u32::<LittleEndian>(self.calc_checksum(alg))?;
186 }
187 }
188
189 Ok(())
190 }
191
192 fn len(&self, _version: BinlogVersion) -> usize {
194 let is_fde = self.header.event_type.0 == EventType::FORMAT_DESCRIPTION_EVENT as u8;
195 let mut len = S(0);
196
197 len += S(BinlogEventHeader::LEN);
198 len += S(self.data.len());
199 if let Ok(Some(alg)) = self.footer.get_checksum_alg() {
200 if is_fde {
201 len += S(BinlogEventFooter::BINLOG_CHECKSUM_ALG_DESC_LEN);
202 }
203 if is_fde || alg != BinlogChecksumAlg::BINLOG_CHECKSUM_ALG_OFF {
204 len += S(BinlogEventFooter::BINLOG_CHECKSUM_LEN);
205 }
206 }
207
208 min(len.0, u32::MAX as usize - BinlogEventHeader::LEN)
209 }
210
211 pub fn fde(&self) -> &FormatDescriptionEvent<'static> {
213 &self.fde
214 }
215
216 pub fn header(&self) -> BinlogEventHeader {
218 self.header
219 }
220
221 pub fn data(&self) -> &[u8] {
223 &self.data
224 }
225
226 pub fn footer(&self) -> BinlogEventFooter {
228 self.footer
229 }
230
231 pub fn checksum(&self) -> Option<[u8; BinlogEventFooter::BINLOG_CHECKSUM_LEN]> {
233 let contains_checksum = self.footer.checksum_alg.is_some()
234 && (self.header.event_type.0 == (EventType::FORMAT_DESCRIPTION_EVENT as u8)
235 || self.footer.checksum_alg != Some(RawConst::new(0)));
236 contains_checksum.then_some(self.checksum)
237 }
238
239 pub fn read_event<'a, T: BinlogEvent<'a>>(&'a self) -> io::Result<T> {
241 let event_size = BinlogEventHeader::LEN + self.data.len();
243 let event_data = &mut ParseBuf(&self.data);
244 let ctx = BinlogCtx::new(event_size, &self.fde, self.header.event_type_raw());
245
246 let event = event_data.parse(ctx)?;
247
248 if !event_data.is_empty() {
250 return Err(io::Error::other(format!(
251 "bytes remaining on stream while reading {}",
252 type_name::<T>()
253 )));
254 }
255
256 Ok(event)
257 }
258
259 pub fn read_data(&self) -> io::Result<Option<EventData<'_>>> {
261 use EventType::*;
262
263 let event_type = match self.header.event_type.get() {
264 Ok(event_type) => event_type,
265 _ => return Ok(None),
266 };
267
268 let event_data = match event_type {
269 ENUM_END_EVENT | UNKNOWN_EVENT => EventData::UnknownEvent,
270 START_EVENT_V3 => EventData::StartEventV3(Cow::Borrowed(&*self.data)),
271 QUERY_EVENT => EventData::QueryEvent(self.read_event()?),
272 STOP_EVENT => EventData::StopEvent,
273 ROTATE_EVENT => EventData::RotateEvent(self.read_event()?),
274 INTVAR_EVENT => EventData::IntvarEvent(self.read_event()?),
275 LOAD_EVENT => EventData::LoadEvent(Cow::Borrowed(&*self.data)),
276 SLAVE_EVENT => EventData::SlaveEvent,
277 CREATE_FILE_EVENT => EventData::CreateFileEvent(Cow::Borrowed(&*self.data)),
278 APPEND_BLOCK_EVENT => EventData::AppendBlockEvent(Cow::Borrowed(&*self.data)),
279 EXEC_LOAD_EVENT => EventData::ExecLoadEvent(Cow::Borrowed(&*self.data)),
280 DELETE_FILE_EVENT => EventData::DeleteFileEvent(Cow::Borrowed(&*self.data)),
281 NEW_LOAD_EVENT => EventData::NewLoadEvent(Cow::Borrowed(&*self.data)),
282 RAND_EVENT => EventData::RandEvent(self.read_event()?),
283 USER_VAR_EVENT => EventData::UserVarEvent(self.read_event()?),
284 FORMAT_DESCRIPTION_EVENT => {
285 let fde = self
286 .read_event::<FormatDescriptionEvent<'_>>()?
287 .with_footer(self.footer);
288 EventData::FormatDescriptionEvent(fde)
289 }
290 XID_EVENT => EventData::XidEvent(self.read_event()?),
291 BEGIN_LOAD_QUERY_EVENT => EventData::BeginLoadQueryEvent(self.read_event()?),
292 EXECUTE_LOAD_QUERY_EVENT => EventData::ExecuteLoadQueryEvent(self.read_event()?),
293 TABLE_MAP_EVENT => EventData::TableMapEvent(self.read_event()?),
294 PRE_GA_WRITE_ROWS_EVENT => EventData::PreGaWriteRowsEvent(Cow::Borrowed(&*self.data)),
295 PRE_GA_UPDATE_ROWS_EVENT => EventData::PreGaUpdateRowsEvent(Cow::Borrowed(&*self.data)),
296 PRE_GA_DELETE_ROWS_EVENT => EventData::PreGaDeleteRowsEvent(Cow::Borrowed(&*self.data)),
297 WRITE_ROWS_EVENT_V1 => {
298 EventData::RowsEvent(RowsEventData::WriteRowsEventV1(self.read_event()?))
299 }
300 UPDATE_ROWS_EVENT_V1 => {
301 EventData::RowsEvent(RowsEventData::UpdateRowsEventV1(self.read_event()?))
302 }
303 DELETE_ROWS_EVENT_V1 => {
304 EventData::RowsEvent(RowsEventData::DeleteRowsEventV1(self.read_event()?))
305 }
306 INCIDENT_EVENT => EventData::IncidentEvent(self.read_event()?),
307 HEARTBEAT_EVENT => EventData::HeartbeatEvent,
308 IGNORABLE_EVENT => EventData::IgnorableEvent(Cow::Borrowed(&*self.data)),
309 ROWS_QUERY_EVENT => EventData::RowsQueryEvent(self.read_event()?),
310 WRITE_ROWS_EVENT => {
311 EventData::RowsEvent(RowsEventData::WriteRowsEvent(self.read_event()?))
312 }
313 UPDATE_ROWS_EVENT => {
314 EventData::RowsEvent(RowsEventData::UpdateRowsEvent(self.read_event()?))
315 }
316 DELETE_ROWS_EVENT => {
317 EventData::RowsEvent(RowsEventData::DeleteRowsEvent(self.read_event()?))
318 }
319 GTID_EVENT => EventData::GtidEvent(self.read_event()?),
320 ANONYMOUS_GTID_EVENT => EventData::AnonymousGtidEvent(self.read_event()?),
321 PREVIOUS_GTIDS_EVENT => EventData::PreviousGtidsEvent(self.read_event()?),
322 TRANSACTION_CONTEXT_EVENT => {
323 EventData::TransactionContextEvent(Cow::Borrowed(&*self.data))
324 }
325 VIEW_CHANGE_EVENT => EventData::ViewChangeEvent(Cow::Borrowed(&*self.data)),
326 XA_PREPARE_LOG_EVENT => EventData::XaPrepareLogEvent(Cow::Borrowed(&*self.data)),
327 PARTIAL_UPDATE_ROWS_EVENT => {
328 EventData::RowsEvent(RowsEventData::PartialUpdateRowsEvent(self.read_event()?))
329 }
330 TRANSACTION_PAYLOAD_EVENT => EventData::TransactionPayloadEvent(self.read_event()?),
331 GTID_TAGGED_LOG_EVENT => EventData::GtidEvent(self.read_event()?),
332 };
333
334 Ok(Some(event_data))
335 }
336
337 pub fn calc_checksum(&self, alg: BinlogChecksumAlg) -> u32 {
339 let is_fde = self.header.event_type.0 == EventType::FORMAT_DESCRIPTION_EVENT as u8;
340
341 let mut hasher = crc32fast::Hasher::new();
342 let mut header = Vec::with_capacity(BinlogEventHeader::LEN);
343 let mut header_struct = self.header;
344 if header_struct
345 .flags
346 .get()
347 .contains(EventFlags::LOG_EVENT_BINLOG_IN_USE_F)
348 {
349 header_struct.flags.0 &= !(EventFlags::LOG_EVENT_BINLOG_IN_USE_F.bits());
354 }
355 header_struct.serialize(&mut header);
356 hasher.update(&header);
357 hasher.update(&self.data);
358 if is_fde {
359 hasher.update(&[alg as u8][..]);
360 }
361 hasher.finalize()
362 }
363}
364
365#[derive(Debug, Clone, Copy, Eq, PartialEq, Hash)]
367pub struct BinlogEventHeader {
368 timestamp: RawInt<LeU32>,
370 event_type: RawConst<u8, EventType>,
372 server_id: RawInt<LeU32>,
376 event_size: RawInt<LeU32>,
378 log_pos: RawInt<LeU32>,
380 flags: RawFlags<EventFlags, LeU16>,
384}
385
386impl BinlogEventHeader {
387 pub const LEN: usize = 19;
389
390 pub fn new(
392 timestamp: u32,
393 event_type: EventType,
394 server_id: u32,
395 event_size: u32,
396 log_pos: u32,
397 flags: EventFlags,
398 ) -> Self {
399 Self {
400 timestamp: RawInt::new(timestamp),
401 event_type: RawConst::new(event_type as u8),
402 server_id: RawInt::new(server_id),
403 event_size: RawInt::new(event_size),
404 log_pos: RawInt::new(log_pos),
405 flags: RawFlags::new(flags.bits()),
406 }
407 }
408
409 pub fn timestamp(&self) -> u32 {
413 self.timestamp.0
414 }
415
416 pub fn event_type_raw(&self) -> u8 {
418 self.event_type.0
419 }
420
421 pub fn event_type(&self) -> Result<EventType, UnknownEventType> {
423 self.event_type.get()
424 }
425
426 pub fn server_id(&self) -> u32 {
430 self.server_id.0
431 }
432
433 pub fn event_size(&self) -> u32 {
435 self.event_size.0
436 }
437
438 pub fn log_pos(&self) -> u32 {
440 self.log_pos.0
441 }
442
443 pub fn flags(&self) -> EventFlags {
445 self.flags.get()
446 }
447
448 pub fn flags_raw(&self) -> u16 {
450 self.flags.0
451 }
452}
453
454impl<'de> MyDeserialize<'de> for BinlogEventHeader {
455 const SIZE: Option<usize> = Some(Self::LEN);
456 type Ctx = ();
457
458 fn deserialize((): Self::Ctx, buf: &mut ParseBuf<'de>) -> io::Result<Self> {
459 let mut buf: ParseBuf<'_> = buf.parse_unchecked(Self::LEN)?;
460 Ok(Self {
461 timestamp: buf.parse_unchecked(())?,
462 event_type: buf.parse_unchecked(())?,
463 server_id: buf.parse_unchecked(())?,
464 event_size: buf.parse_unchecked(())?,
465 log_pos: buf.parse_unchecked(())?,
466 flags: buf.parse_unchecked(())?,
467 })
468 }
469}
470
471impl MySerialize for BinlogEventHeader {
472 fn serialize(&self, buf: &mut Vec<u8>) {
473 self.timestamp.serialize(&mut *buf);
474 self.event_type.serialize(&mut *buf);
475 self.server_id.serialize(&mut *buf);
476 self.event_size.serialize(&mut *buf);
477 self.log_pos.serialize(&mut *buf);
478 self.flags.serialize(&mut *buf);
479 }
480}
481
482#[derive(Debug, Clone, Copy, Eq, PartialEq, Hash)]
484pub struct BinlogEventFooter {
485 checksum_alg: Option<RawConst<u8, BinlogChecksumAlg>>,
487
488 checksum_enabled: bool,
490}
491
492impl BinlogEventFooter {
493 pub const BINLOG_CHECKSUM_ALG_DESC_LEN: usize = 1;
495 pub const BINLOG_CHECKSUM_LEN: usize = 4;
497 pub const CHECKSUM_VERSION_PRODUCT: (u8, u8, u8) = (5, 6, 1);
499
500 pub fn new(checksum_alg: BinlogChecksumAlg) -> Self {
501 Self {
502 checksum_alg: Some(RawConst::new(checksum_alg as u8)),
503 checksum_enabled: true,
504 }
505 }
506
507 pub fn get_checksum_alg(&self) -> Result<Option<BinlogChecksumAlg>, UnknownChecksumAlg> {
509 self.checksum_alg.as_ref().map(RawConst::get).transpose()
510 }
511
512 pub fn get_checksum_enabled(&self) -> bool {
514 self.checksum_enabled
515 }
516
517 pub fn set_checksum_enabled(&mut self, enabled: bool) {
519 self.checksum_enabled = enabled;
520 }
521
522 pub fn read(buf: &[u8]) -> io::Result<Self> {
526 let checksum_alg = if buf.len()
527 >= FormatDescriptionEvent::SERVER_VER_OFFSET + FormatDescriptionEvent::SERVER_VER_LEN
528 {
529 let mut server_version = vec![0_u8; FormatDescriptionEvent::SERVER_VER_LEN];
530 (&buf[FormatDescriptionEvent::SERVER_VER_OFFSET..]).read_exact(&mut server_version)?;
531 server_version[FormatDescriptionEvent::SERVER_VER_LEN - 1] = 0;
532 let version = crate::misc::split_version(&server_version);
533 if version < Self::CHECKSUM_VERSION_PRODUCT {
534 None
535 } else {
536 let offset = buf.len()
537 - (BinlogEventFooter::BINLOG_CHECKSUM_ALG_DESC_LEN
538 + BinlogEventFooter::BINLOG_CHECKSUM_LEN);
539 Some(buf[offset])
540 }
541 } else {
542 None
543 };
544
545 Ok(Self {
546 checksum_alg: checksum_alg.map(RawConst::new),
547 checksum_enabled: true,
548 })
549 }
550}
551
552impl Default for BinlogEventFooter {
553 fn default() -> Self {
554 BinlogEventFooter {
555 checksum_alg: Some(RawConst::new(
556 BinlogChecksumAlg::BINLOG_CHECKSUM_ALG_OFF as u8,
557 )),
558 checksum_enabled: true,
559 }
560 }
561}
562
563#[derive(Debug, Clone, Eq, PartialEq, Hash)]
565pub enum EventData<'a> {
566 UnknownEvent,
567 StartEventV3(Cow<'a, [u8]>),
569 QueryEvent(QueryEvent<'a>),
570 StopEvent,
571 RotateEvent(RotateEvent<'a>),
572 IntvarEvent(IntvarEvent),
573 LoadEvent(Cow<'a, [u8]>),
575 SlaveEvent,
576 CreateFileEvent(Cow<'a, [u8]>),
577 AppendBlockEvent(Cow<'a, [u8]>),
579 ExecLoadEvent(Cow<'a, [u8]>),
581 DeleteFileEvent(Cow<'a, [u8]>),
583 NewLoadEvent(Cow<'a, [u8]>),
585 RandEvent(RandEvent),
586 UserVarEvent(UserVarEvent<'a>),
587 FormatDescriptionEvent(FormatDescriptionEvent<'a>),
588 XidEvent(XidEvent),
589 BeginLoadQueryEvent(BeginLoadQueryEvent<'a>),
590 ExecuteLoadQueryEvent(ExecuteLoadQueryEvent<'a>),
591 TableMapEvent(TableMapEvent<'a>),
592 PreGaWriteRowsEvent(Cow<'a, [u8]>),
594 PreGaUpdateRowsEvent(Cow<'a, [u8]>),
596 PreGaDeleteRowsEvent(Cow<'a, [u8]>),
598 IncidentEvent(IncidentEvent<'a>),
599 HeartbeatEvent,
600 IgnorableEvent(Cow<'a, [u8]>),
601 RowsQueryEvent(RowsQueryEvent<'a>),
602 GtidEvent(GtidEvent),
603 AnonymousGtidEvent(AnonymousGtidEvent),
605 PreviousGtidsEvent(PreviousGtidsEvent<'a>),
606 TransactionContextEvent(Cow<'a, [u8]>),
608 ViewChangeEvent(Cow<'a, [u8]>),
610 XaPrepareLogEvent(Cow<'a, [u8]>),
612 RowsEvent(RowsEventData<'a>),
613 TransactionPayloadEvent(TransactionPayloadEvent<'a>),
614}
615
616impl<'a> EventData<'a> {
617 pub fn into_owned(self) -> EventData<'static> {
618 match self {
619 EventData::UnknownEvent => EventData::UnknownEvent,
620 EventData::StartEventV3(ev) => EventData::StartEventV3(Cow::Owned(ev.into_owned())),
621 Self::QueryEvent(ev) => EventData::QueryEvent(ev.into_owned()),
622 Self::StopEvent => EventData::StopEvent,
623 Self::RotateEvent(ev) => EventData::RotateEvent(ev.into_owned()),
624 Self::IntvarEvent(ev) => EventData::IntvarEvent(ev),
625 Self::LoadEvent(ev) => EventData::LoadEvent(Cow::Owned(ev.into_owned())),
626 Self::SlaveEvent => EventData::SlaveEvent,
627 Self::CreateFileEvent(ev) => EventData::CreateFileEvent(Cow::Owned(ev.into_owned())),
628 Self::AppendBlockEvent(ev) => EventData::AppendBlockEvent(Cow::Owned(ev.into_owned())),
629 Self::ExecLoadEvent(ev) => EventData::ExecLoadEvent(Cow::Owned(ev.into_owned())),
630 Self::DeleteFileEvent(ev) => EventData::DeleteFileEvent(Cow::Owned(ev.into_owned())),
631 Self::NewLoadEvent(ev) => EventData::NewLoadEvent(Cow::Owned(ev.into_owned())),
632 Self::RandEvent(ev) => EventData::RandEvent(ev),
633 Self::UserVarEvent(ev) => EventData::UserVarEvent(ev.into_owned()),
634 Self::FormatDescriptionEvent(ev) => EventData::FormatDescriptionEvent(ev.into_owned()),
635 Self::XidEvent(ev) => EventData::XidEvent(ev),
636 Self::BeginLoadQueryEvent(ev) => EventData::BeginLoadQueryEvent(ev.into_owned()),
637 Self::ExecuteLoadQueryEvent(ev) => EventData::ExecuteLoadQueryEvent(ev.into_owned()),
638 Self::TableMapEvent(ev) => EventData::TableMapEvent(ev.into_owned()),
639 Self::PreGaWriteRowsEvent(ev) => {
640 EventData::PreGaWriteRowsEvent(Cow::Owned(ev.into_owned()))
641 }
642 Self::PreGaUpdateRowsEvent(ev) => {
643 EventData::PreGaUpdateRowsEvent(Cow::Owned(ev.into_owned()))
644 }
645 Self::PreGaDeleteRowsEvent(ev) => {
646 EventData::PreGaDeleteRowsEvent(Cow::Owned(ev.into_owned()))
647 }
648 Self::IncidentEvent(ev) => EventData::IncidentEvent(ev.into_owned()),
649 Self::HeartbeatEvent => EventData::HeartbeatEvent,
650 Self::IgnorableEvent(ev) => EventData::IgnorableEvent(Cow::Owned(ev.into_owned())),
651 Self::RowsQueryEvent(ev) => EventData::RowsQueryEvent(ev.into_owned()),
652 Self::GtidEvent(ev) => EventData::GtidEvent(ev),
653 Self::AnonymousGtidEvent(ev) => EventData::AnonymousGtidEvent(ev),
654 Self::PreviousGtidsEvent(ev) => EventData::PreviousGtidsEvent(ev.into_owned()),
655 Self::TransactionContextEvent(ev) => {
656 EventData::TransactionContextEvent(Cow::Owned(ev.into_owned()))
657 }
658 Self::ViewChangeEvent(ev) => EventData::ViewChangeEvent(Cow::Owned(ev.into_owned())),
659 Self::XaPrepareLogEvent(ev) => {
660 EventData::XaPrepareLogEvent(Cow::Owned(ev.into_owned()))
661 }
662 Self::RowsEvent(ev) => EventData::RowsEvent(ev.into_owned()),
663 Self::TransactionPayloadEvent(ev) => {
664 EventData::TransactionPayloadEvent(ev.into_owned())
665 }
666 }
667 }
668}
669
670impl MySerialize for EventData<'_> {
671 fn serialize(&self, buf: &mut Vec<u8>) {
672 match self {
673 EventData::UnknownEvent => (),
674 EventData::StartEventV3(ev) => buf.put_slice(ev),
675 EventData::QueryEvent(ev) => ev.serialize(buf),
676 EventData::StopEvent => (),
677 EventData::RotateEvent(ev) => ev.serialize(buf),
678 EventData::IntvarEvent(ev) => ev.serialize(buf),
679 EventData::LoadEvent(ev) => buf.put_slice(ev),
680 EventData::SlaveEvent => (),
681 EventData::CreateFileEvent(ev) => buf.put_slice(ev),
682 EventData::AppendBlockEvent(ev) => buf.put_slice(ev),
683 EventData::ExecLoadEvent(ev) => buf.put_slice(ev),
684 EventData::DeleteFileEvent(ev) => buf.put_slice(ev),
685 EventData::NewLoadEvent(ev) => buf.put_slice(ev),
686 EventData::RandEvent(ev) => ev.serialize(buf),
687 EventData::UserVarEvent(ev) => ev.serialize(buf),
688 EventData::FormatDescriptionEvent(ev) => ev.serialize(buf),
689 EventData::XidEvent(ev) => ev.serialize(buf),
690 EventData::BeginLoadQueryEvent(ev) => ev.serialize(buf),
691 EventData::ExecuteLoadQueryEvent(ev) => ev.serialize(buf),
692 EventData::TableMapEvent(ev) => ev.serialize(buf),
693 EventData::PreGaWriteRowsEvent(ev) => buf.put_slice(ev),
694 EventData::PreGaUpdateRowsEvent(ev) => buf.put_slice(ev),
695 EventData::PreGaDeleteRowsEvent(ev) => buf.put_slice(ev),
696 EventData::IncidentEvent(ev) => ev.serialize(buf),
697 EventData::HeartbeatEvent => (),
698 EventData::IgnorableEvent(ev) => buf.put_slice(ev),
699 EventData::RowsQueryEvent(ev) => ev.serialize(buf),
700 EventData::GtidEvent(ev) => ev.serialize(buf),
701 EventData::AnonymousGtidEvent(ev) => ev.serialize(buf),
702 EventData::PreviousGtidsEvent(ev) => ev.serialize(buf),
703 EventData::TransactionContextEvent(ev) => buf.put_slice(ev),
704 EventData::ViewChangeEvent(ev) => buf.put_slice(ev),
705 EventData::XaPrepareLogEvent(ev) => buf.put_slice(ev),
706 EventData::RowsEvent(ev) => ev.serialize(buf),
707 EventData::TransactionPayloadEvent(ev) => ev.serialize(buf),
708 }
709 }
710}
711
712#[derive(Debug, Clone, Eq, PartialEq, Hash)]
714pub enum RowsEventData<'a> {
715 WriteRowsEventV1(WriteRowsEventV1<'a>),
716 UpdateRowsEventV1(UpdateRowsEventV1<'a>),
717 DeleteRowsEventV1(DeleteRowsEventV1<'a>),
718 WriteRowsEvent(WriteRowsEvent<'a>),
719 UpdateRowsEvent(UpdateRowsEvent<'a>),
720 DeleteRowsEvent(DeleteRowsEvent<'a>),
721 PartialUpdateRowsEvent(PartialUpdateRowsEvent<'a>),
722}
723
724impl<'a> RowsEventData<'a> {
725 pub fn table_id(&self) -> u64 {
727 match self {
728 RowsEventData::WriteRowsEventV1(ev) => ev.table_id(),
729 RowsEventData::UpdateRowsEventV1(ev) => ev.table_id(),
730 RowsEventData::DeleteRowsEventV1(ev) => ev.table_id(),
731 RowsEventData::WriteRowsEvent(ev) => ev.table_id(),
732 RowsEventData::UpdateRowsEvent(ev) => ev.table_id(),
733 RowsEventData::DeleteRowsEvent(ev) => ev.table_id(),
734 RowsEventData::PartialUpdateRowsEvent(ev) => ev.table_id(),
735 }
736 }
737
738 pub fn num_columns(&self) -> u64 {
740 match self {
741 RowsEventData::WriteRowsEventV1(ev) => ev.num_columns(),
742 RowsEventData::UpdateRowsEventV1(ev) => ev.num_columns(),
743 RowsEventData::DeleteRowsEventV1(ev) => ev.num_columns(),
744 RowsEventData::WriteRowsEvent(ev) => ev.num_columns(),
745 RowsEventData::UpdateRowsEvent(ev) => ev.num_columns(),
746 RowsEventData::DeleteRowsEvent(ev) => ev.num_columns(),
747 RowsEventData::PartialUpdateRowsEvent(ev) => ev.num_columns(),
748 }
749 }
750
751 pub fn columns_before_image(&'a self) -> Option<&'a BitSlice<u8>> {
755 match self {
756 RowsEventData::WriteRowsEventV1(_) => None,
757 RowsEventData::UpdateRowsEventV1(ev) => Some(ev.columns_before_image()),
758 RowsEventData::DeleteRowsEventV1(ev) => Some(ev.columns_before_image()),
759 RowsEventData::WriteRowsEvent(_) => None,
760 RowsEventData::UpdateRowsEvent(ev) => Some(ev.columns_before_image()),
761 RowsEventData::DeleteRowsEvent(ev) => Some(ev.columns_before_image()),
762 RowsEventData::PartialUpdateRowsEvent(ev) => Some(ev.columns_before_image()),
763 }
764 }
765
766 pub fn columns_after_image(&'a self) -> Option<&'a BitSlice<u8>> {
770 match self {
771 RowsEventData::WriteRowsEventV1(ev) => Some(ev.columns_after_image()),
772 RowsEventData::UpdateRowsEventV1(ev) => Some(ev.columns_after_image()),
773 RowsEventData::DeleteRowsEventV1(_) => None,
774 RowsEventData::WriteRowsEvent(ev) => Some(ev.columns_after_image()),
775 RowsEventData::UpdateRowsEvent(ev) => Some(ev.columns_after_image()),
776 RowsEventData::DeleteRowsEvent(_) => None,
777 RowsEventData::PartialUpdateRowsEvent(ev) => Some(ev.columns_after_image()),
778 }
779 }
780
781 pub fn rows_data(&'a self) -> &'a [u8] {
783 match self {
784 RowsEventData::WriteRowsEventV1(ev) => ev.rows_data(),
785 RowsEventData::UpdateRowsEventV1(ev) => ev.rows_data(),
786 RowsEventData::DeleteRowsEventV1(ev) => ev.rows_data(),
787 RowsEventData::WriteRowsEvent(ev) => ev.rows_data(),
788 RowsEventData::UpdateRowsEvent(ev) => ev.rows_data(),
789 RowsEventData::DeleteRowsEvent(ev) => ev.rows_data(),
790 RowsEventData::PartialUpdateRowsEvent(ev) => ev.rows_data(),
791 }
792 }
793
794 pub fn flags(&self) -> RowsEventFlags {
796 match self {
797 RowsEventData::WriteRowsEventV1(ev) => ev.flags(),
798 RowsEventData::UpdateRowsEventV1(ev) => ev.flags(),
799 RowsEventData::DeleteRowsEventV1(ev) => ev.flags(),
800 RowsEventData::WriteRowsEvent(ev) => ev.flags(),
801 RowsEventData::UpdateRowsEvent(ev) => ev.flags(),
802 RowsEventData::DeleteRowsEvent(ev) => ev.flags(),
803 RowsEventData::PartialUpdateRowsEvent(ev) => ev.flags(),
804 }
805 }
806
807 pub fn rows(&'a self, table_map_event: &'a TableMapEvent<'a>) -> RowsEventRows<'a> {
809 match self {
810 RowsEventData::WriteRowsEventV1(ev) => ev.rows(table_map_event),
811 RowsEventData::UpdateRowsEventV1(ev) => ev.rows(table_map_event),
812 RowsEventData::DeleteRowsEventV1(ev) => ev.rows(table_map_event),
813 RowsEventData::WriteRowsEvent(ev) => ev.rows(table_map_event),
814 RowsEventData::UpdateRowsEvent(ev) => ev.rows(table_map_event),
815 RowsEventData::DeleteRowsEvent(ev) => ev.rows(table_map_event),
816 RowsEventData::PartialUpdateRowsEvent(ev) => ev.rows(table_map_event),
817 }
818 }
819
820 pub fn into_owned(self) -> RowsEventData<'static> {
821 match self {
822 Self::WriteRowsEventV1(ev) => RowsEventData::WriteRowsEventV1(ev.into_owned()),
823 Self::UpdateRowsEventV1(ev) => RowsEventData::UpdateRowsEventV1(ev.into_owned()),
824 Self::DeleteRowsEventV1(ev) => RowsEventData::DeleteRowsEventV1(ev.into_owned()),
825 Self::WriteRowsEvent(ev) => RowsEventData::WriteRowsEvent(ev.into_owned()),
826 Self::UpdateRowsEvent(ev) => RowsEventData::UpdateRowsEvent(ev.into_owned()),
827 Self::DeleteRowsEvent(ev) => RowsEventData::DeleteRowsEvent(ev.into_owned()),
828 Self::PartialUpdateRowsEvent(ev) => {
829 RowsEventData::PartialUpdateRowsEvent(ev.into_owned())
830 }
831 }
832 }
833}
834
835impl MySerialize for RowsEventData<'_> {
836 fn serialize(&self, buf: &mut Vec<u8>) {
837 match self {
838 RowsEventData::WriteRowsEventV1(ev) => ev.serialize(buf),
839 RowsEventData::UpdateRowsEventV1(ev) => ev.serialize(buf),
840 RowsEventData::DeleteRowsEventV1(ev) => ev.serialize(buf),
841 RowsEventData::WriteRowsEvent(ev) => ev.serialize(buf),
842 RowsEventData::UpdateRowsEvent(ev) => ev.serialize(buf),
843 RowsEventData::DeleteRowsEvent(ev) => ev.serialize(buf),
844 RowsEventData::PartialUpdateRowsEvent(ev) => ev.serialize(buf),
845 }
846 }
847}