From 65d855b193867343f473bdcd42e181df9a50c0fc Mon Sep 17 00:00:00 2001 From: Leszek Zalewski Date: Tue, 21 Jul 2026 21:34:09 +0200 Subject: [PATCH 1/6] Persist GTID coordinates in state and enable GTID resume Make Ghostferry's serialized state and state tracker mode-aware so a GTID-mode run can be interrupted and resumed, and route writer/verifier state updates through GTID coordinates. - SerializableState gains BinlogCoordinateMode plus GTID counterparts of the three binlog positions (pointer + omitempty, so file/position states serialize unchanged). StateTracker stores GTID coordinates alongside file/position and selects the read path by mode. - Add mode-aware coordinate accessors and a GTID-aware MinSourceBinlogCoordinate that returns the intersection of the writer and inline-verifier GTID sets (the safe resume floor). Add intersectGTIDSets since go-mysql has no intersection primitive. - DMLEventBase can carry explicit coordinates; the BinlogStreamer stamps GTID coordinates onto events in GTID mode so BinlogWriter, inline verifier, and target verifier advance GTID-based state. - Ferry resume now uses the mode-aware coordinates and no longer rejects GTID resume. - Tests: GTID state round-trip, file/position backward-compatibility, GTID intersection resume, and tracker update/serialize/resume. File/position behavior and on-disk format are unchanged. --- binlog_coordinate.go | 35 ++++++++ binlog_streamer.go | 22 +++++ binlog_writer.go | 2 +- dml_events.go | 26 ++++++ ferry.go | 17 ++-- inline_verifier.go | 2 +- state_tracker.go | 149 ++++++++++++++++++++++++++++++++-- target_verifier.go | 2 +- test/go/binlog_writer_test.go | 5 +- test/go/state_gtid_test.go | 135 ++++++++++++++++++++++++++++++ 10 files changed, 371 insertions(+), 24 deletions(-) create mode 100644 test/go/state_gtid_test.go diff --git a/binlog_coordinate.go b/binlog_coordinate.go index 59ee480a..66e9146f 100644 --- a/binlog_coordinate.go +++ b/binlog_coordinate.go @@ -212,6 +212,41 @@ func (c BinlogCoordinate) String() string { } } +// intersectGTIDSets returns the intersection of two MySQL GTID sets, i.e. the +// GTIDs present in both. go-mysql does not expose an intersection primitive, so +// this computes A ∩ B = A - (A - B). +// +// It operates on clones and does not mutate its inputs. If either input is not +// a *mysql.MysqlGTIDSet, or an operation fails, it falls back to a clone of a +// (the first argument) so callers still get a usable, conservative set. +func intersectGTIDSets(a, b mysql.GTIDSet) mysql.GTIDSet { + aMysql, aOK := a.(*mysql.MysqlGTIDSet) + bMysql, bOK := b.(*mysql.MysqlGTIDSet) + if !aOK || !bOK { + return a.Clone() + } + + // diff = A - B + diff, ok := aMysql.Clone().(*mysql.MysqlGTIDSet) + if !ok { + return a.Clone() + } + if err := diff.Minus(*bMysql); err != nil { + return a.Clone() + } + + // result = A - diff = A ∩ B + result, ok := aMysql.Clone().(*mysql.MysqlGTIDSet) + if !ok { + return a.Clone() + } + if err := result.Minus(*diff); err != nil { + return a.Clone() + } + + return result +} + // serializedBinlogCoordinate is the on-disk / on-wire shape of a // BinlogCoordinate. It is deliberately explicit and self-describing so that a // future GTID variant can be added as additional fields without breaking diff --git a/binlog_streamer.go b/binlog_streamer.go index 2fefa24b..559711a3 100644 --- a/binlog_streamer.go +++ b/binlog_streamer.go @@ -780,6 +780,28 @@ func (s *BinlogStreamer) handleRowsEvent(ev *replication.BinlogEvent, query []by return err } + // In GTID mode, stamp GTID coordinates onto the events so that downstream + // consumers (binlog writer, verifiers) advance GTID-based state rather than + // file/position. The resumable coordinate is the committed set BEFORE the + // current transaction, so an interruption replays the whole transaction. + if s.coordinateMode() == BinlogCoordinateGTID { + var currentCoord BinlogCoordinate + if s.lastStreamedGTIDSet != nil { + currentCoord = NewGTIDCoordinate(s.lastStreamedGTIDSet.String()) + } else { + currentCoord = NewGTIDCoordinate("") + } + var resumableCoord BinlogCoordinate + if s.lastResumableGTIDSet != nil { + resumableCoord = NewGTIDCoordinate(s.lastResumableGTIDSet.String()) + } else { + resumableCoord = NewGTIDCoordinate("") + } + for _, dmlEv := range dmlEvs { + dmlEv.SetCoordinates(currentCoord, resumableCoord) + } + } + events := make([]DMLEvent, 0) for _, dmlEv := range dmlEvs { diff --git a/binlog_writer.go b/binlog_writer.go index 6292700e..09528b74 100644 --- a/binlog_writer.go +++ b/binlog_writer.go @@ -123,7 +123,7 @@ func (b *BinlogWriter) writeEvents(events []DMLEvent) error { } if b.StateTracker != nil { - b.StateTracker.UpdateLastResumableSourceBinlogPosition(events[len(events)-1].ResumableBinlogPosition()) + b.StateTracker.UpdateLastResumableSourceBinlogCoordinate(events[len(events)-1].ResumableBinlogCoordinate()) } b.lastProcessedEventTime = events[len(events)-1].Timestamp() diff --git a/dml_events.go b/dml_events.go index d0a2df9c..c7e7a713 100644 --- a/dml_events.go +++ b/dml_events.go @@ -84,6 +84,9 @@ type DMLEvent interface { // forward-looking API used while Ghostferry migrates off raw mysql.Position. BinlogCoordinate() BinlogCoordinate ResumableBinlogCoordinate() BinlogCoordinate + // SetCoordinates lets the BinlogStreamer stamp non-file/position + // coordinates (e.g. GTID) onto the event. + SetCoordinates(coordinate, resumableCoordinate BinlogCoordinate) Annotation() (string, error) Timestamp() time.Time } @@ -95,6 +98,23 @@ type DMLEventBase struct { resumablePos mysql.Position query []byte timestamp time.Time + + // coordinate and resumableCoordinate optionally carry non-file/position + // coordinates (e.g. GTID). When set (non-nil), they take precedence over + // the file/position fields in the coordinate accessors. They are stamped by + // the BinlogStreamer in GTID mode. + coordinate *BinlogCoordinate + resumableCoordinate *BinlogCoordinate +} + +// SetCoordinates overrides the coordinate-typed accessors with explicit +// coordinates. It is used by the BinlogStreamer to stamp GTID coordinates onto +// events in GTID mode. The file/position accessors are unaffected. +func (e *DMLEventBase) SetCoordinates(coordinate, resumableCoordinate BinlogCoordinate) { + c := coordinate + r := resumableCoordinate + e.coordinate = &c + e.resumableCoordinate = &r } func (e *DMLEventBase) Database() string { @@ -118,10 +138,16 @@ func (e *DMLEventBase) ResumableBinlogPosition() mysql.Position { } func (e *DMLEventBase) BinlogCoordinate() BinlogCoordinate { + if e.coordinate != nil { + return *e.coordinate + } return NewFilePositionCoordinate(e.pos) } func (e *DMLEventBase) ResumableBinlogCoordinate() BinlogCoordinate { + if e.resumableCoordinate != nil { + return *e.resumableCoordinate + } return NewFilePositionCoordinate(e.resumablePos) } diff --git a/ferry.go b/ferry.go index 8087a91a..59bbefb0 100644 --- a/ferry.go +++ b/ferry.go @@ -605,13 +605,6 @@ func (f *Ferry) Start() error { // miss some records that are inserted between the time the // DataIterator determines the range of IDs to copy and the time that // the starting binlog coordinates are determined. - // Resume-from-state currently only supports file/position coordinates. - // GTID-based resume is introduced in a later stage once GTID coordinates are - // persisted in the serialized state. - if f.StateToResumeFrom != nil && f.Config.BinlogCoordinateMode == BinlogCoordinateGTID { - return fmt.Errorf("resuming from state is not yet supported in GTID binlog coordinate mode") - } - var sourceCoord BinlogCoordinate var targetCoord BinlogCoordinate @@ -626,10 +619,12 @@ func (f *Ferry) Start() error { } if !f.Config.SkipTargetVerification { - if f.StateToResumeFrom != nil && f.StateToResumeFrom.LastStoredBinlogPositionForTargetVerifier != zeroPosition { - targetCoord, err = f.TargetVerifier.BinlogStreamer.ConnectBinlogStreamerToMysqlSinceCoordinate( - NewFilePositionCoordinate(f.StateToResumeFrom.LastStoredBinlogPositionForTargetVerifier), - ) + var targetResumeCoord BinlogCoordinate + if f.StateToResumeFrom != nil { + targetResumeCoord = f.StateToResumeFrom.TargetVerifierBinlogCoordinate() + } + if f.StateToResumeFrom != nil && !targetResumeCoord.IsZero() { + targetCoord, err = f.TargetVerifier.BinlogStreamer.ConnectBinlogStreamerToMysqlSinceCoordinate(targetResumeCoord) } else { targetCoord, err = f.TargetVerifier.BinlogStreamer.ConnectBinlogStreamerToMysqlWithCoordinate() } diff --git a/inline_verifier.go b/inline_verifier.go index d69e48ff..909b3f90 100644 --- a/inline_verifier.go +++ b/inline_verifier.go @@ -750,7 +750,7 @@ func (v *InlineVerifier) binlogEventListener(evs []DMLEvent) error { } if v.StateTracker != nil { - v.StateTracker.UpdateLastResumableSourceBinlogPositionForInlineVerifier(evs[len(evs)-1].ResumableBinlogPosition()) + v.StateTracker.UpdateLastResumableSourceBinlogCoordinateForInlineVerifier(evs[len(evs)-1].ResumableBinlogCoordinate()) } return nil diff --git a/state_tracker.go b/state_tracker.go index 6f974cb8..8d2e1634 100644 --- a/state_tracker.go +++ b/state_tracker.go @@ -40,6 +40,17 @@ type SerializableState struct { BinlogVerifyStore BinlogVerifySerializedStore LastStoredBinlogPositionForInlineVerifier mysql.Position LastStoredBinlogPositionForTargetVerifier mysql.Position + + // BinlogCoordinateMode records which coordinate representation the state + // below was produced with. Empty means file/position (legacy states). + BinlogCoordinateMode BinlogCoordinateType `json:",omitempty"` + + // GTID counterparts of the three binlog positions above. They are only + // populated when BinlogCoordinateMode is "gtid". Pointers + omitempty keep + // file/position states byte-for-byte compatible with older versions. + LastWrittenBinlogCoordinate *BinlogCoordinate `json:",omitempty"` + LastStoredBinlogCoordinateForInlineVerifier *BinlogCoordinate `json:",omitempty"` + LastStoredBinlogCoordinateForTargetVerifier *BinlogCoordinate `json:",omitempty"` } func (s *SerializableState) MarshalJSON() ([]byte, error) { @@ -106,11 +117,81 @@ func (s *SerializableState) MinSourceBinlogPosition() mysql.Position { } } -// MinSourceBinlogCoordinate returns the safe source resume coordinate. It is -// the coordinate-typed counterpart of MinSourceBinlogPosition and currently -// delegates to it, wrapping the result as a file/position coordinate. +// coordinateMode returns the effective coordinate mode for this serialized +// state, treating the empty value as file/position for legacy states. +func (s *SerializableState) coordinateMode() BinlogCoordinateType { + if s.BinlogCoordinateMode == "" { + return BinlogCoordinateFilePosition + } + return s.BinlogCoordinateMode +} + +// WrittenSourceBinlogCoordinate returns the last written source coordinate as a +// BinlogCoordinate, matching the state's coordinate mode. +func (s *SerializableState) WrittenSourceBinlogCoordinate() BinlogCoordinate { + if s.coordinateMode() == BinlogCoordinateGTID && s.LastWrittenBinlogCoordinate != nil { + return *s.LastWrittenBinlogCoordinate + } + return NewFilePositionCoordinate(s.LastWrittenBinlogPosition) +} + +// InlineVerifierSourceBinlogCoordinate returns the inline verifier's stored +// source coordinate as a BinlogCoordinate, matching the state's coordinate mode. +func (s *SerializableState) InlineVerifierSourceBinlogCoordinate() BinlogCoordinate { + if s.coordinateMode() == BinlogCoordinateGTID && s.LastStoredBinlogCoordinateForInlineVerifier != nil { + return *s.LastStoredBinlogCoordinateForInlineVerifier + } + return NewFilePositionCoordinate(s.LastStoredBinlogPositionForInlineVerifier) +} + +// TargetVerifierBinlogCoordinate returns the target verifier's stored +// coordinate as a BinlogCoordinate, matching the state's coordinate mode. +func (s *SerializableState) TargetVerifierBinlogCoordinate() BinlogCoordinate { + if s.coordinateMode() == BinlogCoordinateGTID && s.LastStoredBinlogCoordinateForTargetVerifier != nil { + return *s.LastStoredBinlogCoordinateForTargetVerifier + } + return NewFilePositionCoordinate(s.LastStoredBinlogPositionForTargetVerifier) +} + +// MinSourceBinlogCoordinate returns the safe source resume coordinate. +// +// For file/position mode this is the coordinate-typed counterpart of +// MinSourceBinlogPosition (the earlier of the writer and inline-verifier +// positions). For GTID mode the safe resume point is the intersection of the +// writer and inline-verifier GTID sets, since resuming must not skip events +// either consumer had not yet durably processed. When only one side is present, +// that side is used. func (s *SerializableState) MinSourceBinlogCoordinate() BinlogCoordinate { - return NewFilePositionCoordinate(s.MinSourceBinlogPosition()) + if s.coordinateMode() != BinlogCoordinateGTID { + return NewFilePositionCoordinate(s.MinSourceBinlogPosition()) + } + + written := s.LastWrittenBinlogCoordinate + inline := s.LastStoredBinlogCoordinateForInlineVerifier + + if written == nil && inline == nil { + return NewGTIDCoordinate("") + } + if written == nil { + return *inline + } + if inline == nil { + return *written + } + + // Safe resume is the intersection: only GTIDs that both the writer and the + // inline verifier have durably processed can be skipped on resume. + writtenSet, err := written.ParsedGTIDSet() + if err != nil { + return *written + } + inlineSet, err := inline.ParsedGTIDSet() + if err != nil { + return *inline + } + + intersection := intersectGTIDSets(writtenSet, inlineSet) + return NewGTIDCoordinate(intersection.String()) } // For tracking the speed of the copy @@ -146,6 +227,13 @@ type StateTracker struct { lastStoredBinlogPositionForInlineVerifier mysql.Position lastStoredBinlogPositionForTargetVerifier mysql.Position + // GTID coordinates, only populated when the tracker operates in GTID mode. + // They are stored alongside (not instead of) the file/position fields so + // that switching the read path is a mode decision, not a data migration. + lastWrittenBinlogCoordinate *BinlogCoordinate + lastStoredBinlogCoordinateForInlineVerifier *BinlogCoordinate + lastStoredBinlogCoordinateForTargetVerifier *BinlogCoordinate + lastSuccessfulPaginationKeys map[string]PaginationKey completedTables map[string]bool @@ -176,6 +264,9 @@ func NewStateTrackerFromSerializedState(speedLogCount int, serializedState *Seri s.lastWrittenBinlogPosition = serializedState.LastWrittenBinlogPosition s.lastStoredBinlogPositionForInlineVerifier = serializedState.LastStoredBinlogPositionForInlineVerifier s.lastStoredBinlogPositionForTargetVerifier = serializedState.LastStoredBinlogPositionForTargetVerifier + s.lastWrittenBinlogCoordinate = serializedState.LastWrittenBinlogCoordinate + s.lastStoredBinlogCoordinateForInlineVerifier = serializedState.LastStoredBinlogCoordinateForInlineVerifier + s.lastStoredBinlogCoordinateForTargetVerifier = serializedState.LastStoredBinlogCoordinateForTargetVerifier return s } @@ -202,20 +293,41 @@ func (s *StateTracker) UpdateLastResumableBinlogPositionForTargetVerifier(pos my // Coordinate-based accessors and mutators. // -// These are the forward-looking API. They currently delegate to the -// file/position storage above so behavior is unchanged, but they let callers be -// written in terms of BinlogCoordinate ahead of introducing a GTID coordinate -// mode. +// These are the forward-looking API. For file/position coordinates they store +// into the existing file/position fields (unchanged behavior). For GTID +// coordinates they store into dedicated GTID fields, so both representations +// can coexist and the read path is selected by coordinate mode. func (s *StateTracker) UpdateLastResumableSourceBinlogCoordinate(coord BinlogCoordinate) { + if coord.IsGTID() { + s.BinlogRWMutex.Lock() + defer s.BinlogRWMutex.Unlock() + c := coord + s.lastWrittenBinlogCoordinate = &c + return + } s.UpdateLastResumableSourceBinlogPosition(coord.Position()) } func (s *StateTracker) UpdateLastResumableSourceBinlogCoordinateForInlineVerifier(coord BinlogCoordinate) { + if coord.IsGTID() { + s.BinlogRWMutex.Lock() + defer s.BinlogRWMutex.Unlock() + c := coord + s.lastStoredBinlogCoordinateForInlineVerifier = &c + return + } s.UpdateLastResumableSourceBinlogPositionForInlineVerifier(coord.Position()) } func (s *StateTracker) UpdateLastResumableBinlogCoordinateForTargetVerifier(coord BinlogCoordinate) { + if coord.IsGTID() { + s.BinlogRWMutex.Lock() + defer s.BinlogRWMutex.Unlock() + c := coord + s.lastStoredBinlogCoordinateForTargetVerifier = &c + return + } s.UpdateLastResumableBinlogPositionForTargetVerifier(coord.Position()) } @@ -223,6 +335,9 @@ func (s *StateTracker) LastResumableSourceBinlogCoordinate() BinlogCoordinate { s.BinlogRWMutex.RLock() defer s.BinlogRWMutex.RUnlock() + if s.lastWrittenBinlogCoordinate != nil { + return *s.lastWrittenBinlogCoordinate + } return NewFilePositionCoordinate(s.lastWrittenBinlogPosition) } @@ -230,6 +345,9 @@ func (s *StateTracker) LastResumableSourceBinlogCoordinateForInlineVerifier() Bi s.BinlogRWMutex.RLock() defer s.BinlogRWMutex.RUnlock() + if s.lastStoredBinlogCoordinateForInlineVerifier != nil { + return *s.lastStoredBinlogCoordinateForInlineVerifier + } return NewFilePositionCoordinate(s.lastStoredBinlogPositionForInlineVerifier) } @@ -237,6 +355,9 @@ func (s *StateTracker) LastResumableBinlogCoordinateForTargetVerifier() BinlogCo s.BinlogRWMutex.RLock() defer s.BinlogRWMutex.RUnlock() + if s.lastStoredBinlogCoordinateForTargetVerifier != nil { + return *s.lastStoredBinlogCoordinateForTargetVerifier + } return NewFilePositionCoordinate(s.lastStoredBinlogPositionForTargetVerifier) } @@ -367,6 +488,18 @@ func (s *StateTracker) Serialize(lastKnownTableSchemaCache TableSchemaCache, bin LastWrittenBinlogPosition: s.lastWrittenBinlogPosition, LastStoredBinlogPositionForInlineVerifier: s.lastStoredBinlogPositionForInlineVerifier, LastStoredBinlogPositionForTargetVerifier: s.lastStoredBinlogPositionForTargetVerifier, + + LastWrittenBinlogCoordinate: s.lastWrittenBinlogCoordinate, + LastStoredBinlogCoordinateForInlineVerifier: s.lastStoredBinlogCoordinateForInlineVerifier, + LastStoredBinlogCoordinateForTargetVerifier: s.lastStoredBinlogCoordinateForTargetVerifier, + } + + // If any GTID coordinate has been recorded, this state was produced in GTID + // mode. Marking the mode lets readers select the GTID fields on resume. + if s.lastWrittenBinlogCoordinate != nil || + s.lastStoredBinlogCoordinateForInlineVerifier != nil || + s.lastStoredBinlogCoordinateForTargetVerifier != nil { + state.BinlogCoordinateMode = BinlogCoordinateGTID } if binlogVerifyStore != nil { diff --git a/target_verifier.go b/target_verifier.go index 8e99b4d3..d00ebd26 100644 --- a/target_verifier.go +++ b/target_verifier.go @@ -46,7 +46,7 @@ func (t *TargetVerifier) BinlogEventListener(evs []DMLEvent) error { } if t.StateTracker != nil { - t.StateTracker.UpdateLastResumableBinlogPositionForTargetVerifier(evs[len(evs)-1].ResumableBinlogPosition()) + t.StateTracker.UpdateLastResumableBinlogCoordinateForTargetVerifier(evs[len(evs)-1].ResumableBinlogCoordinate()) } return nil diff --git a/test/go/binlog_writer_test.go b/test/go/binlog_writer_test.go index b2908eaa..902f2d49 100644 --- a/test/go/binlog_writer_test.go +++ b/test/go/binlog_writer_test.go @@ -30,8 +30,9 @@ func (e *stubDMLEvent) BinlogCoordinate() ghostferry.BinlogCoordinate { func (e *stubDMLEvent) ResumableBinlogCoordinate() ghostferry.BinlogCoordinate { return ghostferry.NewFilePositionCoordinate(mysql.Position{}) } -func (e *stubDMLEvent) Annotation() (string, error) { return "", nil } -func (e *stubDMLEvent) Timestamp() time.Time { return time.Time{} } +func (e *stubDMLEvent) SetCoordinates(_, _ ghostferry.BinlogCoordinate) {} +func (e *stubDMLEvent) Annotation() (string, error) { return "", nil } +func (e *stubDMLEvent) Timestamp() time.Time { return time.Time{} } // TestBinlogWriterBufferBinlogEventsBeforeRun verifies that BufferBinlogEvents // does not block when called before Run() has started in its own goroutine. diff --git a/test/go/state_gtid_test.go b/test/go/state_gtid_test.go new file mode 100644 index 00000000..d869b648 --- /dev/null +++ b/test/go/state_gtid_test.go @@ -0,0 +1,135 @@ +package test + +import ( + "encoding/json" + "testing" + + "github.com/Shopify/ghostferry" + "github.com/go-mysql-org/go-mysql/mysql" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +const ( + uuidA = "3e11fa47-71ca-11e1-9e33-c80aa9429562" + gtidWritten = uuidA + ":1-100" + gtidInline = uuidA + ":1-80" + gtidTarget = uuidA + ":1-90" + gtidExpected = uuidA + ":1-80" // intersection of written and inline +) + +func gtidCoord(s string) *ghostferry.BinlogCoordinate { + c := ghostferry.NewGTIDCoordinate(s) + return &c +} + +func TestSerializableState_GTIDRoundTrip(t *testing.T) { + state := &ghostferry.SerializableState{ + GhostferryVersion: "test-version", + LastSuccessfulPaginationKeys: map[string]ghostferry.PaginationKey{}, + CompletedTables: map[string]bool{}, + BinlogCoordinateMode: ghostferry.BinlogCoordinateGTID, + LastWrittenBinlogCoordinate: gtidCoord(gtidWritten), + LastStoredBinlogCoordinateForInlineVerifier: gtidCoord(gtidInline), + LastStoredBinlogCoordinateForTargetVerifier: gtidCoord(gtidTarget), + } + + data, err := json.Marshal(state) + require.NoError(t, err) + + var decoded ghostferry.SerializableState + require.NoError(t, json.Unmarshal(data, &decoded)) + + assert.Equal(t, ghostferry.BinlogCoordinateGTID, decoded.BinlogCoordinateMode) + require.NotNil(t, decoded.LastWrittenBinlogCoordinate) + assert.True(t, decoded.LastWrittenBinlogCoordinate.IsGTID()) + assert.Equal(t, gtidWritten, decoded.LastWrittenBinlogCoordinate.GTIDSet) + assert.Equal(t, gtidInline, decoded.LastStoredBinlogCoordinateForInlineVerifier.GTIDSet) + assert.Equal(t, gtidTarget, decoded.LastStoredBinlogCoordinateForTargetVerifier.GTIDSet) +} + +// TestSerializableState_FilePositionOmitsGTIDFields guards backward +// compatibility: a file/position state must not emit GTID fields or a mode. +func TestSerializableState_FilePositionOmitsGTIDFields(t *testing.T) { + state := &ghostferry.SerializableState{ + GhostferryVersion: "test-version", + LastSuccessfulPaginationKeys: map[string]ghostferry.PaginationKey{}, + CompletedTables: map[string]bool{}, + LastWrittenBinlogPosition: mysql.Position{Name: "mysql-bin.000001", Pos: 4}, + } + + data, err := json.Marshal(state) + require.NoError(t, err) + + str := string(data) + assert.NotContains(t, str, "BinlogCoordinateMode") + assert.NotContains(t, str, "LastWrittenBinlogCoordinate") + assert.NotContains(t, str, "GTIDSet") +} + +func TestSerializableState_MinSourceBinlogCoordinate_GTIDIntersection(t *testing.T) { + state := &ghostferry.SerializableState{ + BinlogCoordinateMode: ghostferry.BinlogCoordinateGTID, + LastWrittenBinlogCoordinate: gtidCoord(gtidWritten), + LastStoredBinlogCoordinateForInlineVerifier: gtidCoord(gtidInline), + } + + coord := state.MinSourceBinlogCoordinate() + assert.True(t, coord.IsGTID()) + + // The safe resume point is the intersection: the smaller of the two here. + got, err := coord.ParsedGTIDSet() + require.NoError(t, err) + want, err := mysql.ParseMysqlGTIDSet(gtidExpected) + require.NoError(t, err) + assert.True(t, got.Equal(want), "expected intersection %s, got %s", want.String(), got.String()) +} + +func TestSerializableState_MinSourceBinlogCoordinate_GTIDSingleSide(t *testing.T) { + state := &ghostferry.SerializableState{ + BinlogCoordinateMode: ghostferry.BinlogCoordinateGTID, + LastWrittenBinlogCoordinate: gtidCoord(gtidWritten), + } + + coord := state.MinSourceBinlogCoordinate() + assert.True(t, coord.IsGTID()) + assert.Equal(t, gtidWritten, coord.GTIDSet) +} + +func TestStateTracker_GTIDUpdateAndSerialize(t *testing.T) { + st := ghostferry.NewStateTracker(0) + + st.UpdateLastResumableSourceBinlogCoordinate(ghostferry.NewGTIDCoordinate(gtidWritten)) + st.UpdateLastResumableSourceBinlogCoordinateForInlineVerifier(ghostferry.NewGTIDCoordinate(gtidInline)) + st.UpdateLastResumableBinlogCoordinateForTargetVerifier(ghostferry.NewGTIDCoordinate(gtidTarget)) + + assert.Equal(t, gtidWritten, st.LastResumableSourceBinlogCoordinate().GTIDSet) + assert.Equal(t, gtidInline, st.LastResumableSourceBinlogCoordinateForInlineVerifier().GTIDSet) + assert.Equal(t, gtidTarget, st.LastResumableBinlogCoordinateForTargetVerifier().GTIDSet) + + state := st.Serialize(nil, nil) + assert.Equal(t, ghostferry.BinlogCoordinateGTID, state.BinlogCoordinateMode) + require.NotNil(t, state.LastWrittenBinlogCoordinate) + assert.Equal(t, gtidWritten, state.LastWrittenBinlogCoordinate.GTIDSet) + + // Round-trip through a new tracker resumed from the serialized state. + resumed := ghostferry.NewStateTrackerFromSerializedState(0, state) + assert.Equal(t, gtidWritten, resumed.LastResumableSourceBinlogCoordinate().GTIDSet) + assert.Equal(t, gtidTarget, resumed.LastResumableBinlogCoordinateForTargetVerifier().GTIDSet) +} + +func TestStateTracker_FilePositionUpdateStaysFilePosition(t *testing.T) { + st := ghostferry.NewStateTracker(0) + + st.UpdateLastResumableSourceBinlogCoordinate( + ghostferry.NewFilePositionCoordinate(mysql.Position{Name: "mysql-bin.000009", Pos: 42}), + ) + + coord := st.LastResumableSourceBinlogCoordinate() + assert.True(t, coord.IsFilePosition()) + assert.Equal(t, "mysql-bin.000009", coord.Position().Name) + + state := st.Serialize(nil, nil) + assert.Equal(t, ghostferry.BinlogCoordinateType(""), state.BinlogCoordinateMode) + assert.Nil(t, state.LastWrittenBinlogCoordinate) +} From 0e330bc07fdd9297907fabf2a75510e5abd437da Mon Sep 17 00:00:00 2001 From: Leszek Zalewski Date: Tue, 21 Jul 2026 21:34:30 +0200 Subject: [PATCH 2/6] Test both binlog coordinate modes in Ruby suite and CI Allow the integration suite to run against file/position or GTID coordinate mode, and exercise both in CI. - integrationferry and the Go test helper read GHOSTFERRY_BINLOG_COORDINATE_MODE so every test can run in either mode. - Ruby harness accepts a per-test :binlog_coordinate_mode and falls back to the suite-wide GHOSTFERRY_BINLOG_COORDINATE_MODE env var. - GitHub Actions ruby-test matrix runs file_position on 5.7/8.0/8.4 and gtid on 8.0/8.4 only (GTID mode is MySQL 8+). - test_helper gains binlog_coordinate_mode / gtid_coordinate_mode? helpers; dumped-state and progress assertions in callbacks and interrupt/resume tests are made mode-aware. - Progress reports mode-aware LastSuccessfulBinlogCoordinate and FinalBinlogCoordinate (alongside the legacy file/position fields), and BinlogStreamer exposes GetStopBinlogCoordinate, so GTID runs report coherent progress. Verified locally against MySQL 8.0: full Go and Ruby suites pass in both file_position and gtid modes. --- .github/workflows/tests.yml | 11 +++++ ferry.go | 2 + progress.go | 15 +++++-- test/helpers/ghostferry_helper.rb | 9 ++++ test/integration/callbacks_test.rb | 15 +++++-- test/integration/interrupt_resume_test.rb | 52 +++++++++++++++++------ test/lib/go/integrationferry/ferry.go | 7 +++ test/test_helper.rb | 17 +++++++- testhelpers/test_ferry.go | 6 +++ 9 files changed, 113 insertions(+), 21 deletions(-) diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml index 39cde138..a8b7d8a4 100644 --- a/.github/workflows/tests.yml +++ b/.github/workflows/tests.yml @@ -68,11 +68,21 @@ jobs: ruby-test: strategy: matrix: + # file_position is exercised on every supported MySQL version. mysql: ["5.7", "8.0", "8.4"] log_backend: ["logrus"] + binlog_coordinate_mode: ["file_position"] include: - mysql: "8.4" log_backend: "zerolog" + binlog_coordinate_mode: "file_position" + # GTID coordinate mode is only supported/tested on MySQL 8+. + - mysql: "8.0" + log_backend: "logrus" + binlog_coordinate_mode: "gtid" + - mysql: "8.4" + log_backend: "logrus" + binlog_coordinate_mode: "gtid" runs-on: ubuntu-latest timeout-minutes: 15 @@ -83,6 +93,7 @@ jobs: BUNDLE_WITHOUT: "development" MYSQL_VERSION: ${{ matrix.mysql }} GHOSTFERRY_LOG_BACKEND: ${{ matrix.log_backend }} + GHOSTFERRY_BINLOG_COORDINATE_MODE: ${{ matrix.binlog_coordinate_mode }} steps: - uses: actions/checkout@v6.0.2 diff --git a/ferry.go b/ferry.go index 59bbefb0..90288bc9 100644 --- a/ferry.go +++ b/ferry.go @@ -993,9 +993,11 @@ func (f *Ferry) Progress() *Progress { // Binlog Progress s.LastSuccessfulBinlogPos = f.BinlogStreamer.lastStreamedBinlogPosition + s.LastSuccessfulBinlogCoordinate = f.BinlogStreamer.GetLastStreamedBinlogCoordinate() s.BinlogStreamerLag = now.Sub(f.BinlogStreamer.lastProcessedEventTime).Seconds() s.BinlogWriterLag = now.Sub(f.BinlogWriter.lastProcessedEventTime).Seconds() s.FinalBinlogPos = f.BinlogStreamer.stopAtBinlogPosition + s.FinalBinlogCoordinate = f.BinlogStreamer.GetStopBinlogCoordinate() if f.TargetVerifier != nil { s.TargetBinlogStreamerLag = now.Sub(f.TargetVerifier.BinlogStreamer.lastProcessedEventTime).Seconds() diff --git a/progress.go b/progress.go index 85105452..db920417 100644 --- a/progress.go +++ b/progress.go @@ -32,9 +32,15 @@ type Progress struct { Tables map[string]TableProgress LastSuccessfulBinlogPos mysql.Position - BinlogStreamerLag float64 // This is the amount of seconds the binlog streamer is lagging by (seconds) - BinlogWriterLag float64 // This is the amount of seconds the binlog writer is lagging by (seconds) - Throttled bool + // LastSuccessfulBinlogCoordinate is the coordinate-typed counterpart of + // LastSuccessfulBinlogPos. It reflects the active BinlogCoordinateMode, so + // in GTID mode it carries the streamed GTID set. The legacy + // LastSuccessfulBinlogPos field is retained for backward compatibility and + // is only meaningful in file/position mode. + LastSuccessfulBinlogCoordinate BinlogCoordinate + BinlogStreamerLag float64 // This is the amount of seconds the binlog streamer is lagging by (seconds) + BinlogWriterLag float64 // This is the amount of seconds the binlog writer is lagging by (seconds) + Throttled bool // if the TargetVerifier is enabled, we emit this lag, otherwise this number will be 0 TargetBinlogStreamerLag float64 @@ -51,6 +57,9 @@ type Progress struct { // These are some variables that are only filled when CurrentState == done. FinalBinlogPos mysql.Position + // FinalBinlogCoordinate is the coordinate-typed counterpart of + // FinalBinlogPos, reflecting the active BinlogCoordinateMode. + FinalBinlogCoordinate BinlogCoordinate // A best estimate on the speed at which the copying is taking place. If // there are large gaps in the PaginationKey space, this probably will be inaccurate. diff --git a/test/helpers/ghostferry_helper.rb b/test/helpers/ghostferry_helper.rb index 5d2e6162..736a3a06 100644 --- a/test/helpers/ghostferry_helper.rb +++ b/test/helpers/ghostferry_helper.rb @@ -270,6 +270,15 @@ def start_ghostferry(resuming_state = nil) environment["GHOSTFERRY_MARGINALIA"] = @config[:marginalia] end + # Binlog coordinate mode: prefer an explicit per-test config, otherwise + # fall back to the suite-wide GHOSTFERRY_BINLOG_COORDINATE_MODE env var. + # This lets CI run the entire suite against "file_position" or "gtid" + # while individual tests can still pin a mode. + binlog_coordinate_mode = @config[:binlog_coordinate_mode] || ENV["GHOSTFERRY_BINLOG_COORDINATE_MODE"] + if binlog_coordinate_mode && !binlog_coordinate_mode.empty? + environment["GHOSTFERRY_BINLOG_COORDINATE_MODE"] = binlog_coordinate_mode + end + @logger.debug("starting ghostferry test binary #{@compiled_binary_path}") Open3.popen3(environment, @compiled_binary_path) do |stdin, stdout, stderr, wait_thr| stdin.puts(resuming_state) unless resuming_state.nil? diff --git a/test/integration/callbacks_test.rb b/test/integration/callbacks_test.rb index 4ce644fd..1f7ad8bb 100644 --- a/test/integration/callbacks_test.rb +++ b/test/integration/callbacks_test.rb @@ -32,10 +32,17 @@ def test_progress_callback assert_equal 0, progress.last["ActiveDataIterators"] - refute progress.last["LastSuccessfulBinlogPos"]["Name"].nil? - refute progress.last["LastSuccessfulBinlogPos"]["Pos"].nil? - assert progress.last["BinlogStreamerLag"] > 0 - assert_equal progress.last["LastSuccessfulBinlogPos"], progress.last["FinalBinlogPos"] + if gtid_coordinate_mode? + refute progress.last["LastSuccessfulBinlogCoordinate"].nil? + refute progress.last["LastSuccessfulBinlogCoordinate"]["GTIDSet"].nil? + assert progress.last["BinlogStreamerLag"] > 0 + assert_equal progress.last["LastSuccessfulBinlogCoordinate"], progress.last["FinalBinlogCoordinate"] + else + refute progress.last["LastSuccessfulBinlogPos"]["Name"].nil? + refute progress.last["LastSuccessfulBinlogPos"]["Pos"].nil? + assert progress.last["BinlogStreamerLag"] > 0 + assert_equal progress.last["LastSuccessfulBinlogPos"], progress.last["FinalBinlogPos"] + end assert progress.last["VerifierMessage"].include?("currentRowCount =") assert progress.last["VerifierMessage"].include?("currentEntryCount =") diff --git a/test/integration/interrupt_resume_test.rb b/test/integration/interrupt_resume_test.rb index 6a6e9033..709bed83 100644 --- a/test/integration/interrupt_resume_test.rb +++ b/test/integration/interrupt_resume_test.rb @@ -121,8 +121,13 @@ def test_interrupt_resume_will_not_emit_binlog_position_for_inline_verifier_if_n dumped_state = ghostferry.run_expecting_interrupt assert_basic_fields_exist_in_dumped_state(dumped_state) - assert_equal "", dumped_state["LastStoredBinlogPositionForInlineVerifier"]["Name"] - assert_equal 0, dumped_state["LastStoredBinlogPositionForInlineVerifier"]["Pos"] + if gtid_coordinate_mode? + # Without a verifier, no inline-verifier coordinate should be recorded. + assert_nil dumped_state["LastStoredBinlogCoordinateForInlineVerifier"] + else + assert_equal "", dumped_state["LastStoredBinlogPositionForInlineVerifier"]["Name"] + assert_equal 0, dumped_state["LastStoredBinlogPositionForInlineVerifier"]["Pos"] + end end def test_interrupt_resume_inline_verifier_with_datawriter @@ -350,17 +355,38 @@ def test_interrupt_resume_between_consecutive_rows_events dumped_state = ghostferry.run_expecting_interrupt - refute_nil dumped_state['LastWrittenBinlogPosition']['Name'] - refute_nil dumped_state['LastWrittenBinlogPosition']['Pos'] - refute_nil dumped_state['LastStoredBinlogPositionForInlineVerifier']['Name'] - refute_nil dumped_state['LastStoredBinlogPositionForInlineVerifier']['Pos'] - refute_nil dumped_state['LastStoredBinlogPositionForTargetVerifier']['Name'] - refute_nil dumped_state['LastStoredBinlogPositionForTargetVerifier']['Pos'] - - # assert the resumable position is not the start position - if dumped_state['LastWrittenBinlogPosition']['Name'] == start_binlog_status['File'] - refute_equal dumped_state['LastWrittenBinlogPosition']['Pos'], start_binlog_status['Position'] - refute_equal dumped_state['LastStoredBinlogPositionForInlineVerifier']['Pos'], start_binlog_status['Position'] + if gtid_coordinate_mode? + written = dumped_state['LastWrittenBinlogCoordinate'] + inline = dumped_state['LastStoredBinlogCoordinateForInlineVerifier'] + target = dumped_state['LastStoredBinlogCoordinateForTargetVerifier'] + + refute_nil written + refute_nil inline + refute_nil target + refute_nil written['GTIDSet'] + refute_nil inline['GTIDSet'] + refute_nil target['GTIDSet'] + + # The resumable GTID set is conservatively the committed set BEFORE the + # in-flight transaction, so when the interrupt lands inside the first + # transaction after start it can still equal the start set. The key + # invariants are that a non-empty GTID set was recorded and that resume + # (asserted below) succeeds. It must never go backwards from the start set. + refute_empty written['GTIDSet'] + refute_empty inline['GTIDSet'] + else + refute_nil dumped_state['LastWrittenBinlogPosition']['Name'] + refute_nil dumped_state['LastWrittenBinlogPosition']['Pos'] + refute_nil dumped_state['LastStoredBinlogPositionForInlineVerifier']['Name'] + refute_nil dumped_state['LastStoredBinlogPositionForInlineVerifier']['Pos'] + refute_nil dumped_state['LastStoredBinlogPositionForTargetVerifier']['Name'] + refute_nil dumped_state['LastStoredBinlogPositionForTargetVerifier']['Pos'] + + # assert the resumable position is not the start position + if dumped_state['LastWrittenBinlogPosition']['Name'] == start_binlog_status['File'] + refute_equal dumped_state['LastWrittenBinlogPosition']['Pos'], start_binlog_status['Position'] + refute_equal dumped_state['LastStoredBinlogPositionForInlineVerifier']['Pos'], start_binlog_status['Position'] + end end ghostferry = new_ghostferry(MINIMAL_GHOSTFERRY) diff --git a/test/lib/go/integrationferry/ferry.go b/test/lib/go/integrationferry/ferry.go index fe85962d..aa72b181 100644 --- a/test/lib/go/integrationferry/ferry.go +++ b/test/lib/go/integrationferry/ferry.go @@ -241,6 +241,13 @@ func NewStandardConfig() (*ghostferry.Config, error) { } } + // GHOSTFERRY_BINLOG_COORDINATE_MODE selects the binlog coordinate mode used + // by the integration run: "file_position" (default) or "gtid". It lets the + // Ruby suite run every test against either coordinate representation. + if binlogCoordinateMode := os.Getenv("GHOSTFERRY_BINLOG_COORDINATE_MODE"); binlogCoordinateMode != "" { + config.BinlogCoordinateMode = ghostferry.BinlogCoordinateType(binlogCoordinateMode) + } + verifierType := os.Getenv("GHOSTFERRY_VERIFIER_TYPE") if verifierType == ghostferry.VerifierTypeIterative { config.VerifierType = ghostferry.VerifierTypeIterative diff --git a/test/test_helper.rb b/test/test_helper.rb index bb752ded..545e3e76 100644 --- a/test/test_helper.rb +++ b/test/test_helper.rb @@ -203,13 +203,28 @@ def assert_uuid_table_is_identical # # To actually assert the validity of the data within the dumped state, you # have to do it manually. + # Returns the binlog coordinate mode the suite is currently running under. + # Defaults to "file_position" when unset. + def binlog_coordinate_mode + mode = ENV["GHOSTFERRY_BINLOG_COORDINATE_MODE"] + (mode.nil? || mode.empty?) ? "file_position" : mode + end + + def gtid_coordinate_mode? + binlog_coordinate_mode == "gtid" + end + def assert_basic_fields_exist_in_dumped_state(dumped_state) refute dumped_state.nil? refute dumped_state["GhostferryVersion"].nil? refute dumped_state["LastKnownTableSchemaCache"].nil? refute dumped_state["LastSuccessfulPaginationKeys"].nil? refute dumped_state["CompletedTables"].nil? - refute dumped_state["LastWrittenBinlogPosition"].nil? + if gtid_coordinate_mode? + refute dumped_state["LastWrittenBinlogCoordinate"].nil? + else + refute dumped_state["LastWrittenBinlogPosition"].nil? + end end def assert_ghostferry_completed(instance, times:) diff --git a/testhelpers/test_ferry.go b/testhelpers/test_ferry.go index 4e09516c..aae49aba 100644 --- a/testhelpers/test_ferry.go +++ b/testhelpers/test_ferry.go @@ -59,6 +59,12 @@ func NewTestConfig() *ghostferry.Config { }, } + // Allow the test suite to run against either binlog coordinate mode by + // setting GHOSTFERRY_BINLOG_COORDINATE_MODE. Defaults to file/position. + if mode := os.Getenv("GHOSTFERRY_BINLOG_COORDINATE_MODE"); mode != "" { + config.BinlogCoordinateMode = ghostferry.BinlogCoordinateType(mode) + } + err := config.ValidateConfig() PanicIfError(err) From fcf454f6a35bdb3437a6e40da4c6bf01e36f1534 Mon Sep 17 00:00:00 2001 From: Leszek Zalewski Date: Wed, 22 Jul 2026 20:50:24 +0200 Subject: [PATCH 3/6] Harden GTID resume: mode guard, presence checks, fail-closed intersection - Reject resume when the serialized state's coordinate mode differs from the configured mode, preventing cross-mode mis-routing. - Check target-verifier coordinate presence explicitly (HasTargetVerifierBinlogCoordinate) rather than via IsZero(), so an empty (but valid) GTID set on resume is not mistaken for "absent" and does not skip pre-interrupt target events. - Make intersectGTIDSets and MinSourceBinlogCoordinate fail closed: a parse/set-math error returns an error instead of silently falling back to one side, which could advance the resume floor and skip events. Ferry.Start propagates the error. - Remove SetCoordinates from the exported DMLEvent interface; stamp GTID coordinates via an internal coordinateStamper type assertion so external DMLEvent implementers are unaffected. Use the cached NewGTIDCoordinateFromSet when stamping. - CI: run go-test in gtid mode on MySQL 8.0/8.4 (file_position on all). - Tests: multi-UUID intersection, fail-closed parse error, target verifier presence with empty GTID set. --- .github/workflows/tests.yml | 9 ++++++ binlog_coordinate.go | 21 +++++++------ binlog_streamer.go | 21 +++++-------- dml_events.go | 28 ++++++++++------- ferry.go | 31 ++++++++++++++----- state_tracker.go | 47 +++++++++++++++++++++------- test/go/state_gtid_test.go | 62 +++++++++++++++++++++++++++++++++++-- 7 files changed, 165 insertions(+), 54 deletions(-) diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml index a8b7d8a4..3076c9a7 100644 --- a/.github/workflows/tests.yml +++ b/.github/workflows/tests.yml @@ -43,7 +43,15 @@ jobs: go-test: strategy: matrix: + # file_position is exercised on every supported MySQL version. mysql: ["5.7", "8.0", "8.4"] + binlog_coordinate_mode: ["file_position"] + include: + # GTID coordinate mode is only supported/tested on MySQL 8+. + - mysql: "8.0" + binlog_coordinate_mode: "gtid" + - mysql: "8.4" + binlog_coordinate_mode: "gtid" runs-on: ubuntu-latest timeout-minutes: 15 @@ -52,6 +60,7 @@ jobs: env: CI: "true" MYSQL_VERSION: ${{ matrix.mysql }} + GHOSTFERRY_BINLOG_COORDINATE_MODE: ${{ matrix.binlog_coordinate_mode }} steps: - uses: actions/checkout@v6.0.2 diff --git a/binlog_coordinate.go b/binlog_coordinate.go index 66e9146f..f3ece121 100644 --- a/binlog_coordinate.go +++ b/binlog_coordinate.go @@ -216,35 +216,36 @@ func (c BinlogCoordinate) String() string { // GTIDs present in both. go-mysql does not expose an intersection primitive, so // this computes A ∩ B = A - (A - B). // -// It operates on clones and does not mutate its inputs. If either input is not -// a *mysql.MysqlGTIDSet, or an operation fails, it falls back to a clone of a -// (the first argument) so callers still get a usable, conservative set. -func intersectGTIDSets(a, b mysql.GTIDSet) mysql.GTIDSet { +// It operates on clones and does not mutate its inputs. It is fail-closed: +// rather than silently falling back to one side (which could advance a resume +// floor past what a consumer durably processed and skip events), it returns an +// error so the caller can refuse to resume. +func intersectGTIDSets(a, b mysql.GTIDSet) (mysql.GTIDSet, error) { aMysql, aOK := a.(*mysql.MysqlGTIDSet) bMysql, bOK := b.(*mysql.MysqlGTIDSet) if !aOK || !bOK { - return a.Clone() + return nil, fmt.Errorf("GTID intersection requires MySQL GTID sets, got %T and %T", a, b) } // diff = A - B diff, ok := aMysql.Clone().(*mysql.MysqlGTIDSet) if !ok { - return a.Clone() + return nil, fmt.Errorf("GTID intersection: unexpected clone type %T", aMysql.Clone()) } if err := diff.Minus(*bMysql); err != nil { - return a.Clone() + return nil, fmt.Errorf("GTID intersection (A - B): %w", err) } // result = A - diff = A ∩ B result, ok := aMysql.Clone().(*mysql.MysqlGTIDSet) if !ok { - return a.Clone() + return nil, fmt.Errorf("GTID intersection: unexpected clone type %T", aMysql.Clone()) } if err := result.Minus(*diff); err != nil { - return a.Clone() + return nil, fmt.Errorf("GTID intersection (A - diff): %w", err) } - return result + return result, nil } // serializedBinlogCoordinate is the on-disk / on-wire shape of a diff --git a/binlog_streamer.go b/binlog_streamer.go index 559711a3..f20d0ddf 100644 --- a/binlog_streamer.go +++ b/binlog_streamer.go @@ -785,20 +785,15 @@ func (s *BinlogStreamer) handleRowsEvent(ev *replication.BinlogEvent, query []by // file/position. The resumable coordinate is the committed set BEFORE the // current transaction, so an interruption replays the whole transaction. if s.coordinateMode() == BinlogCoordinateGTID { - var currentCoord BinlogCoordinate - if s.lastStreamedGTIDSet != nil { - currentCoord = NewGTIDCoordinate(s.lastStreamedGTIDSet.String()) - } else { - currentCoord = NewGTIDCoordinate("") - } - var resumableCoord BinlogCoordinate - if s.lastResumableGTIDSet != nil { - resumableCoord = NewGTIDCoordinate(s.lastResumableGTIDSet.String()) - } else { - resumableCoord = NewGTIDCoordinate("") - } + currentCoord := NewGTIDCoordinateFromSet(s.lastStreamedGTIDSet) + resumableCoord := NewGTIDCoordinateFromSet(s.lastResumableGTIDSet) for _, dmlEv := range dmlEvs { - dmlEv.SetCoordinates(currentCoord, resumableCoord) + // SetCoordinates is an internal capability (coordinateStamper), not + // part of the exported DMLEvent interface; all built-in events + // satisfy it via DMLEventBase. + if stamper, ok := dmlEv.(coordinateStamper); ok { + stamper.SetCoordinates(currentCoord, resumableCoord) + } } } diff --git a/dml_events.go b/dml_events.go index c7e7a713..83e87f5e 100644 --- a/dml_events.go +++ b/dml_events.go @@ -42,14 +42,14 @@ type RowData []interface{} // https://github.com/Shopify/ghostferry/issues/165. // // In summary: -// - This code receives values from both go-sql-driver/mysql and -// go-mysql-org/go-mysql. -// - go-sql-driver/mysql gives us int64 for signed integer, and uint64 in a byte -// slice for unsigned integer. -// - go-mysql-org/go-mysql gives us int64 for signed integer, and uint64 for -// unsigned integer. -// - We currently make this function deal with both cases. In the future we can -// investigate alternative solutions. +// - This code receives values from both go-sql-driver/mysql and +// go-mysql-org/go-mysql. +// - go-sql-driver/mysql gives us int64 for signed integer, and uint64 in a byte +// slice for unsigned integer. +// - go-mysql-org/go-mysql gives us int64 for signed integer, and uint64 for +// unsigned integer. +// - We currently make this function deal with both cases. In the future we can +// investigate alternative solutions. func (r RowData) GetUint64(colIdx int) (uint64, error) { u64, ok := Uint64Value(r[colIdx]) if ok { @@ -84,13 +84,19 @@ type DMLEvent interface { // forward-looking API used while Ghostferry migrates off raw mysql.Position. BinlogCoordinate() BinlogCoordinate ResumableBinlogCoordinate() BinlogCoordinate - // SetCoordinates lets the BinlogStreamer stamp non-file/position - // coordinates (e.g. GTID) onto the event. - SetCoordinates(coordinate, resumableCoordinate BinlogCoordinate) Annotation() (string, error) Timestamp() time.Time } +// coordinateStamper is an internal capability used by the BinlogStreamer to +// stamp non-file/position coordinates (e.g. GTID) onto events. It is +// deliberately NOT part of the exported DMLEvent interface so that external +// implementations of DMLEvent do not have to implement it. All built-in DML +// events satisfy it via DMLEventBase. +type coordinateStamper interface { + SetCoordinates(coordinate, resumableCoordinate BinlogCoordinate) +} + // The base of DMLEvent to provide the necessary methods. type DMLEventBase struct { table *TableSchema diff --git a/ferry.go b/ferry.go index 90288bc9..ea8cf676 100644 --- a/ferry.go +++ b/ferry.go @@ -610,7 +610,22 @@ func (f *Ferry) Start() error { var err error if f.StateToResumeFrom != nil { - sourceCoord, err = f.BinlogStreamer.ConnectBinlogStreamerToMysqlSinceCoordinate(f.StateToResumeFrom.MinSourceBinlogCoordinate()) + // Fail fast if the serialized state was produced in a different binlog + // coordinate mode than the current run; resuming across modes would + // mis-route the streamer (e.g. seeking a GTID set against a + // file/position stream) and could silently skip or replay events. + if stateMode := f.StateToResumeFrom.CoordinateMode(); stateMode != f.Config.BinlogCoordinateMode { + return fmt.Errorf( + "cannot resume: state was saved in %q binlog coordinate mode but this run is configured for %q", + stateMode, f.Config.BinlogCoordinateMode, + ) + } + + resumeCoord, resumeErr := f.StateToResumeFrom.MinSourceBinlogCoordinate() + if resumeErr != nil { + return fmt.Errorf("computing resume coordinate: %w", resumeErr) + } + sourceCoord, err = f.BinlogStreamer.ConnectBinlogStreamerToMysqlSinceCoordinate(resumeCoord) } else { sourceCoord, err = f.BinlogStreamer.ConnectBinlogStreamerToMysqlWithCoordinate() } @@ -619,12 +634,14 @@ func (f *Ferry) Start() error { } if !f.Config.SkipTargetVerification { - var targetResumeCoord BinlogCoordinate - if f.StateToResumeFrom != nil { - targetResumeCoord = f.StateToResumeFrom.TargetVerifierBinlogCoordinate() - } - if f.StateToResumeFrom != nil && !targetResumeCoord.IsZero() { - targetCoord, err = f.TargetVerifier.BinlogStreamer.ConnectBinlogStreamerToMysqlSinceCoordinate(targetResumeCoord) + // Presence must be checked explicitly, not via IsZero(): in GTID mode an + // empty executed set is a valid target-verifier coordinate, and treating + // it as absent would reconnect the target verifier at the target's + // current position and skip events that occurred before the interrupt. + if f.StateToResumeFrom != nil && f.StateToResumeFrom.HasTargetVerifierBinlogCoordinate() { + targetCoord, err = f.TargetVerifier.BinlogStreamer.ConnectBinlogStreamerToMysqlSinceCoordinate( + f.StateToResumeFrom.TargetVerifierBinlogCoordinate(), + ) } else { targetCoord, err = f.TargetVerifier.BinlogStreamer.ConnectBinlogStreamerToMysqlWithCoordinate() } diff --git a/state_tracker.go b/state_tracker.go index 8d2e1634..507d71c7 100644 --- a/state_tracker.go +++ b/state_tracker.go @@ -3,6 +3,7 @@ package ghostferry import ( "container/ring" "encoding/json" + "fmt" "sync" "time" @@ -117,15 +118,32 @@ func (s *SerializableState) MinSourceBinlogPosition() mysql.Position { } } -// coordinateMode returns the effective coordinate mode for this serialized +// CoordinateMode returns the effective coordinate mode for this serialized // state, treating the empty value as file/position for legacy states. -func (s *SerializableState) coordinateMode() BinlogCoordinateType { +func (s *SerializableState) CoordinateMode() BinlogCoordinateType { if s.BinlogCoordinateMode == "" { return BinlogCoordinateFilePosition } return s.BinlogCoordinateMode } +// coordinateMode is the unexported alias kept for existing internal callers. +func (s *SerializableState) coordinateMode() BinlogCoordinateType { + return s.CoordinateMode() +} + +// HasTargetVerifierBinlogCoordinate reports whether the serialized state +// carries a target-verifier resume coordinate at all. This is distinct from +// whether that coordinate is the zero/empty value: in GTID mode an empty +// executed set is a valid coordinate, so presence must be checked via the +// pointer/field rather than IsZero(). +func (s *SerializableState) HasTargetVerifierBinlogCoordinate() bool { + if s.coordinateMode() == BinlogCoordinateGTID { + return s.LastStoredBinlogCoordinateForTargetVerifier != nil + } + return s.LastStoredBinlogPositionForTargetVerifier != (mysql.Position{}) +} + // WrittenSourceBinlogCoordinate returns the last written source coordinate as a // BinlogCoordinate, matching the state's coordinate mode. func (s *SerializableState) WrittenSourceBinlogCoordinate() BinlogCoordinate { @@ -161,37 +179,44 @@ func (s *SerializableState) TargetVerifierBinlogCoordinate() BinlogCoordinate { // writer and inline-verifier GTID sets, since resuming must not skip events // either consumer had not yet durably processed. When only one side is present, // that side is used. -func (s *SerializableState) MinSourceBinlogCoordinate() BinlogCoordinate { +// +// It is fail-closed for GTID: a parse or intersection failure returns an error +// rather than silently falling back to one side, because a wrong (too-advanced) +// resume floor would skip binlog events and corrupt the copy. +func (s *SerializableState) MinSourceBinlogCoordinate() (BinlogCoordinate, error) { if s.coordinateMode() != BinlogCoordinateGTID { - return NewFilePositionCoordinate(s.MinSourceBinlogPosition()) + return NewFilePositionCoordinate(s.MinSourceBinlogPosition()), nil } written := s.LastWrittenBinlogCoordinate inline := s.LastStoredBinlogCoordinateForInlineVerifier if written == nil && inline == nil { - return NewGTIDCoordinate("") + return NewGTIDCoordinate(""), nil } if written == nil { - return *inline + return *inline, nil } if inline == nil { - return *written + return *written, nil } // Safe resume is the intersection: only GTIDs that both the writer and the // inline verifier have durably processed can be skipped on resume. writtenSet, err := written.ParsedGTIDSet() if err != nil { - return *written + return BinlogCoordinate{}, fmt.Errorf("parsing written GTID coordinate %q: %w", written.GTIDSet, err) } inlineSet, err := inline.ParsedGTIDSet() if err != nil { - return *inline + return BinlogCoordinate{}, fmt.Errorf("parsing inline-verifier GTID coordinate %q: %w", inline.GTIDSet, err) } - intersection := intersectGTIDSets(writtenSet, inlineSet) - return NewGTIDCoordinate(intersection.String()) + intersection, err := intersectGTIDSets(writtenSet, inlineSet) + if err != nil { + return BinlogCoordinate{}, fmt.Errorf("computing safe GTID resume floor: %w", err) + } + return NewGTIDCoordinateFromSet(intersection), nil } // For tracking the speed of the copy diff --git a/test/go/state_gtid_test.go b/test/go/state_gtid_test.go index d869b648..e99d13a2 100644 --- a/test/go/state_gtid_test.go +++ b/test/go/state_gtid_test.go @@ -74,7 +74,8 @@ func TestSerializableState_MinSourceBinlogCoordinate_GTIDIntersection(t *testing LastStoredBinlogCoordinateForInlineVerifier: gtidCoord(gtidInline), } - coord := state.MinSourceBinlogCoordinate() + coord, err := state.MinSourceBinlogCoordinate() + require.NoError(t, err) assert.True(t, coord.IsGTID()) // The safe resume point is the intersection: the smaller of the two here. @@ -85,17 +86,74 @@ func TestSerializableState_MinSourceBinlogCoordinate_GTIDIntersection(t *testing assert.True(t, got.Equal(want), "expected intersection %s, got %s", want.String(), got.String()) } +// TestSerializableState_MinSourceBinlogCoordinate_MultiUUIDIntersection covers +// the intersection across two server UUIDs with non-contiguous ranges. +func TestSerializableState_MinSourceBinlogCoordinate_MultiUUIDIntersection(t *testing.T) { + uuidB := "8e12fa47-71ca-11e1-9e33-c80aa9429999" + written := uuidA + ":1-100," + uuidB + ":1-40" + inline := uuidA + ":1-70," + uuidB + ":1-60" + // Intersection = min ranges per UUID. + wantStr := uuidA + ":1-70," + uuidB + ":1-40" + + state := &ghostferry.SerializableState{ + BinlogCoordinateMode: ghostferry.BinlogCoordinateGTID, + LastWrittenBinlogCoordinate: gtidCoord(written), + LastStoredBinlogCoordinateForInlineVerifier: gtidCoord(inline), + } + + coord, err := state.MinSourceBinlogCoordinate() + require.NoError(t, err) + + got, err := coord.ParsedGTIDSet() + require.NoError(t, err) + want, err := mysql.ParseMysqlGTIDSet(wantStr) + require.NoError(t, err) + assert.True(t, got.Equal(want), "expected intersection %s, got %s", want.String(), got.String()) +} + +// TestSerializableState_MinSourceBinlogCoordinate_FailClosed verifies that an +// unparseable stored GTID coordinate produces an error rather than silently +// falling back to a possibly-too-advanced resume floor. +func TestSerializableState_MinSourceBinlogCoordinate_FailClosed(t *testing.T) { + state := &ghostferry.SerializableState{ + BinlogCoordinateMode: ghostferry.BinlogCoordinateGTID, + LastWrittenBinlogCoordinate: gtidCoord("not-a-valid-gtid-set"), + LastStoredBinlogCoordinateForInlineVerifier: gtidCoord(gtidInline), + } + + _, err := state.MinSourceBinlogCoordinate() + assert.Error(t, err) +} + func TestSerializableState_MinSourceBinlogCoordinate_GTIDSingleSide(t *testing.T) { state := &ghostferry.SerializableState{ BinlogCoordinateMode: ghostferry.BinlogCoordinateGTID, LastWrittenBinlogCoordinate: gtidCoord(gtidWritten), } - coord := state.MinSourceBinlogCoordinate() + coord, err := state.MinSourceBinlogCoordinate() + require.NoError(t, err) assert.True(t, coord.IsGTID()) assert.Equal(t, gtidWritten, coord.GTIDSet) } +func TestSerializableState_HasTargetVerifierBinlogCoordinate(t *testing.T) { + // GTID mode: empty-set coordinate is present (not absent). + emptyGTID := ghostferry.NewGTIDCoordinate("") + state := &ghostferry.SerializableState{ + BinlogCoordinateMode: ghostferry.BinlogCoordinateGTID, + LastStoredBinlogCoordinateForTargetVerifier: &emptyGTID, + } + assert.True(t, state.HasTargetVerifierBinlogCoordinate(), + "an empty GTID set coordinate must count as present") + + // GTID mode with no coordinate: absent. + stateNone := &ghostferry.SerializableState{ + BinlogCoordinateMode: ghostferry.BinlogCoordinateGTID, + } + assert.False(t, stateNone.HasTargetVerifierBinlogCoordinate()) +} + func TestStateTracker_GTIDUpdateAndSerialize(t *testing.T) { st := ghostferry.NewStateTracker(0) From 8b10648969217148d56c2e7b0a994230d5d1f102 Mon Sep 17 00:00:00 2001 From: Leszek Zalewski Date: Tue, 28 Jul 2026 17:09:22 +0200 Subject: [PATCH 4/6] Guard GTID coordinate stamping under gtidMu Stage 4 stamps GTID coordinates onto DML events in handleRowsEvent by reading lastStreamedGTIDSet / lastResumableGTIDSet directly. Route those reads through gtidMu-guarded snapshot accessors (streamedGTIDCoordinate / resumableGTIDCoordinate) so the stamping read cannot race the streaming goroutine's set mutations or Progress()'s concurrent reads, consistent with the GTID-set guarding introduced in the parent PR. --- binlog_streamer.go | 23 +++++++++++++++++++++-- 1 file changed, 21 insertions(+), 2 deletions(-) diff --git a/binlog_streamer.go b/binlog_streamer.go index f20d0ddf..232983df 100644 --- a/binlog_streamer.go +++ b/binlog_streamer.go @@ -574,6 +574,23 @@ func (s *BinlogStreamer) resumableGTIDClone() mysql.GTIDSet { return s.lastResumableGTIDSet.Clone() } +// streamedGTIDCoordinate and resumableGTIDCoordinate return GTID coordinates +// snapshotted under gtidMu, for stamping onto DML events. They read the sets +// while holding the lock so the snapshot never races the streaming goroutine's +// writes. NewGTIDCoordinateFromSet clones internally, so the returned +// coordinate does not alias the streamer's set. +func (s *BinlogStreamer) streamedGTIDCoordinate() BinlogCoordinate { + s.gtidMu.RLock() + defer s.gtidMu.RUnlock() + return NewGTIDCoordinateFromSet(s.lastStreamedGTIDSet) +} + +func (s *BinlogStreamer) resumableGTIDCoordinate() BinlogCoordinate { + s.gtidMu.RLock() + defer s.gtidMu.RUnlock() + return NewGTIDCoordinateFromSet(s.lastResumableGTIDSet) +} + // stopGTIDString returns the stop set as a string under the lock, or "" when // unset. func (s *BinlogStreamer) stopGTIDString() string { @@ -785,8 +802,10 @@ func (s *BinlogStreamer) handleRowsEvent(ev *replication.BinlogEvent, query []by // file/position. The resumable coordinate is the committed set BEFORE the // current transaction, so an interruption replays the whole transaction. if s.coordinateMode() == BinlogCoordinateGTID { - currentCoord := NewGTIDCoordinateFromSet(s.lastStreamedGTIDSet) - resumableCoord := NewGTIDCoordinateFromSet(s.lastResumableGTIDSet) + // Snapshot both coordinates under gtidMu so the read does not race the + // streaming goroutine's set mutations / Progress()'s reads. + currentCoord := s.streamedGTIDCoordinate() + resumableCoord := s.resumableGTIDCoordinate() for _, dmlEv := range dmlEvs { // SetCoordinates is an internal capability (coordinateStamper), not // part of the exported DMLEvent interface; all built-in events From a885722262217157a638e8a4b56b4546bf046ebb Mon Sep 17 00:00:00 2001 From: Leszek Zalewski Date: Mon, 21 Sep 2026 20:02:24 +0200 Subject: [PATCH 5/6] Fix GTID resume checkpoints during reconnect replay Exclude the incoming transaction from the committed GTID set when deriving the resume floor, preserving all other UUIDs and intervals. Stop streaming on invalid GTID identities even when Fatal returns. Add persisted-replay, subtraction-boundary, and fail-closed regressions. --- CHANGELOG.md | 1 + binlog_streamer.go | 55 +++++--- binlog_streamer_gtid_test.go | 235 ++++++++++++++++++++++++++++++++++- 3 files changed, 273 insertions(+), 18 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index cbf2d14f..ed64da1c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -25,6 +25,7 @@ All notable changes to this project will be documented in this file. - Keep GTID cutover aligned with transaction boundaries, including savepoints, DDL, empty transaction commits, and XA prepare/one-phase commit markers, without bypassing coordinate tracking when custom event handlers are registered. - Synchronize GTID coordinate snapshots and updates across goroutines. - Finish replayed GTID transactions before stopping at cutover, even when their GTIDs were already streamed before a reconnect. +- Exclude the current GTID from serialized resume checkpoints during reconnect replay, preserving other committed GTIDs. Stop streaming on invalid GTID identities before delivering further rows. - Publish binlog stop requests atomically after recording the stop coordinate, including reads from the control server. ## [1.3.1 - 2026-04-15] diff --git a/binlog_streamer.go b/binlog_streamer.go index 232983df..93180174 100644 --- a/binlog_streamer.go +++ b/binlog_streamer.go @@ -6,6 +6,7 @@ import ( sqlorig "database/sql" "errors" "fmt" + "math" "strings" "sync" "sync/atomic" @@ -15,6 +16,7 @@ import ( "github.com/go-mysql-org/go-mysql/mysql" "github.com/go-mysql-org/go-mysql/replication" + "github.com/google/uuid" ) const caughtUpThreshold = 10 * time.Second @@ -64,9 +66,9 @@ type BinlogStreamer struct { // GTID tracking, only maintained when BinlogCoordinateMode is // BinlogCoordinateGTID. lastStreamedGTIDSet is the committed GTID set seen - // so far. lastResumableGTIDSet is the committed GTID set at the last - // transaction boundary (a safe resume point). stopAtGTIDSet is the target - // executed set to stop at during cutover. + // so far. lastResumableGTIDSet excludes the current transaction from that + // committed set, even during replay. stopAtGTIDSet is the target executed + // set to stop at during cutover. // // These are *mysql.MysqlGTIDSet values (Go maps under the hood), mutated on // the streaming goroutine while Ferry.Progress() reads them (via the Get* @@ -467,7 +469,11 @@ func (s *BinlogStreamer) Run() { // Coordinate tracking must run even when a custom handler replaces the // default handler, but only after the event has been handled successfully. if err == nil && s.coordinateMode() == BinlogCoordinateGTID { - s.updateGTIDState(ev, &es) + if err := s.updateGTIDState(ev, &es); err != nil { + s.logger.WithError(err).Error("failed to update GTID state") + s.ErrorHandler.Fatal("binlog_streamer", err) + return + } } if es.isEventPositionValid { @@ -526,14 +532,32 @@ func (s *BinlogStreamer) setLastStreamedGTIDSet(set mysql.GTIDSet) { s.lastStreamedGTIDSet = set.Clone() } -// setResumableToStreamed records the current streamed set as the resumable -// point (the pre-transaction committed set), cloning under the lock. -func (s *BinlogStreamer) setResumableToStreamed() { +// setResumableBeforeGTID excludes the incoming transaction from the committed +// set so checkpoints remain safe even when go-mysql replays a committed GTID. +func (s *BinlogStreamer) setResumableBeforeGTID(event *replication.GTIDEvent) error { + sid, err := uuid.FromBytes(event.SID) + if err != nil { + return fmt.Errorf("invalid GTID SID: %w", err) + } + if event.GNO <= 0 || event.GNO == math.MaxInt64 { + return fmt.Errorf("invalid GTID sequence number %d", event.GNO) + } + s.gtidMu.Lock() defer s.gtidMu.Unlock() - if s.lastStreamedGTIDSet != nil { - s.lastResumableGTIDSet = s.lastStreamedGTIDSet.Clone() + if s.lastStreamedGTIDSet == nil { + s.lastResumableGTIDSet = &mysql.MysqlGTIDSet{Sets: make(map[string]*mysql.UUIDSet)} + return nil + } + committed, ok := s.lastStreamedGTIDSet.(*mysql.MysqlGTIDSet) + if !ok || committed == nil { + return fmt.Errorf("GTID resume floor requires a MySQL GTID set, got %T", s.lastStreamedGTIDSet) } + floor := committed.Clone().(*mysql.MysqlGTIDSet) + singleton := mysql.UUIDSet{SID: sid, Intervals: mysql.IntervalSlice{{Start: event.GNO, Stop: event.GNO + 1}}} + floor.MinusSet(&singleton) + s.lastResumableGTIDSet = floor + return nil } // seedGTIDSets initialises both the streamed and resumable sets to set. Used @@ -625,26 +649,26 @@ func (s *BinlogStreamer) GetStopBinlogCoordinate() BinlogCoordinate { // updateGTIDState tracks transaction boundaries independently of event handlers. // go-mysql's GSet includes the in-flight GTID even on BEGIN and SAVEPOINT, so // QueryEvents inside an explicit transaction must not advance the committed set. -func (s *BinlogStreamer) updateGTIDState(ev *replication.BinlogEvent, es *BinlogEventState) { +func (s *BinlogStreamer) updateGTIDState(ev *replication.BinlogEvent, es *BinlogEventState) error { // go-mysql decodes XA_PREPARE_LOG_EVENT as a GenericEvent. It terminates // the GTID group for both XA PREPARE and XA COMMIT ONE PHASE. XA END does // not; retain its GTID snapshot until this marker has been consumed. if ev.Header.EventType == replication.XA_PREPARE_LOG_EVENT { s.commitGTIDState(es.queryGTIDSet, es) - return + return nil } switch e := ev.Event.(type) { case *replication.GTIDEvent: es.transactionOpen = true es.inTransaction = false es.queryGTIDSet = nil - s.setResumableToStreamed() + return s.setResumableBeforeGTID(e) case *replication.QueryEvent: es.queryGTIDSet = e.GSet query := strings.TrimSpace(string(e.Query)) if strings.EqualFold(query, "BEGIN") || isXAStart(query) { es.inTransaction = true - return + return nil } if strings.EqualFold(query, "COMMIT") || strings.EqualFold(query, "ROLLBACK") { es.inTransaction = false @@ -655,6 +679,7 @@ func (s *BinlogStreamer) updateGTIDState(ev *replication.BinlogEvent, es *Binlog case *replication.XIDEvent: s.commitGTIDState(e.GSet, es) } + return nil } func isXAStart(query string) bool { @@ -799,8 +824,8 @@ func (s *BinlogStreamer) handleRowsEvent(ev *replication.BinlogEvent, query []by // In GTID mode, stamp GTID coordinates onto the events so that downstream // consumers (binlog writer, verifiers) advance GTID-based state rather than - // file/position. The resumable coordinate is the committed set BEFORE the - // current transaction, so an interruption replays the whole transaction. + // file/position. The resumable coordinate excludes the current transaction + // from the committed set, so an interruption replays the whole transaction. if s.coordinateMode() == BinlogCoordinateGTID { // Snapshot both coordinates under gtidMu so the read does not race the // streaming goroutine's set mutations / Progress()'s reads. diff --git a/binlog_streamer_gtid_test.go b/binlog_streamer_gtid_test.go index 2c6348c5..9993351e 100644 --- a/binlog_streamer_gtid_test.go +++ b/binlog_streamer_gtid_test.go @@ -1,11 +1,15 @@ package ghostferry import ( + "encoding/json" + "math" "sync" "testing" "github.com/go-mysql-org/go-mysql/mysql" "github.com/go-mysql-org/go-mysql/replication" + "github.com/go-mysql-org/go-mysql/schema" + "github.com/google/uuid" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) @@ -25,10 +29,11 @@ func mustParseGTID(t *testing.T, s string) mysql.GTIDSet { func TestRunStopsAtGTIDTransactionBoundary(t *testing.T) { target := mustParseGTID(t, gtidSetTarget) + sid := uuid.MustParse("3e11fa47-71ca-11e1-9e33-c80aa9429562") gtid := func() *replication.BinlogEvent { return &replication.BinlogEvent{ Header: &replication.EventHeader{EventType: replication.GTID_EVENT}, - Event: &replication.GTIDEvent{}, + Event: &replication.GTIDEvent{SID: sid[:], GNO: 100}, } } query := func(statement string) *replication.BinlogEvent { @@ -269,6 +274,8 @@ func TestGTIDCoordinateAccessorsAreRaceFree(t *testing.T) { s := &BinlogStreamer{BinlogCoordinateMode: BinlogCoordinateGTID} s.logger = LogWithField("tag", "test") s.seedGTIDSets(mustParseGTID(t, gtidSetLower)) + sid := uuid.MustParse("3e11fa47-71ca-11e1-9e33-c80aa9429562") + event := &replication.GTIDEvent{SID: sid[:], GNO: 100} const iterations = 500 var wg sync.WaitGroup @@ -278,7 +285,10 @@ func TestGTIDCoordinateAccessorsAreRaceFree(t *testing.T) { go func() { defer wg.Done() for i := 0; i < iterations; i++ { - s.setResumableToStreamed() + if err := s.setResumableBeforeGTID(event); err != nil { + t.Error(err) + return + } s.setLastStreamedGTIDSet(mustParseGTID(t, gtidSetTarget)) } }() @@ -307,6 +317,7 @@ func TestGTIDCoordinateAccessorsAreRaceFree(t *testing.T) { } func TestRunFinishesReplayedTransactionBeforeCutover(t *testing.T) { + sid := uuid.MustParse("3e11fa47-71ca-11e1-9e33-c80aa9429562") for _, stopAt := range []string{"GTID", "first replayed update"} { t.Run(stopAt, func(t *testing.T) { stream := replication.NewBinlogStreamer() @@ -346,7 +357,7 @@ func TestRunFinishesReplayedTransactionBeforeCutover(t *testing.T) { })) } for range 2 { - add(replication.GTID_EVENT, &replication.GTIDEvent{}) + add(replication.GTID_EVENT, &replication.GTIDEvent{SID: sid[:], GNO: 100}) add(replication.QUERY_EVENT, &replication.QueryEvent{Query: []byte("BEGIN"), GSet: mustParseGTID(t, gtidSetTarget)}) add(replication.UPDATE_ROWS_EVENTv2, &replication.RowsEvent{Rows: [][]interface{}{{0}, {1}}}) add(replication.UPDATE_ROWS_EVENTv2, &replication.RowsEvent{Rows: [][]interface{}{{1}, {0}}}) @@ -359,3 +370,221 @@ func TestRunFinishesReplayedTransactionBeforeCutover(t *testing.T) { }) } } + +func gtidReplayTableSchema(database string) TableSchemaCache { + return TableSchemaCache{ + database + ".t": &TableSchema{Table: &schema.Table{ + Schema: database, Name: "t", + Columns: []schema.TableColumn{{Name: "id", Type: schema.TYPE_NUMBER}, {Name: "v", Type: schema.TYPE_NUMBER}}, + }}, + } +} + +func queueGTIDReplayTransaction(t *testing.T, stream *replication.BinlogStreamer) { + t.Helper() + sid := uuid.MustParse("3e11fa47-71ca-11e1-9e33-c80aa9429562") + target := mustParseGTID(t, gtidSetTarget) + events := []*replication.BinlogEvent{ + {Header: &replication.EventHeader{EventType: replication.GTID_EVENT}, Event: &replication.GTIDEvent{SID: sid[:], GNO: 100}}, + {Header: &replication.EventHeader{EventType: replication.QUERY_EVENT}, Event: &replication.QueryEvent{Query: []byte("BEGIN"), GSet: target}}, + {Header: &replication.EventHeader{EventType: replication.UPDATE_ROWS_EVENTv2}, Event: &replication.RowsEvent{ + Table: &replication.TableMapEvent{Schema: []byte("review"), Table: []byte("t")}, Rows: [][]interface{}{{1, 0}, {1, 1}}, + }}, + {Header: &replication.EventHeader{EventType: replication.UPDATE_ROWS_EVENTv2}, Event: &replication.RowsEvent{ + Table: &replication.TableMapEvent{Schema: []byte("review"), Table: []byte("t")}, Rows: [][]interface{}{{1, 1}, {1, 0}}, + }}, + {Header: &replication.EventHeader{EventType: replication.XID_EVENT}, Event: &replication.XIDEvent{GSet: target}}, + } + for i, event := range events { + event.Header.LogPos = uint32(100 + i) + require.NoError(t, stream.AddEventToStreamer(event)) + } +} + +func TestRunResumesInterruptedGTIDReplay(t *testing.T) { + const start = "3e11fa47-71ca-11e1-9e33-c80aa9429562:1-99" + for _, inline := range []bool{false, true} { + name := "writer_only" + if inline { + name = "writer_and_inline" + } + t.Run(name, func(t *testing.T) { + newStreamer := func(coord BinlogCoordinate, groups int) *BinlogStreamer { + stream := replication.NewBinlogStreamer() + s := &BinlogStreamer{ + BinlogCoordinateMode: BinlogCoordinateGTID, + binlogStreamer: stream, + binlogSyncer: replication.NewBinlogSyncer(replication.BinlogSyncerConfig{ServerID: 1}), + TableSchema: gtidReplayTableSchema("review"), + } + s.seedGTIDSets(mustParseGTID(t, coord.GTIDSet)) + s.setStopGTIDSet(mustParseGTID(t, gtidSetTarget)) + for range groups { + queueGTIDReplayTransaction(t, stream) + } + require.NoError(t, s.AddBinlogEventHandler(replication.HEARTBEAT_EVENT, func(_ *replication.BinlogEvent, _ []byte, _ *BinlogEventState) ([]byte, error) { + t.Fatal("stream continued past the replayed commit") + return nil, nil + })) + require.NoError(t, stream.AddEventToStreamer(&replication.BinlogEvent{ + Header: &replication.EventHeader{EventType: replication.HEARTBEAT_EVENT}, Event: &replication.GenericEvent{}, + })) + return s + } + tracker := NewStateTracker(0) + advance := func(coord BinlogCoordinate) { + tracker.UpdateLastResumableSourceBinlogCoordinate(coord) + if inline { + tracker.UpdateLastResumableSourceBinlogCoordinateForInlineVerifier(coord) + } + } + advance(NewGTIDCoordinate(start)) + s := newStreamer(NewGTIDCoordinate(start), 2) + consumer, updates, savedConsumer := 0, 0, 0 + var saved []byte + apply := func(ev DMLEvent) { + if consumer == ev.OldValues()[1].(int) { + consumer = ev.NewValues()[1].(int) + } + advance(ev.ResumableBinlogCoordinate()) + updates++ + } + s.AddEventListener(func(events []DMLEvent) error { + for _, ev := range events { + apply(ev) + if updates == 3 { + assert.Equal(t, 1, consumer) + savedConsumer = consumer + var err error + saved, err = json.Marshal(tracker.Serialize(nil, nil)) + require.NoError(t, err) + s.stopRequested.Store(true) + } + } + return nil + }) + s.Run() + assert.Equal(t, 4, updates, "graceful cutover must finish the replay") + assert.Equal(t, 0, consumer) + + var state SerializableState + require.NoError(t, json.Unmarshal(saved, &state)) + tracker = NewStateTrackerFromSerializedState(0, &state) + resume, err := state.MinSourceBinlogCoordinate() + require.NoError(t, err) + assert.True(t, mustParseGTID(t, start).Equal(mustParseGTID(t, resume.GTIDSet)), "saved floor: %s", resume.GTIDSet) + + consumer, updates = savedConsumer, 0 + resumed := newStreamer(resume, 1) + resumed.stopRequested.Store(true) + resumed.AddEventListener(func(events []DMLEvent) error { + for _, ev := range events { + apply(ev) + } + return nil + }) + resumed.Run() + assert.Equal(t, 2, updates, "resume must deliver the interrupted transaction") + assert.Equal(t, 0, consumer, "resume must repair the saved intermediate row") + }) + } +} + +func TestGTIDResumeFloorExcludesOnlyCurrentTransaction(t *testing.T) { + const a = "3e11fa47-71ca-11e1-9e33-c80aa9429562" + const b = "8e12fa47-71ca-11e1-9e33-c80aa9429999" + sid := uuid.MustParse(a) + for _, tc := range []struct { + name, committed, want string + }{ + {"new transaction", a + ":1-99", a + ":1-99"}, + {"interior replay preserves later and unrelated GTIDs", a + ":1-150," + b + ":1-3:8-10", a + ":1-99:101-150," + b + ":1-3:8-10"}, + {"sole transaction", a + ":100", ""}, + {"nil committed progress", "", ""}, + } { + t.Run(tc.name, func(t *testing.T) { + s := &BinlogStreamer{BinlogCoordinateMode: BinlogCoordinateGTID} + // Start with an existing floor to ensure nil committed progress clears it. + s.seedGTIDSets(mustParseGTID(t, gtidSetTarget)) + var committed mysql.GTIDSet + if tc.committed != "" { + committed = mustParseGTID(t, tc.committed) + } + s.setLastStreamedGTIDSet(committed) + s.setStopGTIDSet(mustParseGTID(t, gtidSetPast)) + beforeCommitted := s.GetLastStreamedBinlogCoordinate() + beforeStop := s.GetStopBinlogCoordinate() + + require.NoError(t, s.setResumableBeforeGTID(&replication.GTIDEvent{SID: sid[:], GNO: 100})) + floor := s.resumableGTIDCoordinate() + assert.True(t, mustParseGTID(t, tc.want).Equal(mustParseGTID(t, floor.GTIDSet)), "resume floor: %s", floor.GTIDSet) + if tc.want == "" { + assert.Empty(t, floor.GTIDSet, "empty floor must not serialize a bare UUID") + } + assert.True(t, mustParseGTID(t, beforeCommitted.GTIDSet).Equal(mustParseGTID(t, s.GetLastStreamedBinlogCoordinate().GTIDSet))) + assert.True(t, mustParseGTID(t, beforeStop.GTIDSet).Equal(mustParseGTID(t, s.GetStopBinlogCoordinate().GTIDSet))) + }) + } +} + +type recordingGTIDErrorHandler struct { + err error +} + +func (h *recordingGTIDErrorHandler) ReportError(_ string, err error) { + h.err = err +} + +func (h *recordingGTIDErrorHandler) Fatal(from string, err error) { + h.ReportError(from, err) +} + +func TestRunRejectsInvalidGTIDBeforeRows(t *testing.T) { + sid := uuid.MustParse("3e11fa47-71ca-11e1-9e33-c80aa9429562") + for _, tc := range []struct { + name string + event *replication.GTIDEvent + }{ + {"malformed SID", &replication.GTIDEvent{SID: []byte{1}, GNO: 100}}, + {"zero sequence", &replication.GTIDEvent{SID: sid[:], GNO: 0}}, + {"overflowing sequence", &replication.GTIDEvent{SID: sid[:], GNO: math.MaxInt64}}, + } { + t.Run(tc.name, func(t *testing.T) { + stream := replication.NewBinlogStreamer() + handler := &recordingGTIDErrorHandler{} + s := &BinlogStreamer{ + BinlogCoordinateMode: BinlogCoordinateGTID, + binlogStreamer: stream, + binlogSyncer: replication.NewBinlogSyncer(replication.BinlogSyncerConfig{ServerID: 1}), + TableSchema: gtidReplayTableSchema("review"), + ErrorHandler: handler, + } + s.seedGTIDSets(mustParseGTID(t, gtidSetLower)) + beforeResumable := s.resumableGTIDCoordinate() + beforeCommitted := s.GetLastStreamedBinlogCoordinate() + rowsDelivered := false + s.AddEventListener(func(_ []DMLEvent) error { + rowsDelivered = true + return nil + }) + require.NoError(t, s.AddBinlogEventHandler(replication.HEARTBEAT_EVENT, func(_ *replication.BinlogEvent, _ []byte, _ *BinlogEventState) ([]byte, error) { + t.Fatal("stream continued after invalid GTID") + return nil, nil + })) + for _, event := range []*replication.BinlogEvent{ + {Header: &replication.EventHeader{EventType: replication.GTID_EVENT, LogPos: 100}, Event: tc.event}, + {Header: &replication.EventHeader{EventType: replication.UPDATE_ROWS_EVENTv2, LogPos: 101}, Event: &replication.RowsEvent{ + Table: &replication.TableMapEvent{Schema: []byte("review"), Table: []byte("t")}, Rows: [][]interface{}{{1, 0}, {1, 1}}, + }}, + {Header: &replication.EventHeader{EventType: replication.HEARTBEAT_EVENT}, Event: &replication.GenericEvent{}}, + } { + require.NoError(t, stream.AddEventToStreamer(event)) + } + s.Run() + assert.Error(t, handler.err) + assert.False(t, rowsDelivered) + assert.True(t, mustParseGTID(t, beforeResumable.GTIDSet).Equal(mustParseGTID(t, s.resumableGTIDCoordinate().GTIDSet))) + assert.True(t, mustParseGTID(t, beforeCommitted.GTIDSet).Equal(mustParseGTID(t, s.GetLastStreamedBinlogCoordinate().GTIDSet))) + }) + } +} From 335886bd5015a3567a23d7af67ca951e1de45393 Mon Sep 17 00:00:00 2001 From: Leszek Zalewski Date: Mon, 21 Sep 2026 23:11:22 +0200 Subject: [PATCH 6/6] Add google/uuid as direct dep --- go.mod | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/go.mod b/go.mod index 0dbdbfd0..bd8d0b00 100644 --- a/go.mod +++ b/go.mod @@ -8,6 +8,7 @@ require ( github.com/go-mysql-org/go-mysql v1.13.0 github.com/go-sql-driver/mysql v1.7.1 github.com/golang/snappy v0.0.0-20180518054509-2e65f85255db + github.com/google/uuid v1.3.0 github.com/gorilla/mux v1.6.1 github.com/rs/zerolog v1.35.0 github.com/shopspring/decimal v1.2.0 @@ -20,7 +21,6 @@ require ( github.com/Microsoft/go-winio v0.5.0 // indirect github.com/davecgh/go-spew v1.1.1 // indirect github.com/goccy/go-json v0.10.2 // indirect - github.com/google/uuid v1.3.0 // indirect github.com/gorilla/context v1.1.1 // indirect github.com/klauspost/compress v1.17.8 // indirect github.com/lann/builder v0.0.0-20180216234317-1b87b36280d0 // indirect