@@ -119,7 +119,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_Complete_Flow_No_Data() {
119
119
Destination : s .Peer ().Name ,
120
120
}
121
121
122
- flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s . t )
122
+ flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s )
123
123
flowConnConfig .MaxBatchSize = 1
124
124
125
125
env := e2e .ExecutePeerflow (tc , peerflow .CDCFlowWorkflow , flowConnConfig , nil )
@@ -150,7 +150,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_Char_ColType_Error() {
150
150
Destination : s .Peer ().Name ,
151
151
}
152
152
153
- flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s . t )
153
+ flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s )
154
154
flowConnConfig .MaxBatchSize = 1
155
155
156
156
env := e2e .ExecutePeerflow (tc , peerflow .CDCFlowWorkflow , flowConnConfig , nil )
@@ -182,7 +182,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_Toast_BQ() {
182
182
Destination : s .Peer ().Name ,
183
183
}
184
184
185
- flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s . t )
185
+ flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s )
186
186
flowConnConfig .MaxBatchSize = 100
187
187
188
188
// wait for PeerFlowStatusQuery to finish setup
@@ -233,7 +233,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_Toast_Advance_1_BQ() {
233
233
Destination : s .Peer ().Name ,
234
234
}
235
235
236
- flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s . t )
236
+ flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s )
237
237
flowConnConfig .MaxBatchSize = 100
238
238
239
239
// wait for PeerFlowStatusQuery to finish setup
@@ -289,7 +289,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_Toast_Advance_2_BQ() {
289
289
Destination : s .Peer ().Name ,
290
290
}
291
291
292
- flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s . t )
292
+ flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s )
293
293
flowConnConfig .MaxBatchSize = 100
294
294
295
295
// wait for PeerFlowStatusQuery to finish setup
@@ -339,7 +339,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_Toast_Advance_3_BQ() {
339
339
Destination : s .Peer ().Name ,
340
340
}
341
341
342
- flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s . t )
342
+ flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s )
343
343
flowConnConfig .MaxBatchSize = 100
344
344
345
345
// wait for PeerFlowStatusQuery to finish setup
@@ -395,7 +395,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_Types_BQ() {
395
395
Destination : s .Peer ().Name ,
396
396
}
397
397
398
- flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s . t )
398
+ flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s )
399
399
flowConnConfig .MaxBatchSize = 100
400
400
401
401
// wait for PeerFlowStatusQuery to finish setup
@@ -477,7 +477,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_NaN_Doubles_BQ() {
477
477
Destination : s .Peer ().Name ,
478
478
}
479
479
480
- flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s . t )
480
+ flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s )
481
481
flowConnConfig .MaxBatchSize = 100
482
482
483
483
// wait for PeerFlowStatusQuery to finish setup
@@ -529,7 +529,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_Invalid_Geo_BQ_Avro_CDC() {
529
529
Destination : s .Peer ().Name ,
530
530
}
531
531
532
- flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s . t )
532
+ flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s )
533
533
flowConnConfig .MaxBatchSize = 100
534
534
535
535
// wait for PeerFlowStatusQuery to finish setup
@@ -605,7 +605,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_Multi_Table_BQ() {
605
605
Destination : s .Peer ().Name ,
606
606
}
607
607
608
- flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s . t )
608
+ flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s )
609
609
flowConnConfig .MaxBatchSize = 100
610
610
611
611
// wait for PeerFlowStatusQuery to finish setup
@@ -659,7 +659,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_Simple_Schema_Changes_BQ() {
659
659
Destination : s .Peer ().Name ,
660
660
}
661
661
662
- flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s . t )
662
+ flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s )
663
663
flowConnConfig .MaxBatchSize = 100
664
664
665
665
// wait for PeerFlowStatusQuery to finish setup
@@ -737,7 +737,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_All_Types_Schema_Changes_BQ() {
737
737
Destination : s .Peer ().Name ,
738
738
}
739
739
740
- flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s . t )
740
+ flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s )
741
741
flowConnConfig .MaxBatchSize = 100
742
742
743
743
// wait for PeerFlowStatusQuery to finish setup
@@ -805,7 +805,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_Composite_PKey_BQ() {
805
805
Destination : s .Peer ().Name ,
806
806
}
807
807
808
- flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s . t )
808
+ flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s )
809
809
flowConnConfig .MaxBatchSize = 100
810
810
811
811
// wait for PeerFlowStatusQuery to finish setup
@@ -861,7 +861,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_Composite_PKey_Toast_1_BQ() {
861
861
Destination : s .Peer ().Name ,
862
862
}
863
863
864
- flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s . t )
864
+ flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s )
865
865
flowConnConfig .MaxBatchSize = 100
866
866
flowConnConfig .SoftDeleteColName = ""
867
867
flowConnConfig .SyncedAtColName = ""
@@ -921,7 +921,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_Composite_PKey_Toast_2_BQ() {
921
921
Destination : s .Peer ().Name ,
922
922
}
923
923
924
- flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s . t )
924
+ flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s )
925
925
flowConnConfig .MaxBatchSize = 100
926
926
927
927
// wait for PeerFlowStatusQuery to finish setup
@@ -972,7 +972,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_Columns_BQ() {
972
972
SoftDelete : true ,
973
973
}
974
974
975
- flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s . t )
975
+ flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s )
976
976
flowConnConfig .MaxBatchSize = 100
977
977
978
978
env := e2e .ExecutePeerflow (tc , peerflow .CDCFlowWorkflow , flowConnConfig , nil )
@@ -1022,7 +1022,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_Multi_Table_Multi_Dataset_BQ() {
1022
1022
Destination : s .Peer ().Name ,
1023
1023
}
1024
1024
1025
- flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s . t )
1025
+ flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s )
1026
1026
flowConnConfig .MaxBatchSize = 100
1027
1027
1028
1028
// wait for PeerFlowStatusQuery to finish setup
@@ -1103,7 +1103,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_Soft_Delete_Basic() {
1103
1103
_ , err = s .Conn ().Exec (context .Background (), fmt .Sprintf (`DELETE FROM %s WHERE id=1` , srcTableName ))
1104
1104
e2e .EnvNoError (s .t , env , err )
1105
1105
e2e .EnvWaitFor (s .t , env , 3 * time .Minute , "normalize delete" , func () bool {
1106
- pgRows , err := e2e . GetPgRows ( s . conn , s .bqSuffix , srcName , "id,c1,c2,t" )
1106
+ pgRows , err := s . Source (). GetRows ( s .bqSuffix , srcName , "id,c1,c2,t" )
1107
1107
if err != nil {
1108
1108
return false
1109
1109
}
@@ -1248,7 +1248,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_Soft_Delete_UD_Same_Batch() {
1248
1248
e2e .EnvNoError (s .t , env , insertTx .Commit (context .Background ()))
1249
1249
1250
1250
e2e .EnvWaitFor (s .t , env , 3 * time .Minute , "normalize transaction" , func () bool {
1251
- pgRows , err := e2e . GetPgRows ( s . conn , s .bqSuffix , srcName , "id,c1,c2,t" )
1251
+ pgRows , err := s . Source (). GetRows ( s .bqSuffix , srcName , "id,c1,c2,t" )
1252
1252
e2e .EnvNoError (s .t , env , err )
1253
1253
rows , err := s .GetRowsWhere (dstName , "id,c1,c2,t" , "NOT _PEERDB_IS_DELETED" )
1254
1254
if err != nil {
@@ -1312,7 +1312,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_Soft_Delete_Insert_After_Delete() {
1312
1312
"DELETE FROM %s WHERE id=1" , srcTableName ))
1313
1313
e2e .EnvNoError (s .t , env , err )
1314
1314
e2e .EnvWaitFor (s .t , env , 3 * time .Minute , "normalize delete" , func () bool {
1315
- pgRows , err := e2e . GetPgRows ( s . conn , s .bqSuffix , tableName , "id,c1,c2,t" )
1315
+ pgRows , err := s . Source (). GetRows ( s .bqSuffix , tableName , "id,c1,c2,t" )
1316
1316
if err != nil {
1317
1317
return false
1318
1318
}
@@ -1362,7 +1362,7 @@ func (s PeerFlowE2ETestSuiteBQ) Test_JSON_PKey_BQ() {
1362
1362
Destination : s .Peer ().Name ,
1363
1363
}
1364
1364
1365
- flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s . t )
1365
+ flowConnConfig := connectionGen .GenerateFlowConnectionConfigs (s )
1366
1366
flowConnConfig .MaxBatchSize = 100
1367
1367
flowConnConfig .SoftDeleteColName = ""
1368
1368
flowConnConfig .SyncedAtColName = ""
0 commit comments