Skip to content

Commit d4f4bfe

Browse files
committed
fix: per-node origin name translation in table diff and merkle tree
NodeOriginNames was loaded from one random pool and used to translate roidents from all nodes. In native PG, roidents are node-local, so the same roident maps to different subscription names on different nodes. This caused intermittent preserve-origin repair failures when the diff happened to load origin names from the wrong node. Change NodeOriginNames to map[string]map[string]string (nodeName → roident → name). Load from every pool. Pass the originating node name into withSpockMetadata so each row's roident is resolved correctly.
1 parent 0157cc8 commit d4f4bfe

4 files changed

Lines changed: 105 additions & 43 deletions

File tree

internal/consistency/diff/table_diff.go

Lines changed: 42 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -102,7 +102,7 @@ type TableDiffTask struct {
102102
blockHashSQLCache map[hashBoundsKey]string
103103
blockHashSQLMu sync.Mutex
104104

105-
NodeOriginNames map[string]string
105+
NodeOriginNames map[string]map[string]string
106106

107107
CompareUnitSize int
108108
MaxDiffRows int64
@@ -204,53 +204,66 @@ func (t *TableDiffTask) loadNodeOriginNames() error {
204204
return nil
205205
}
206206

207-
var firstPool *pgxpool.Pool
208-
for _, pool := range t.Pools {
209-
firstPool = pool
210-
break
211-
}
207+
t.NodeOriginNames = make(map[string]map[string]string)
212208

213-
if firstPool == nil {
214-
t.NodeOriginNames = make(map[string]string)
215-
return fmt.Errorf("no connection pool available to load node origin names")
209+
var lastErr error
210+
for name, pool := range t.Pools {
211+
names, err := queries.GetNodeOriginNames(t.Ctx, pool)
212+
if err != nil {
213+
lastErr = err
214+
continue
215+
}
216+
t.NodeOriginNames[name] = names
216217
}
217218

218-
names, err := queries.GetNodeOriginNames(t.Ctx, firstPool)
219-
if err != nil {
220-
t.NodeOriginNames = make(map[string]string)
221-
return err
219+
if len(t.NodeOriginNames) == 0 && lastErr != nil {
220+
return lastErr
222221
}
223-
224-
t.NodeOriginNames = names
225222
return nil
226223
}
227224

225+
// flatNodeOriginNames merges all per-node origin maps into a single map for
226+
// lookup purposes (e.g. resolving --against-origin). If the same roident
227+
// appears on multiple nodes, the last one wins — this is acceptable because
228+
// resolveAgainstOrigin is only used with spock (where roidents are global)
229+
// or for user-facing name resolution where any match suffices.
230+
func (t *TableDiffTask) flatNodeOriginNames() map[string]string {
231+
flat := make(map[string]string)
232+
for _, nodeMap := range t.NodeOriginNames {
233+
for id, name := range nodeMap {
234+
flat[id] = name
235+
}
236+
}
237+
return flat
238+
}
239+
228240
func (t *TableDiffTask) resolveAgainstOrigin() error {
229241
if strings.TrimSpace(t.AgainstOrigin) == "" {
230242
return nil
231243
}
232-
if len(t.NodeOriginNames) == 0 {
244+
flat := t.flatNodeOriginNames()
245+
if len(flat) == 0 {
233246
return fmt.Errorf("unable to resolve --against-origin: no node origin names available")
234247
}
235248

236249
orig := strings.TrimSpace(t.AgainstOrigin)
237250
// direct match on origin id
238-
if _, ok := t.NodeOriginNames[orig]; ok {
251+
if _, ok := flat[orig]; ok {
239252
t.resolvedAgainstOrigin = orig
240253
return nil
241254
}
242255

243256
// match on origin name
244-
for id, name := range t.NodeOriginNames {
257+
for id, name := range flat {
245258
if name == orig {
246259
t.resolvedAgainstOrigin = id
247260
return nil
248261
}
249262
}
250263

251264
// build a list of available origins for the error message
252-
available := make([]string, 0, len(t.NodeOriginNames))
253-
for id, name := range t.NodeOriginNames {
265+
available := make([]string, 0, len(flat))
266+
for id, name := range flat {
254267
if id != name {
255268
available = append(available, fmt.Sprintf("%s (%s)", id, name))
256269
} else {
@@ -294,8 +307,8 @@ func (t *TableDiffTask) buildEffectiveFilter() (string, error) {
294307
return strings.Join(parts, " AND "), nil
295308
}
296309

297-
func (t *TableDiffTask) withSpockMetadata(row map[string]any) map[string]any {
298-
row["node_origin"] = utils.TranslateNodeOrigin(row["node_origin"], t.NodeOriginNames)
310+
func (t *TableDiffTask) withSpockMetadata(row map[string]any, nodeName string) map[string]any {
311+
row["node_origin"] = utils.TranslateNodeOrigin(row["node_origin"], t.NodeOriginNames[nodeName])
299312
return utils.AddSpockMetadata(row)
300313
}
301314

@@ -1375,8 +1388,9 @@ func (t *TableDiffTask) ExecuteTask() (err error) {
13751388
DiffRowsCount: make(map[string]int),
13761389
AgainstOrigin: t.AgainstOrigin,
13771390
AgainstOriginResolved: func() string {
1378-
if t.resolvedAgainstOrigin != "" && t.NodeOriginNames != nil {
1379-
if name, ok := t.NodeOriginNames[t.resolvedAgainstOrigin]; ok {
1391+
if t.resolvedAgainstOrigin != "" {
1392+
flat := t.flatNodeOriginNames()
1393+
if name, ok := flat[t.resolvedAgainstOrigin]; ok {
13801394
return name
13811395
}
13821396
}
@@ -1995,7 +2009,7 @@ func (t *TableDiffTask) recursiveDiff(
19952009
break
19962010
}
19972011
rowAsMap := utils.OrderedMapToMap(row)
1998-
rowWithMeta := t.withSpockMetadata(rowAsMap)
2012+
rowWithMeta := t.withSpockMetadata(rowAsMap, node1Name)
19992013
rowAsOrderedMap := utils.MapToOrderedMap(rowWithMeta, t.Cols)
20002014
t.DiffResult.NodeDiffs[pairKey].Rows[node1Name] = append(t.DiffResult.NodeDiffs[pairKey].Rows[node1Name], rowAsOrderedMap)
20012015
currentDiffRowsForPair++
@@ -2012,7 +2026,7 @@ func (t *TableDiffTask) recursiveDiff(
20122026
break
20132027
}
20142028
rowAsMap := utils.OrderedMapToMap(row)
2015-
rowWithMeta := t.withSpockMetadata(rowAsMap)
2029+
rowWithMeta := t.withSpockMetadata(rowAsMap, node2Name)
20162030
rowAsOrderedMap := utils.MapToOrderedMap(rowWithMeta, t.Cols)
20172031
t.DiffResult.NodeDiffs[pairKey].Rows[node2Name] = append(t.DiffResult.NodeDiffs[pairKey].Rows[node2Name], rowAsOrderedMap)
20182032
currentDiffRowsForPair++
@@ -2030,12 +2044,12 @@ func (t *TableDiffTask) recursiveDiff(
20302044
break
20312045
}
20322046
node1DataAsMap := utils.OrderedMapToMap(modRow.Node1Data)
2033-
node1DataWithMeta := t.withSpockMetadata(node1DataAsMap)
2047+
node1DataWithMeta := t.withSpockMetadata(node1DataAsMap, node1Name)
20342048
node1DataAsOrderedMap := utils.MapToOrderedMap(node1DataWithMeta, t.Cols)
20352049
t.DiffResult.NodeDiffs[pairKey].Rows[node1Name] = append(t.DiffResult.NodeDiffs[pairKey].Rows[node1Name], node1DataAsOrderedMap)
20362050

20372051
node2DataAsMap := utils.OrderedMapToMap(modRow.Node2Data)
2038-
node2DataWithMeta := t.withSpockMetadata(node2DataAsMap)
2052+
node2DataWithMeta := t.withSpockMetadata(node2DataAsMap, node2Name)
20392053
node2DataAsOrderedMap := utils.MapToOrderedMap(node2DataWithMeta, t.Cols)
20402054
t.DiffResult.NodeDiffs[pairKey].Rows[node2Name] = append(t.DiffResult.NodeDiffs[pairKey].Rows[node2Name], node2DataAsOrderedMap)
20412055
currentDiffRowsForPair++

internal/consistency/diff/table_diff_origin_test.go

Lines changed: 50 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -40,7 +40,7 @@ func TestResolveAgainstOrigin_NoNodeOriginNames(t *testing.T) {
4040
func TestResolveAgainstOrigin_MatchByID(t *testing.T) {
4141
task := &TableDiffTask{
4242
AgainstOrigin: "3",
43-
NodeOriginNames: map[string]string{"3": "n1", "4": "n2"},
43+
NodeOriginNames: map[string]map[string]string{"node1": {"3": "n1", "4": "n2"}},
4444
}
4545
if err := task.resolveAgainstOrigin(); err != nil {
4646
t.Fatalf("unexpected error: %v", err)
@@ -53,7 +53,7 @@ func TestResolveAgainstOrigin_MatchByID(t *testing.T) {
5353
func TestResolveAgainstOrigin_MatchByName_SpockNodeName(t *testing.T) {
5454
task := &TableDiffTask{
5555
AgainstOrigin: "n1",
56-
NodeOriginNames: map[string]string{"3": "n1", "4": "n2"},
56+
NodeOriginNames: map[string]map[string]string{"node1": {"3": "n1", "4": "n2"}},
5757
}
5858
if err := task.resolveAgainstOrigin(); err != nil {
5959
t.Fatalf("unexpected error: %v", err)
@@ -64,10 +64,12 @@ func TestResolveAgainstOrigin_MatchByName_SpockNodeName(t *testing.T) {
6464
}
6565

6666
func TestResolveAgainstOrigin_MatchByName_SubscriptionName(t *testing.T) {
67-
// Native PG: NodeOriginNames maps roident -> subscription name
67+
// Native PG: NodeOriginNames maps roident -> subscription name, per node
6868
task := &TableDiffTask{
69-
AgainstOrigin: "sub_n1_to_n2",
70-
NodeOriginNames: map[string]string{"5": "sub_n1_to_n2", "6": "sub_n3_to_n2"},
69+
AgainstOrigin: "sub_n1_to_n2",
70+
NodeOriginNames: map[string]map[string]string{
71+
"node1": {"5": "sub_n1_to_n2", "6": "sub_n3_to_n2"},
72+
},
7173
}
7274
if err := task.resolveAgainstOrigin(); err != nil {
7375
t.Fatalf("unexpected error: %v", err)
@@ -80,7 +82,7 @@ func TestResolveAgainstOrigin_MatchByName_SubscriptionName(t *testing.T) {
8082
func TestResolveAgainstOrigin_NoMatch(t *testing.T) {
8183
task := &TableDiffTask{
8284
AgainstOrigin: "nonexistent",
83-
NodeOriginNames: map[string]string{"3": "n1", "4": "n2"},
85+
NodeOriginNames: map[string]map[string]string{"node1": {"3": "n1", "4": "n2"}},
8486
}
8587
err := task.resolveAgainstOrigin()
8688
if err == nil {
@@ -133,3 +135,45 @@ func TestBuildEffectiveFilter_Empty(t *testing.T) {
133135
t.Fatalf("expected empty filter, got: %s", filter)
134136
}
135137
}
138+
139+
func TestWithSpockMetadata_PerNodeTranslation(t *testing.T) {
140+
// Simulate native PG: same roident "1" maps to different names on different nodes
141+
task := &TableDiffTask{
142+
NodeOriginNames: map[string]map[string]string{
143+
"n1": {"1": "sub_n3_to_n1"},
144+
"n2": {"1": "sub_n3_to_n2"},
145+
},
146+
}
147+
148+
row1 := map[string]any{"node_origin": "1", "id": 1}
149+
result1 := task.withSpockMetadata(row1, "n1")
150+
meta1 := result1["_spock_metadata_"].(map[string]any)
151+
if meta1["node_origin"] != "sub_n3_to_n1" {
152+
t.Fatalf("n1 row: expected origin sub_n3_to_n1, got %v", meta1["node_origin"])
153+
}
154+
155+
row2 := map[string]any{"node_origin": "1", "id": 2}
156+
result2 := task.withSpockMetadata(row2, "n2")
157+
meta2 := result2["_spock_metadata_"].(map[string]any)
158+
if meta2["node_origin"] != "sub_n3_to_n2" {
159+
t.Fatalf("n2 row: expected origin sub_n3_to_n2, got %v", meta2["node_origin"])
160+
}
161+
}
162+
163+
func TestFlatNodeOriginNames_MergesAllNodes(t *testing.T) {
164+
task := &TableDiffTask{
165+
NodeOriginNames: map[string]map[string]string{
166+
"n1": {"1": "sub_n3_to_n1"},
167+
"n2": {"1": "sub_n3_to_n2", "2": "sub_n4_to_n2"},
168+
},
169+
}
170+
flat := task.flatNodeOriginNames()
171+
// roident "2" should always be present
172+
if flat["2"] != "sub_n4_to_n2" {
173+
t.Fatalf("expected flat[2]=sub_n4_to_n2, got %q", flat["2"])
174+
}
175+
// roident "1" will be one of the two — either is valid for flattened lookup
176+
if flat["1"] != "sub_n3_to_n1" && flat["1"] != "sub_n3_to_n2" {
177+
t.Fatalf("expected flat[1] to be one of the subscription names, got %q", flat["1"])
178+
}
179+
}

internal/consistency/diff/table_rerun.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -374,12 +374,12 @@ func (t *TableDiffTask) reCompareDiffs(fetchedRowsByNode map[string]map[string]t
374374
persistentDiffCount++
375375
if nowOnNode1 {
376376
rowAsMap := utils.OrderedMapToMap(newRow1)
377-
rowWithMeta := t.withSpockMetadata(rowAsMap)
377+
rowWithMeta := t.withSpockMetadata(rowAsMap, node1)
378378
newDiffsForPair.Rows[node1] = append(newDiffsForPair.Rows[node1], utils.MapToOrderedMap(rowWithMeta, t.Cols))
379379
}
380380
if nowOnNode2 {
381381
rowAsMap := utils.OrderedMapToMap(newRow2)
382-
rowWithMeta := t.withSpockMetadata(rowAsMap)
382+
rowWithMeta := t.withSpockMetadata(rowAsMap, node2)
383383
newDiffsForPair.Rows[node2] = append(newDiffsForPair.Rows[node2], utils.MapToOrderedMap(rowWithMeta, t.Cols))
384384
}
385385
}

internal/consistency/mtree/merkle.go

Lines changed: 11 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -99,7 +99,7 @@ type MerkleTreeTask struct {
9999
diffMutex sync.Mutex
100100
diffRowKeySets map[string]map[string]map[string]struct{}
101101
StartTime time.Time
102-
NodeOriginNames map[string]string
102+
NodeOriginNames map[string]map[string]string
103103

104104
Ctx context.Context
105105
}
@@ -510,8 +510,14 @@ func (m *MerkleTreeTask) loadNodeOriginNames() error {
510510
return nil
511511
}
512512

513+
m.NodeOriginNames = make(map[string]map[string]string)
514+
513515
var lastErr error
514516
for _, nodeInfo := range m.ClusterNodes {
517+
nodeName, _ := nodeInfo["Name"].(string)
518+
if nodeName == "" {
519+
continue
520+
}
515521
pool, err := auth.GetClusterNodeConnection(m.Ctx, nodeInfo, m.connOpts())
516522
if err != nil {
517523
lastErr = err
@@ -523,15 +529,13 @@ func (m *MerkleTreeTask) loadNodeOriginNames() error {
523529
lastErr = err
524530
continue
525531
}
526-
m.NodeOriginNames = names
527-
return nil
532+
m.NodeOriginNames[nodeName] = names
528533
}
529534

530-
m.NodeOriginNames = make(map[string]string)
531-
if lastErr != nil {
535+
if len(m.NodeOriginNames) == 0 && lastErr != nil {
532536
return lastErr
533537
}
534-
return fmt.Errorf("no nodes available to load node origin names")
538+
return nil
535539
}
536540

537541
func (m *MerkleTreeTask) appendDiffs(nodePairKey string, work CompareRangesWorkItem, pr1, pr2 []types.OrderedMap) error {
@@ -603,7 +607,7 @@ func (m *MerkleTreeTask) addRowToDiff(nodePairKey, nodeName string, row types.Or
603607
}
604608

605609
rowMap := utils.OrderedMapToMap(row)
606-
rowMap["node_origin"] = utils.TranslateNodeOrigin(rowMap["node_origin"], m.NodeOriginNames)
610+
rowMap["node_origin"] = utils.TranslateNodeOrigin(rowMap["node_origin"], m.NodeOriginNames[nodeName])
607611
rowWithMeta := utils.AddSpockMetadata(rowMap)
608612
orderedRow := utils.MapToOrderedMap(rowWithMeta, m.Cols)
609613

0 commit comments

Comments
 (0)