1use std::{borrow::Cow, cmp::min, io};
31
32use saturating::Saturating as S;
33
34use crate::{
35 binlog::{
36 BinlogCtx, BinlogEvent, BinlogStruct,
37 consts::{BinlogVersion, EventType, Gno, GtidFlags},
38 },
39 io::ParseBuf,
40 misc::{
41 raw::{RawConst, RawFlags, int::*},
42 read_varlen_uint, varlen_uint_size, write_varlen_uint,
43 },
44 packets::Tag,
45 proto::{MyDeserialize, MySerialize},
46};
47
48use super::BinlogEventHeader;
49
50define_const!(
51 ConstU8,
52 LogicalTimestampTypecode,
53 InvalidLogicalTimestampTypecode("Invalid logical timestamp typecode value for GTID event"),
54 2
55);
56
57mod field_id {
61 pub const GTID_FLAGS: u64 = 0;
62 pub const SID: u64 = 1;
63 pub const GNO: u64 = 2;
64 pub const TAG: u64 = 3;
65 pub const LAST_COMMITTED: u64 = 4;
66 pub const SEQUENCE_NUMBER: u64 = 5;
67 pub const IMMEDIATE_COMMIT_TIMESTAMP: u64 = 6;
68 pub const ORIGINAL_COMMIT_TIMESTAMP: u64 = 7;
69 pub const TRANSACTION_LENGTH: u64 = 8;
70 pub const IMMEDIATE_SERVER_VERSION: u64 = 9;
71 pub const ORIGINAL_SERVER_VERSION: u64 = 10;
72 pub const COMMIT_GROUP_TICKET: u64 = 11;
73}
74
75#[derive(Debug, Clone, Eq, PartialEq, Hash)]
84pub struct GtidEvent {
85 flags: RawFlags<GtidFlags, u8>,
87 sid: [u8; Self::ENCODED_SID_LENGTH],
89 gno: RawConst<LeU64, Gno>,
94 tag: Option<Tag<'static>>,
96 lc_typecode: Option<LogicalTimestampTypecode>,
101 last_committed: RawInt<LeU64>,
103 sequence_number: RawInt<LeU64>,
107 immediate_commit_timestamp: RawInt<LeU56>,
109 original_commit_timestamp: RawInt<LeU56>,
111 tx_length: RawInt<LenEnc>,
113 original_server_version: RawInt<LeU32>,
115 immediate_server_version: RawInt<LeU32>,
117 serialization_version: u8,
119 commit_group_ticket: u64,
121}
122
123impl GtidEvent {
124 pub const POST_HEADER_LENGTH: usize = 1 + Self::ENCODED_SID_LENGTH + 8 + 1 + 16;
125 pub const ENCODED_SID_LENGTH: usize = 16;
126 pub const LOGICAL_TIMESTAMP_TYPECODE: u8 = 2;
127 pub const IMMEDIATE_COMMIT_TIMESTAMP_LENGTH: usize = 7;
128 pub const ORIGINAL_COMMIT_TIMESTAMP_LENGTH: usize = 7;
129 pub const UNDEFINED_SERVER_VERSION: u32 = 999_999;
130 pub const IMMEDIATE_SERVER_VERSION_LENGTH: usize = 4;
131 pub const COMMIT_GROUP_TICKET_UNSET: u64 = 0;
132 pub const TAGGED_SERIALIZATION_VERSION_V2: u8 = 2;
134
135 pub fn new(sid: [u8; Self::ENCODED_SID_LENGTH], gno: u64) -> Self {
137 Self {
138 flags: Default::default(),
139 sid,
140 gno: RawConst::new(gno),
141 tag: None,
142 lc_typecode: Some(LogicalTimestampTypecode::default()),
143 last_committed: Default::default(),
144 sequence_number: Default::default(),
145 immediate_commit_timestamp: Default::default(),
146 original_commit_timestamp: Default::default(),
147 tx_length: Default::default(),
148 original_server_version: Default::default(),
149 immediate_server_version: Default::default(),
150 serialization_version: 0,
151 commit_group_ticket: Self::COMMIT_GROUP_TICKET_UNSET,
152 }
153 }
154
155 pub fn new_tagged(sid: [u8; Self::ENCODED_SID_LENGTH], tag: Tag<'static>, gno: u64) -> Self {
159 Self {
160 flags: Default::default(),
161 sid,
162 gno: RawConst::new(gno),
163 tag: Some(tag),
164 lc_typecode: None,
165 last_committed: Default::default(),
166 sequence_number: Default::default(),
167 immediate_commit_timestamp: Default::default(),
168 original_commit_timestamp: Default::default(),
169 tx_length: Default::default(),
170 original_server_version: RawInt::new(Self::UNDEFINED_SERVER_VERSION),
171 immediate_server_version: RawInt::new(Self::UNDEFINED_SERVER_VERSION),
172 serialization_version: Self::TAGGED_SERIALIZATION_VERSION_V2,
173 commit_group_ticket: Self::COMMIT_GROUP_TICKET_UNSET,
174 }
175 }
176
177 pub fn is_tagged(&self) -> bool {
179 self.tag.is_some()
180 }
181
182 pub fn event_type(&self) -> EventType {
190 if self.tag.is_some() {
191 EventType::GTID_TAGGED_LOG_EVENT
192 } else {
193 EventType::GTID_EVENT
194 }
195 }
196
197 pub fn with_flags(mut self, flags: GtidFlags) -> Self {
199 self.flags = RawFlags::new(flags.bits());
200 self
201 }
202
203 pub fn flags_raw(&self) -> u8 {
205 self.flags.0
206 }
207
208 pub fn flags(&self) -> GtidFlags {
216 self.flags.get()
217 }
218
219 pub fn with_sid(mut self, sid: [u8; Self::ENCODED_SID_LENGTH]) -> Self {
221 self.sid = sid;
222 self
223 }
224
225 pub fn sid(&self) -> [u8; Self::ENCODED_SID_LENGTH] {
229 self.sid
230 }
231
232 pub fn with_gno(mut self, gno: u64) -> Self {
234 self.gno = RawConst::new(gno);
235 self
236 }
237
238 pub fn gno(&self) -> u64 {
242 self.gno.0
243 }
244
245 pub fn tag(&self) -> Option<&Tag<'static>> {
247 self.tag.as_ref()
248 }
249
250 pub fn with_tag(mut self, tag: Tag<'static>) -> Self {
252 self.tag = Some(tag);
253 self.serialization_version = Self::TAGGED_SERIALIZATION_VERSION_V2;
254 self
255 }
256
257 pub fn lc_typecode(&self) -> Option<u8> {
262 self.lc_typecode.as_ref().map(|x| x.value())
263 }
264
265 pub fn with_lc_typecode(mut self) -> Self {
270 self.lc_typecode = Some(LogicalTimestampTypecode::default());
271 self
272 }
273
274 pub fn with_last_committed(mut self, last_committed: u64) -> Self {
276 self.last_committed = RawInt::new(last_committed);
277 self
278 }
279
280 pub fn last_committed(&self) -> u64 {
284 self.last_committed.0
285 }
286
287 pub fn with_sequence_number(mut self, sequence_number: u64) -> Self {
289 self.sequence_number = RawInt::new(sequence_number);
290 self
291 }
292
293 pub fn sequence_number(&self) -> u64 {
297 self.sequence_number.0
298 }
299
300 pub fn with_immediate_commit_timestamp(mut self, immediate_commit_timestamp: u64) -> Self {
302 self.immediate_commit_timestamp = RawInt::new(immediate_commit_timestamp);
303 self
304 }
305
306 pub fn immediate_commit_timestamp(&self) -> u64 {
310 self.immediate_commit_timestamp.0
311 }
312
313 pub fn with_original_commit_timestamp(mut self, original_commit_timestamp: u64) -> Self {
315 self.original_commit_timestamp = RawInt::new(original_commit_timestamp);
316 self
317 }
318
319 pub fn original_commit_timestamp(&self) -> u64 {
323 self.original_commit_timestamp.0
324 }
325
326 pub fn with_tx_length(mut self, tx_length: u64) -> Self {
328 self.tx_length = RawInt::new(tx_length);
329 self
330 }
331
332 pub fn tx_length(&self) -> u64 {
336 self.tx_length.0
337 }
338
339 pub fn with_original_server_version(mut self, original_server_version: u32) -> Self {
341 self.original_server_version = RawInt::new(original_server_version);
342 self
343 }
344
345 pub fn original_server_version(&self) -> u32 {
350 self.original_server_version.0
351 }
352
353 pub fn with_immediate_server_version(mut self, immediate_server_version: u32) -> Self {
355 self.immediate_server_version = RawInt::new(immediate_server_version);
356 self
357 }
358
359 pub fn immediate_server_version(&self) -> u32 {
363 self.immediate_server_version.0
364 }
365
366 pub fn commit_group_ticket(&self) -> u64 {
370 self.commit_group_ticket
371 }
372
373 pub fn with_commit_group_ticket(mut self, ticket: u64) -> Self {
375 self.commit_group_ticket = ticket;
376 self
377 }
378}
379
380fn compute_self_inclusive_payload_size(
395 extra_overhead: usize,
396 fields_size: usize,
397 lnif: u64,
398) -> u64 {
399 let fixed = extra_overhead as u64 + varlen_uint_size(lnif) as u64 + fields_size as u64;
400 let mut ps = fixed + varlen_uint_size(fields_size as u64) as u64;
403 loop {
404 let next = fixed + varlen_uint_size(ps) as u64;
405 if next == ps {
406 return ps;
407 }
408 ps = next;
409 }
410}
411
412fn read_varlen_int(buf: &mut ParseBuf<'_>) -> io::Result<i64> {
420 let unsigned = read_varlen_uint(buf)?;
421 let sign = unsigned & 1;
422 let magnitude = (unsigned >> 1) as i64;
423 if sign != 0 {
424 Ok(-magnitude - 1)
427 } else {
428 Ok(magnitude)
429 }
430}
431
432fn write_varlen_int(buf: &mut Vec<u8>, value: i64) {
439 let unsigned = if value >= 0 {
440 (value as u64) << 1
441 } else {
442 ((-(value + 1)) as u64) << 1 | 1
443 };
444 write_varlen_uint(buf, unsigned);
445}
446
447fn read_serialized_uuid(buf: &mut ParseBuf<'_>) -> io::Result<[u8; 16]> {
453 let mut uuid = [0u8; 16];
454 for (i, byte) in uuid.iter_mut().enumerate() {
455 let val = read_varlen_uint(buf)?;
456 if val > u8::MAX as u64 {
457 return Err(io::Error::new(
458 io::ErrorKind::InvalidData,
459 format!("UUID byte {} out of range: {}", i, val),
460 ));
461 }
462 *byte = val as u8;
463 }
464 Ok(uuid)
465}
466
467fn write_serialized_uuid(buf: &mut Vec<u8>, uuid: &[u8; 16]) {
469 for &byte in uuid.iter() {
470 write_varlen_uint(buf, byte as u64);
471 }
472}
473
474fn read_varlen_string<'a>(buf: &mut ParseBuf<'a>) -> io::Result<Cow<'a, str>> {
476 let raw_len = read_varlen_uint(buf)?;
477 let len: usize = raw_len.try_into().map_err(|_| {
478 io::Error::new(
479 io::ErrorKind::InvalidData,
480 format!("varlen string length {} exceeds platform usize", raw_len),
481 )
482 })?;
483
484 if buf.len() < len {
485 return Err(io::Error::new(
486 io::ErrorKind::UnexpectedEof,
487 "unexpected end of buffer reading varlen string",
488 ));
489 }
490
491 let bytes = &buf.0[..len];
492 buf.0 = &buf.0[len..];
493
494 let s = std::str::from_utf8(bytes)
495 .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, format!("invalid UTF-8: {}", e)))?;
496
497 Ok(Cow::Borrowed(s))
498}
499
500impl<'de> MyDeserialize<'de> for GtidEvent {
505 const SIZE: Option<usize> = None;
506 type Ctx = BinlogCtx<'de>;
507
508 fn deserialize(ctx: Self::Ctx, buf: &mut ParseBuf<'de>) -> io::Result<Self> {
509 let event_type = EventType::try_from(ctx.event_type_raw)
510 .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
511
512 match event_type {
513 EventType::GTID_EVENT | EventType::ANONYMOUS_GTID_EVENT => {
514 Self::deserialize_untagged(buf)
515 }
516 EventType::GTID_TAGGED_LOG_EVENT => Self::deserialize_tagged(buf),
517 _ => Err(io::Error::other("unexpected event type for GtidEvent")),
518 }
519 }
520}
521
522impl GtidEvent {
527 fn deserialize_untagged(buf: &mut ParseBuf<'_>) -> io::Result<Self> {
530 let mut sbuf: ParseBuf<'_> = buf.parse(1 + Self::ENCODED_SID_LENGTH + 8)?;
531 let flags = sbuf.parse_unchecked(())?;
532 let sid: [u8; Self::ENCODED_SID_LENGTH] = sbuf.parse_unchecked(())?;
533 let gno = sbuf.parse_unchecked(())?;
534
535 let mut lc_typecode = None;
536 let mut last_committed = RawInt::new(0);
537 let mut sequence_number = RawInt::new(0);
538 let mut immediate_commit_timestamp = RawInt::new(0);
539 let mut original_commit_timestamp = RawInt::new(0);
540 let mut tx_length = RawInt::new(0);
541
542 let mut original_server_version = RawInt::new(Self::UNDEFINED_SERVER_VERSION);
543 let mut immediate_server_version = RawInt::new(Self::UNDEFINED_SERVER_VERSION);
544
545 if !buf.is_empty() && buf.0[0] == Self::LOGICAL_TIMESTAMP_TYPECODE {
547 lc_typecode = Some(buf.parse_unchecked(())?);
548
549 let mut sbuf: ParseBuf<'_> = buf.parse(16)?;
550 last_committed = sbuf.parse_unchecked(())?;
551 sequence_number = sbuf.parse_unchecked(())?;
552
553 if buf.len() >= Self::IMMEDIATE_COMMIT_TIMESTAMP_LENGTH {
554 immediate_commit_timestamp = buf.parse_unchecked(())?;
555 if immediate_commit_timestamp.0 & (1 << 55) != 0 {
556 immediate_commit_timestamp.0 &= !(1 << 55);
557 original_commit_timestamp = buf.parse(())?;
558 } else {
559 original_commit_timestamp = immediate_commit_timestamp;
561 }
562 }
563
564 if !buf.is_empty() {
565 tx_length = buf.parse_unchecked(())?;
566 }
567
568 if buf.len() >= Self::IMMEDIATE_SERVER_VERSION_LENGTH {
569 immediate_server_version = buf.parse_unchecked(())?;
570 if immediate_server_version.0 & (1 << 31) != 0 {
571 immediate_server_version.0 &= !(1 << 31);
572 original_server_version = buf.parse(())?;
573 } else {
574 original_server_version = immediate_server_version;
575 }
576 }
577 }
578
579 Ok(Self {
580 flags,
581 sid,
582 gno,
583 tag: None,
584 lc_typecode,
585 last_committed,
586 sequence_number,
587 immediate_commit_timestamp,
588 original_commit_timestamp,
589 tx_length,
590 original_server_version,
591 immediate_server_version,
592 serialization_version: 0,
593 commit_group_ticket: Self::COMMIT_GROUP_TICKET_UNSET,
594 })
595 }
596
597 fn deserialize_tagged(buf: &mut ParseBuf<'_>) -> io::Result<Self> {
600 if buf.is_empty() {
601 return Err(io::Error::new(
602 io::ErrorKind::UnexpectedEof,
603 "unexpected end reading serialization version",
604 ));
605 }
606 let serialization_version = buf.0[0];
607 buf.0 = &buf.0[1..];
608
609 match serialization_version {
610 Self::TAGGED_SERIALIZATION_VERSION_V2 => {
611 Self::deserialize_tagged_v2(serialization_version, buf)
612 }
613 v => Err(io::Error::new(
614 io::ErrorKind::InvalidData,
615 format!("unsupported tagged GTID serialization version {v}"),
616 )),
617 }
618 }
619
620 fn deserialize_tagged_v2(
625 serialization_version: u8,
626 buf: &mut ParseBuf<'_>,
627 ) -> io::Result<Self> {
628 let buf_len_at_start = buf.len() + std::mem::size_of_val(&serialization_version);
632
633 let payload_size = read_varlen_uint(buf)?;
634 let last_non_ignorable_field_id = read_varlen_uint(buf)?;
635
636 let envelope_consumed = (buf_len_at_start - buf.len()) as u64;
641 if payload_size < envelope_consumed {
642 return Err(io::Error::new(
643 io::ErrorKind::InvalidData,
644 format!(
645 "payload_size ({}) is smaller than envelope overhead ({})",
646 payload_size, envelope_consumed,
647 ),
648 ));
649 }
650 let fields_len = (payload_size - envelope_consumed) as usize;
651 if buf.len() < fields_len {
652 return Err(io::Error::new(
653 io::ErrorKind::UnexpectedEof,
654 format!(
655 "buffer has {} bytes but payload declares {} field bytes",
656 buf.len(),
657 fields_len,
658 ),
659 ));
660 }
661 let mut payload_buf = ParseBuf(&buf.0[..fields_len]);
662 buf.0 = &buf.0[fields_len..];
663
664 let mut flags = RawFlags::new(0);
666 let mut sid = [0u8; 16];
667 let mut tag: Option<Tag<'static>> = None;
668 let mut gno = 0u64;
669 let mut last_committed = 0u64;
670 let mut sequence_number = 0u64;
671 let mut immediate_commit_timestamp = 0u64;
672 let mut original_commit_timestamp = 0u64;
673 let mut tx_length = 0u64;
674 let mut original_server_version = Self::UNDEFINED_SERVER_VERSION;
675 let mut immediate_server_version = Self::UNDEFINED_SERVER_VERSION;
676 let mut commit_group_ticket = Self::COMMIT_GROUP_TICKET_UNSET;
677 let mut seen_original_commit_timestamp = false;
678 let mut seen_original_server_version = false;
679
680 while !payload_buf.is_empty() {
682 let fid = read_varlen_uint(&mut payload_buf)?;
683
684 match fid {
685 field_id::GTID_FLAGS => {
686 let val = read_varlen_uint(&mut payload_buf)?;
687 if val > u8::MAX as u64 {
688 return Err(io::Error::new(
689 io::ErrorKind::InvalidData,
690 format!("gtid_flags value out of u8 range: {}", val),
691 ));
692 }
693 flags = RawFlags::new(val as u8);
694 }
695 field_id::SID => {
696 sid = read_serialized_uuid(&mut payload_buf)?;
697 }
698 field_id::GNO => {
699 let val = read_varlen_int(&mut payload_buf)?;
701 if val < 0 {
702 return Err(io::Error::new(
703 io::ErrorKind::InvalidData,
704 format!("negative gno value: {}", val),
705 ));
706 }
707 gno = val as u64;
708 }
709 field_id::TAG => {
710 let tag_str = read_varlen_string(&mut payload_buf)?;
711 tag = Some(
712 Tag::new(tag_str.into_owned())
713 .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?,
714 );
715 }
716 field_id::LAST_COMMITTED => {
717 let val = read_varlen_int(&mut payload_buf)?;
719 if val < 0 {
720 return Err(io::Error::new(
721 io::ErrorKind::InvalidData,
722 format!("negative last_committed value: {}", val),
723 ));
724 }
725 last_committed = val as u64;
726 }
727 field_id::SEQUENCE_NUMBER => {
728 let val = read_varlen_int(&mut payload_buf)?;
730 if val < 0 {
731 return Err(io::Error::new(
732 io::ErrorKind::InvalidData,
733 format!("negative sequence_number value: {}", val),
734 ));
735 }
736 sequence_number = val as u64;
737 }
738 field_id::IMMEDIATE_COMMIT_TIMESTAMP => {
739 immediate_commit_timestamp = read_varlen_uint(&mut payload_buf)?;
740 }
741 field_id::ORIGINAL_COMMIT_TIMESTAMP => {
742 original_commit_timestamp = read_varlen_uint(&mut payload_buf)?;
743 seen_original_commit_timestamp = true;
744 }
745 field_id::TRANSACTION_LENGTH => {
746 tx_length = read_varlen_uint(&mut payload_buf)?;
747 }
748 field_id::IMMEDIATE_SERVER_VERSION => {
749 let val = read_varlen_uint(&mut payload_buf)?;
750 if val > u32::MAX as u64 {
751 return Err(io::Error::new(
752 io::ErrorKind::InvalidData,
753 format!("immediate_server_version out of u32 range: {}", val),
754 ));
755 }
756 immediate_server_version = val as u32;
757 }
758 field_id::ORIGINAL_SERVER_VERSION => {
759 let val = read_varlen_uint(&mut payload_buf)?;
760 if val > u32::MAX as u64 {
761 return Err(io::Error::new(
762 io::ErrorKind::InvalidData,
763 format!("original_server_version out of u32 range: {}", val),
764 ));
765 }
766 original_server_version = val as u32;
767 seen_original_server_version = true;
768 }
769 field_id::COMMIT_GROUP_TICKET => {
770 commit_group_ticket = read_varlen_uint(&mut payload_buf)?;
771 }
772 _ => {
773 if fid <= last_non_ignorable_field_id {
776 return Err(io::Error::new(
779 io::ErrorKind::InvalidData,
780 format!(
781 "unknown non-ignorable field {} in GTID_TAGGED_LOG_EVENT \
782 (last_non_ignorable_field_id = {})",
783 fid, last_non_ignorable_field_id,
784 ),
785 ));
786 }
787 break;
790 }
791 }
792 }
793
794 if !seen_original_commit_timestamp {
796 original_commit_timestamp = immediate_commit_timestamp;
797 }
798 if !seen_original_server_version {
799 original_server_version = immediate_server_version;
800 }
801
802 let tag = tag.ok_or_else(|| {
804 io::Error::new(
805 io::ErrorKind::InvalidData,
806 "GTID_TAGGED_LOG_EVENT missing tag field",
807 )
808 })?;
809
810 Ok(Self {
811 flags,
812 sid,
813 gno: RawConst::new(gno),
814 tag: Some(tag),
815 lc_typecode: None,
816 last_committed: RawInt::new(last_committed),
817 sequence_number: RawInt::new(sequence_number),
818 immediate_commit_timestamp: RawInt::new(immediate_commit_timestamp),
819 original_commit_timestamp: RawInt::new(original_commit_timestamp),
820 tx_length: RawInt::new(tx_length),
821 original_server_version: RawInt::new(original_server_version),
822 immediate_server_version: RawInt::new(immediate_server_version),
823 serialization_version,
824 commit_group_ticket,
825 })
826 }
827}
828
829impl MySerialize for GtidEvent {
834 fn serialize(&self, buf: &mut Vec<u8>) {
835 if self.tag.is_some() {
836 self.serialize_tagged(buf);
837 } else {
838 self.serialize_untagged(buf);
839 }
840 }
841}
842
843impl GtidEvent {
844 fn serialize_untagged(&self, buf: &mut Vec<u8>) {
845 self.flags.serialize(&mut *buf);
846 self.sid.serialize(&mut *buf);
847 self.gno.serialize(&mut *buf);
848 match self.lc_typecode {
849 Some(lc_typecode) => lc_typecode.serialize(&mut *buf),
850 None => return,
851 };
852 self.last_committed.serialize(&mut *buf);
853 self.sequence_number.serialize(&mut *buf);
854
855 let mut immediate_commit_timestamp_with_flag = *self.immediate_commit_timestamp;
856 if self.immediate_commit_timestamp != self.original_commit_timestamp {
857 immediate_commit_timestamp_with_flag |= 1 << 55;
858 } else {
859 immediate_commit_timestamp_with_flag &= !(1 << 55);
860 }
861 RawInt::<LeU56>::new(immediate_commit_timestamp_with_flag).serialize(&mut *buf);
862
863 if self.immediate_commit_timestamp != self.original_commit_timestamp {
864 self.original_commit_timestamp.serialize(&mut *buf);
865 }
866
867 self.tx_length.serialize(&mut *buf);
868
869 let mut immediate_server_version_with_flag = *self.immediate_server_version;
870 if self.immediate_server_version != self.original_server_version {
871 immediate_server_version_with_flag |= 1 << 31;
872 } else {
873 immediate_server_version_with_flag &= !(1 << 31);
874 }
875 RawInt::<LeU32>::new(immediate_server_version_with_flag).serialize(&mut *buf);
876
877 if self.immediate_server_version != self.original_server_version {
878 self.original_server_version.serialize(&mut *buf);
879 }
880 }
881
882 fn serialize_tagged(&self, buf: &mut Vec<u8>) {
883 match self.serialization_version {
884 Self::TAGGED_SERIALIZATION_VERSION_V2 => self.serialize_tagged_v2(buf),
885 _ => unreachable!(
888 "unsupported tagged GTID serialization version {}",
889 self.serialization_version
890 ),
891 }
892 }
893
894 fn write_tagged_fields_v2(&self, fields: &mut Vec<u8>) {
896 write_varlen_uint(fields, field_id::GTID_FLAGS);
898 write_varlen_uint(fields, self.flags.0 as u64);
899
900 write_varlen_uint(fields, field_id::SID);
902 write_serialized_uuid(fields, &self.sid);
903
904 assert!(self.gno.0 <= i64::MAX as u64, "gno exceeds i64::MAX");
906 write_varlen_uint(fields, field_id::GNO);
907 write_varlen_int(fields, self.gno.0 as i64);
908
909 write_varlen_uint(fields, field_id::TAG);
911 let tag_bytes = self
912 .tag
913 .as_ref()
914 .map(|t| t.as_str().as_bytes())
915 .unwrap_or(b"");
916 write_varlen_uint(fields, tag_bytes.len() as u64);
917 fields.extend_from_slice(tag_bytes);
918
919 assert!(
921 self.last_committed.0 <= i64::MAX as u64,
922 "last_committed exceeds i64::MAX"
923 );
924 write_varlen_uint(fields, field_id::LAST_COMMITTED);
925 write_varlen_int(fields, self.last_committed.0 as i64);
926
927 assert!(
929 self.sequence_number.0 <= i64::MAX as u64,
930 "sequence_number exceeds i64::MAX"
931 );
932 write_varlen_uint(fields, field_id::SEQUENCE_NUMBER);
933 write_varlen_int(fields, self.sequence_number.0 as i64);
934
935 write_varlen_uint(fields, field_id::IMMEDIATE_COMMIT_TIMESTAMP);
937 write_varlen_uint(fields, self.immediate_commit_timestamp.0);
938
939 if self.original_commit_timestamp != self.immediate_commit_timestamp {
941 write_varlen_uint(fields, field_id::ORIGINAL_COMMIT_TIMESTAMP);
942 write_varlen_uint(fields, self.original_commit_timestamp.0);
943 }
944
945 write_varlen_uint(fields, field_id::TRANSACTION_LENGTH);
947 write_varlen_uint(fields, self.tx_length.0);
948
949 write_varlen_uint(fields, field_id::IMMEDIATE_SERVER_VERSION);
951 write_varlen_uint(fields, self.immediate_server_version.0 as u64);
952
953 if self.original_server_version != self.immediate_server_version {
955 write_varlen_uint(fields, field_id::ORIGINAL_SERVER_VERSION);
956 write_varlen_uint(fields, self.original_server_version.0 as u64);
957 }
958
959 if self.commit_group_ticket != Self::COMMIT_GROUP_TICKET_UNSET {
961 write_varlen_uint(fields, field_id::COMMIT_GROUP_TICKET);
962 write_varlen_uint(fields, self.commit_group_ticket);
963 }
964 }
965
966 fn serialize_tagged_v2(&self, buf: &mut Vec<u8>) {
967 let mut fields = Vec::new();
968 self.write_tagged_fields_v2(&mut fields);
969
970 let last_non_ignorable_field_id: u64 = 0;
973 let version_byte_size = std::mem::size_of_val(&self.serialization_version);
974 let payload_size = compute_self_inclusive_payload_size(
975 version_byte_size,
976 fields.len(),
977 last_non_ignorable_field_id,
978 );
979
980 buf.reserve(
982 1 + varlen_uint_size(payload_size)
983 + varlen_uint_size(last_non_ignorable_field_id)
984 + fields.len(),
985 );
986 buf.push(self.serialization_version);
987 write_varlen_uint(buf, payload_size);
988 write_varlen_uint(buf, last_non_ignorable_field_id);
989 buf.extend_from_slice(&fields);
990 }
991}
992
993impl<'a> BinlogStruct<'a> for GtidEvent {
998 fn len(&self, _version: BinlogVersion) -> usize {
999 if self.tag.is_some() {
1000 self.len_tagged()
1001 } else {
1002 self.len_untagged()
1003 }
1004 }
1005}
1006
1007impl GtidEvent {
1008 fn len_untagged(&self) -> usize {
1009 let mut len = S(0);
1010
1011 len += S(1); len += S(Self::ENCODED_SID_LENGTH); len += S(8); len += S(1); len += S(8); len += S(8); len += S(7); if self.immediate_commit_timestamp != self.original_commit_timestamp {
1021 len += S(7); }
1023
1024 len += S(crate::misc::lenenc_int_len(*self.tx_length) as usize); len += S(4); if self.immediate_server_version != self.original_server_version {
1027 len += S(4); }
1029
1030 min(len.0, u32::MAX as usize - BinlogEventHeader::LEN)
1031 }
1032
1033 fn len_tagged(&self) -> usize {
1034 match self.serialization_version {
1035 Self::TAGGED_SERIALIZATION_VERSION_V2 => self.len_tagged_v2(),
1036 _ => unreachable!(
1039 "unsupported tagged GTID serialization version {}",
1040 self.serialization_version
1041 ),
1042 }
1043 }
1044
1045 fn len_tagged_v2(&self) -> usize {
1046 let mut fields = Vec::new();
1047 self.write_tagged_fields_v2(&mut fields);
1048
1049 let last_non_ignorable_field_id: u64 = 0;
1050 let version_byte_size = std::mem::size_of_val(&self.serialization_version);
1051 let total_size = compute_self_inclusive_payload_size(
1052 version_byte_size,
1053 fields.len(),
1054 last_non_ignorable_field_id,
1055 ) as usize;
1056
1057 min(total_size, u32::MAX as usize - BinlogEventHeader::LEN)
1058 }
1059}
1060
1061impl<'a> BinlogEvent<'a> for GtidEvent {
1062 const EVENT_TYPE: EventType = EventType::GTID_EVENT;
1070}
1071
1072#[cfg(test)]
1077mod tests {
1078 use super::*;
1079
1080 #[test]
1081 fn varlen_uint_roundtrip() {
1082 let test_cases: &[u64] = &[
1083 0,
1084 1,
1085 63,
1086 64,
1087 127,
1088 128,
1089 255,
1090 256,
1091 0x3FFF,
1092 0x4000,
1093 0x1F_FFFF,
1094 0x20_0000,
1095 0x0FFF_FFFF,
1096 0x1000_0000,
1097 0x07_FFFF_FFFF,
1098 0x08_0000_0000,
1099 0x03FF_FFFF_FFFF,
1100 0x0400_0000_0000,
1101 0x01_FFFF_FFFF_FFFF,
1102 0x02_0000_0000_0000,
1103 0x00FF_FFFF_FFFF_FFFF,
1104 0x0100_0000_0000_0000,
1105 u64::MAX / 2,
1106 u64::MAX,
1107 ];
1108
1109 for &value in test_cases {
1110 let mut buf = Vec::new();
1111 write_varlen_uint(&mut buf, value);
1112
1113 let mut parse_buf = ParseBuf(&buf);
1114 let decoded = read_varlen_uint(&mut parse_buf).unwrap();
1115
1116 assert_eq!(
1117 value, decoded,
1118 "roundtrip failed for {} (0x{:X}), encoded {:?}",
1119 value, value, buf
1120 );
1121 assert!(
1122 parse_buf.is_empty(),
1123 "leftover bytes for {} (0x{:X})",
1124 value,
1125 value
1126 );
1127 }
1128 }
1129
1130 #[test]
1131 fn varlen_uint_byte_count() {
1132 let mut buf = Vec::new();
1134 write_varlen_uint(&mut buf, 0);
1135 assert_eq!(buf.len(), 1);
1136
1137 buf.clear();
1138 write_varlen_uint(&mut buf, 127);
1139 assert_eq!(buf.len(), 1);
1140
1141 buf.clear();
1143 write_varlen_uint(&mut buf, 128);
1144 assert_eq!(buf.len(), 2);
1145
1146 buf.clear();
1148 write_varlen_uint(&mut buf, u64::MAX);
1149 assert_eq!(buf.len(), 9);
1150 }
1151
1152 #[test]
1153 fn varlen_int_roundtrip() {
1154 let test_cases: &[i64] = &[
1155 0,
1156 1,
1157 -1,
1158 127,
1159 -128,
1160 i64::MAX,
1161 i64::MIN,
1162 i64::MAX / 2,
1163 i64::MIN / 2,
1164 42,
1165 -42,
1166 ];
1167
1168 for &value in test_cases {
1169 let mut buf = Vec::new();
1170 write_varlen_int(&mut buf, value);
1171
1172 let mut parse_buf = ParseBuf(&buf);
1173 let decoded = read_varlen_int(&mut parse_buf).unwrap();
1174
1175 assert_eq!(
1176 value, decoded,
1177 "signed roundtrip failed for {} (0x{:X})",
1178 value, value
1179 );
1180 assert!(parse_buf.is_empty(), "leftover bytes for signed {}", value);
1181 }
1182 }
1183
1184 #[test]
1185 fn tagged_event_creation() {
1186 let sid = [1u8; 16];
1187 let tag = Tag::new("domain_1").unwrap();
1188 let event = GtidEvent::new_tagged(sid, tag, 42);
1189
1190 assert!(event.is_tagged());
1191 assert_eq!(event.sid(), sid);
1192 assert_eq!(event.tag().unwrap().as_str(), "domain_1");
1193 assert_eq!(event.gno(), 42);
1194 assert_eq!(event.flags_raw(), 0);
1195 }
1196
1197 #[test]
1198 fn untagged_event_creation() {
1199 let sid = [1u8; 16];
1200 let event = GtidEvent::new(sid, 42);
1201
1202 assert!(!event.is_tagged());
1203 assert_eq!(event.sid(), sid);
1204 assert!(event.tag().is_none());
1205 assert_eq!(event.gno(), 42);
1206 }
1207
1208 #[test]
1209 fn tagged_event_serialize_deserialize_roundtrip() {
1210 let sid = *uuid::Uuid::parse_str("3E11FA47-71CA-11E1-9E33-C80AA9429562")
1211 .unwrap()
1212 .as_bytes();
1213
1214 let tag = Tag::new("domain_1").unwrap();
1215 let event = GtidEvent::new_tagged(sid, tag, 42)
1216 .with_last_committed(10)
1217 .with_sequence_number(11)
1218 .with_immediate_commit_timestamp(1234567890)
1219 .with_original_commit_timestamp(1234567890)
1220 .with_tx_length(256)
1221 .with_immediate_server_version(90200)
1222 .with_original_server_version(90200);
1223
1224 let mut buf = Vec::new();
1226 event.serialize(&mut buf);
1227
1228 let mut parse_buf = ParseBuf(&buf);
1230 let decoded = GtidEvent::deserialize_tagged(&mut parse_buf).unwrap();
1231
1232 assert!(decoded.is_tagged());
1233 assert_eq!(decoded.sid(), event.sid());
1234 assert_eq!(
1235 decoded.tag().unwrap().as_str(),
1236 event.tag().unwrap().as_str()
1237 );
1238 assert_eq!(decoded.gno(), event.gno());
1239 assert_eq!(decoded.last_committed(), event.last_committed());
1240 assert_eq!(decoded.sequence_number(), event.sequence_number());
1241 assert_eq!(
1242 decoded.immediate_commit_timestamp(),
1243 event.immediate_commit_timestamp()
1244 );
1245 assert_eq!(
1246 decoded.original_commit_timestamp(),
1247 event.original_commit_timestamp()
1248 );
1249 assert_eq!(decoded.tx_length(), event.tx_length());
1250 assert_eq!(
1251 decoded.immediate_server_version(),
1252 event.immediate_server_version()
1253 );
1254 assert_eq!(
1255 decoded.original_server_version(),
1256 event.original_server_version()
1257 );
1258 }
1259
1260 #[test]
1261 fn tagged_event_with_different_timestamps() {
1262 let sid = [0xABu8; 16];
1263 let tag = Tag::new("backup").unwrap();
1264 let event = GtidEvent::new_tagged(sid, tag, 1)
1265 .with_immediate_commit_timestamp(2000)
1266 .with_original_commit_timestamp(1000)
1267 .with_immediate_server_version(90200)
1268 .with_original_server_version(90100);
1269
1270 let mut buf = Vec::new();
1271 event.serialize(&mut buf);
1272
1273 let mut parse_buf = ParseBuf(&buf);
1274 let decoded = GtidEvent::deserialize_tagged(&mut parse_buf).unwrap();
1275
1276 assert_eq!(decoded.immediate_commit_timestamp(), 2000);
1277 assert_eq!(decoded.original_commit_timestamp(), 1000);
1278 assert_eq!(decoded.immediate_server_version(), 90200);
1279 assert_eq!(decoded.original_server_version(), 90100);
1280 }
1281
1282 #[test]
1283 fn tagged_event_roundtrip_preserves_explicit_zero_originals() {
1284 let sid = [0xCDu8; 16];
1285 let tag = Tag::new("roundtrip").unwrap();
1286
1287 let event = GtidEvent::new_tagged(sid, tag.clone(), 5)
1289 .with_immediate_commit_timestamp(99999)
1290 .with_original_commit_timestamp(0)
1291 .with_immediate_server_version(90200)
1292 .with_original_server_version(90200);
1293
1294 let mut buf = Vec::new();
1295 event.serialize(&mut buf);
1296
1297 let mut parse_buf = ParseBuf(&buf);
1298 let decoded = GtidEvent::deserialize_tagged(&mut parse_buf).unwrap();
1299
1300 assert_eq!(decoded.original_commit_timestamp(), 0);
1301 assert_eq!(decoded.immediate_commit_timestamp(), 99999);
1302
1303 let event = GtidEvent::new_tagged(sid, tag, 5)
1305 .with_immediate_commit_timestamp(50000)
1306 .with_original_commit_timestamp(50000)
1307 .with_immediate_server_version(90200)
1308 .with_original_server_version(GtidEvent::UNDEFINED_SERVER_VERSION);
1309
1310 buf.clear();
1311 event.serialize(&mut buf);
1312
1313 let mut parse_buf = ParseBuf(&buf);
1314 let decoded = GtidEvent::deserialize_tagged(&mut parse_buf).unwrap();
1315
1316 assert_eq!(
1317 decoded.original_server_version(),
1318 GtidEvent::UNDEFINED_SERVER_VERSION,
1319 );
1320 assert_eq!(decoded.immediate_server_version(), 90200);
1321 }
1322
1323 #[test]
1324 fn deserialize_rejects_unknown_non_ignorable_field() {
1325 let mut fields = Vec::new();
1326
1327 write_varlen_uint(&mut fields, field_id::GTID_FLAGS);
1329 write_varlen_uint(&mut fields, 0);
1330
1331 write_varlen_uint(&mut fields, 99);
1333 write_varlen_uint(&mut fields, 0); let last_non_ignorable: u64 = 100;
1336 let fields_size = fields.len() as u64;
1337 let payload_size = fields_size
1338 + varlen_uint_size(fields_size) as u64
1339 + varlen_uint_size(last_non_ignorable) as u64;
1340
1341 let mut buf = Vec::new();
1342 buf.push(GtidEvent::TAGGED_SERIALIZATION_VERSION_V2);
1343 write_varlen_uint(&mut buf, payload_size);
1344 write_varlen_uint(&mut buf, last_non_ignorable);
1345 buf.extend_from_slice(&fields);
1346
1347 let mut parse_buf = ParseBuf(&buf);
1348 let result = GtidEvent::deserialize_tagged(&mut parse_buf);
1349
1350 assert!(result.is_err());
1351 let err = result.unwrap_err();
1352 assert_eq!(err.kind(), io::ErrorKind::InvalidData);
1353 assert!(err.to_string().contains("unknown non-ignorable field"));
1354 }
1355
1356 #[test]
1357 fn deserialize_rejects_unknown_serialization_version() {
1358 let mut buf = Vec::new();
1359 buf.push(0u8); write_varlen_uint(&mut buf, 10); write_varlen_uint(&mut buf, 0); let mut parse_buf = ParseBuf(&buf);
1364 let result = GtidEvent::deserialize_tagged(&mut parse_buf);
1365
1366 assert!(result.is_err());
1367 let err = result.unwrap_err();
1368 assert_eq!(err.kind(), io::ErrorKind::InvalidData);
1369 assert!(
1370 err.to_string()
1371 .contains("unsupported tagged GTID serialization version"),
1372 );
1373 }
1374
1375 fn build_raw_tagged_payload(field_pairs: &[(u64, Vec<u8>)]) -> Vec<u8> {
1378 build_raw_tagged_payload_with_lnif(field_pairs, field_id::IMMEDIATE_SERVER_VERSION)
1379 }
1380
1381 fn build_raw_tagged_payload_with_lnif(
1382 field_pairs: &[(u64, Vec<u8>)],
1383 last_non_ignorable: u64,
1384 ) -> Vec<u8> {
1385 let mut fields = Vec::new();
1386 for (fid, data) in field_pairs {
1387 write_varlen_uint(&mut fields, *fid);
1388 fields.extend_from_slice(data);
1389 }
1390
1391 let version_byte: u8 = GtidEvent::TAGGED_SERIALIZATION_VERSION_V2;
1392 let version_byte_size = std::mem::size_of_val(&version_byte);
1393 let payload_size = compute_self_inclusive_payload_size(
1394 version_byte_size,
1395 fields.len(),
1396 last_non_ignorable,
1397 );
1398
1399 let mut buf = Vec::new();
1400 buf.push(version_byte); write_varlen_uint(&mut buf, payload_size);
1402 write_varlen_uint(&mut buf, last_non_ignorable);
1403 buf.extend_from_slice(&fields);
1404 buf
1405 }
1406
1407 fn encode_varlen_uint(value: u64) -> Vec<u8> {
1408 let mut buf = Vec::new();
1409 write_varlen_uint(&mut buf, value);
1410 buf
1411 }
1412
1413 fn encode_varlen_int(value: i64) -> Vec<u8> {
1414 let mut buf = Vec::new();
1415 write_varlen_int(&mut buf, value);
1416 buf
1417 }
1418
1419 fn encode_uuid(uuid: &[u8; 16]) -> Vec<u8> {
1420 let mut buf = Vec::new();
1421 write_serialized_uuid(&mut buf, uuid);
1422 buf
1423 }
1424
1425 fn encode_varlen_string(s: &str) -> Vec<u8> {
1426 let mut buf = Vec::new();
1427 write_varlen_uint(&mut buf, s.len() as u64);
1428 buf.extend_from_slice(s.as_bytes());
1429 buf
1430 }
1431
1432 #[test]
1433 fn deserialize_rejects_negative_signed_fields() {
1434 let sid = [0u8; 16];
1435
1436 let buf = build_raw_tagged_payload(&[
1438 (field_id::GTID_FLAGS, encode_varlen_uint(0)),
1439 (field_id::SID, encode_uuid(&sid)),
1440 (field_id::GNO, encode_varlen_int(-1)),
1441 (field_id::TAG, encode_varlen_string("test")),
1442 (field_id::LAST_COMMITTED, encode_varlen_int(0)),
1443 (field_id::SEQUENCE_NUMBER, encode_varlen_int(0)),
1444 (field_id::IMMEDIATE_COMMIT_TIMESTAMP, encode_varlen_uint(0)),
1445 (field_id::TRANSACTION_LENGTH, encode_varlen_uint(0)),
1446 (field_id::IMMEDIATE_SERVER_VERSION, encode_varlen_uint(0)),
1447 ]);
1448
1449 let mut parse_buf = ParseBuf(&buf);
1450 let err = GtidEvent::deserialize_tagged(&mut parse_buf).unwrap_err();
1451 assert!(err.to_string().contains("negative gno"), "{}", err);
1452
1453 let buf = build_raw_tagged_payload(&[
1455 (field_id::GTID_FLAGS, encode_varlen_uint(0)),
1456 (field_id::SID, encode_uuid(&sid)),
1457 (field_id::GNO, encode_varlen_int(1)),
1458 (field_id::TAG, encode_varlen_string("test")),
1459 (field_id::LAST_COMMITTED, encode_varlen_int(-5)),
1460 (field_id::SEQUENCE_NUMBER, encode_varlen_int(0)),
1461 (field_id::IMMEDIATE_COMMIT_TIMESTAMP, encode_varlen_uint(0)),
1462 (field_id::TRANSACTION_LENGTH, encode_varlen_uint(0)),
1463 (field_id::IMMEDIATE_SERVER_VERSION, encode_varlen_uint(0)),
1464 ]);
1465
1466 let mut parse_buf = ParseBuf(&buf);
1467 let err = GtidEvent::deserialize_tagged(&mut parse_buf).unwrap_err();
1468 assert!(
1469 err.to_string().contains("negative last_committed"),
1470 "{}",
1471 err
1472 );
1473
1474 let buf = build_raw_tagged_payload(&[
1476 (field_id::GTID_FLAGS, encode_varlen_uint(0)),
1477 (field_id::SID, encode_uuid(&sid)),
1478 (field_id::GNO, encode_varlen_int(1)),
1479 (field_id::TAG, encode_varlen_string("test")),
1480 (field_id::LAST_COMMITTED, encode_varlen_int(0)),
1481 (field_id::SEQUENCE_NUMBER, encode_varlen_int(-100)),
1482 (field_id::IMMEDIATE_COMMIT_TIMESTAMP, encode_varlen_uint(0)),
1483 (field_id::TRANSACTION_LENGTH, encode_varlen_uint(0)),
1484 (field_id::IMMEDIATE_SERVER_VERSION, encode_varlen_uint(0)),
1485 ]);
1486
1487 let mut parse_buf = ParseBuf(&buf);
1488 let err = GtidEvent::deserialize_tagged(&mut parse_buf).unwrap_err();
1489 assert!(
1490 err.to_string().contains("negative sequence_number"),
1491 "{}",
1492 err
1493 );
1494 }
1495
1496 #[test]
1497 fn deserialize_rejects_missing_tag() {
1498 let sid = [0u8; 16];
1499 let buf = build_raw_tagged_payload(&[
1501 (field_id::GTID_FLAGS, encode_varlen_uint(0)),
1502 (field_id::SID, encode_uuid(&sid)),
1503 (field_id::GNO, encode_varlen_int(1)),
1504 (field_id::LAST_COMMITTED, encode_varlen_int(0)),
1506 (field_id::SEQUENCE_NUMBER, encode_varlen_int(0)),
1507 (field_id::IMMEDIATE_COMMIT_TIMESTAMP, encode_varlen_uint(0)),
1508 (field_id::TRANSACTION_LENGTH, encode_varlen_uint(0)),
1509 (field_id::IMMEDIATE_SERVER_VERSION, encode_varlen_uint(0)),
1510 ]);
1511
1512 let mut parse_buf = ParseBuf(&buf);
1513 let err = GtidEvent::deserialize_tagged(&mut parse_buf).unwrap_err();
1514 assert!(err.to_string().contains("missing tag field"), "{}", err);
1515 }
1516
1517 #[test]
1518 fn deserialize_rejects_empty_buffer() {
1519 let buf: &[u8] = &[];
1520 let mut parse_buf = ParseBuf(buf);
1521 let err = GtidEvent::deserialize_tagged(&mut parse_buf).unwrap_err();
1522 assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1523 }
1524
1525 #[test]
1526 fn deserialize_rejects_truncated_varlen() {
1527 let buf: &[u8] = &[GtidEvent::TAGGED_SERIALIZATION_VERSION_V2, 0xFF, 0x01, 0x02];
1530 let mut parse_buf = ParseBuf(buf);
1531 let err = GtidEvent::deserialize_tagged(&mut parse_buf).unwrap_err();
1532 assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1533 }
1534
1535 #[test]
1536 fn tagged_event_roundtrip_with_max_gno() {
1537 let sid = [0xFFu8; 16];
1538 let tag = Tag::new("maxgno").unwrap();
1539 let event = GtidEvent::new_tagged(sid, tag, i64::MAX as u64)
1540 .with_last_committed(i64::MAX as u64)
1541 .with_sequence_number(i64::MAX as u64);
1542
1543 let mut buf = Vec::new();
1544 event.serialize(&mut buf);
1545
1546 let mut parse_buf = ParseBuf(&buf);
1547 let decoded = GtidEvent::deserialize_tagged(&mut parse_buf).unwrap();
1548
1549 assert_eq!(decoded.gno(), i64::MAX as u64);
1550 assert_eq!(decoded.last_committed(), i64::MAX as u64);
1551 assert_eq!(decoded.sequence_number(), i64::MAX as u64);
1552 }
1553}