Documentation
¶
Index ¶
- func AddSuffix(s Suite, str string) string
- func AttachSchema(s Suite, table string) string
- func CatalogTestAccessPool() (*pgxpool.Pool, error)
- func CockroachDBTestContainerConfig(t *testing.T) *protos.CockroachDBConfig
- func ConsumeAllMessages(ctx context.Context, namespaceName string, eventhubName string, ...) ([]string, error)
- func CreatePeer(t *testing.T, peer *protos.Peer)
- func CreateQRepWorkflowConfig(t *testing.T, flowJobName string, sourceTable string, dstTable string, ...) *protos.QRepConfig
- func CreateTableForQRep(ctx context.Context, conn *pgx.Conn, suffix string, tableName string) error
- func DeleteEventhub(ctx context.Context, eventhubName string, ...) error
- func EnvEqualRecordBatches(t *testing.T, env WorkflowRun, q *model.QRecordBatch, ...)
- func EnvEqualTables[TSource connectors.Connector](env WorkflowRun, suite RowSource, table string, cols string)
- func EnvEqualTablesWithNames(env WorkflowRun, suite RowSource, srcTable string, dstTable string, ...)
- func EnvGetRunID(t *testing.T, env WorkflowRun) string
- func EnvGetWorkflowState(t *testing.T, env WorkflowRun) cdc_state.CDCFlowWorkflowState
- func EnvNoError(t *testing.T, env WorkflowRun, err error)
- func EnvTrue(t *testing.T, env WorkflowRun, val bool)
- func EnvWaitFor(t *testing.T, env WorkflowRun, timeout time.Duration, reason string, ...)
- func EnvWaitForCount(env WorkflowRun, suite RowSource, reason string, dstTable string, cols string, ...)
- func EnvWaitForEqualTables(env WorkflowRun, suite RowSource, reason string, table string, cols string)
- func EnvWaitForEqualTablesWithNames(env WorkflowRun, suite RowSource, reason string, srcTable string, ...)
- func EnvWaitForEqualTablesWithNames_Only(env WorkflowRun, suite RowSource, reason string, srcTable string, ...)
- func EnvWaitForFinished(t *testing.T, env WorkflowRun, timeout time.Duration)
- func ExpectedDestinationIdentifier(s GenericSuite, ident string) string
- func ExpectedDestinationSchema(s GenericSuite, dstTable string, userColumns []*protos.FieldDescription) *protos.TableSchema
- func ExpectedDestinationTableName(s GenericSuite, table string) string
- func GeneratePostgresPeer(t *testing.T) *protos.Peer
- func GetLogCount(ctx context.Context, catalog shared.CatalogPool, ...) (int, error)
- func GetOwnersSchema() *types.QRecordSchema
- func GetOwnersSelectorStringsSF() [2]string
- func GetPostgresToxicProxy(t *testing.T, suffix string, port uint32) (*tp.Proxy, error)
- func GetTestDatabase(suffix string) string
- func InitToxiproxy() error
- func InsertScript(t *testing.T, name string, lang string, source string)
- func NewApiClient() (protos.FlowServiceClient, error)
- func NewTemporalClient(t *testing.T) client.Client
- func PopulateSourceTable(ctx context.Context, conn *pgx.Conn, suffix string, tableName string, ...) error
- func RequireEmptyDestinationTable(suite RowSource, dstTable string, cols string)
- func RequireEnvCanceled(t *testing.T, env WorkflowRun)
- func RequireEqualRecordBatches(t *testing.T, q *model.QRecordBatch, other *model.QRecordBatch)
- func RequireEqualTableSchemas(t *testing.T, expected *protos.TableSchema, actual *protos.TableSchema) bool
- func RequireEqualTables(suite RowSource, table string, cols string)
- func RequireEqualTablesWithNames(suite RowSource, srcTable string, dstTable string, cols string)
- func RevokePermissionForTableColumns(ctx context.Context, conn *pgx.Conn, tableIdentifier string, ...) error
- func RunApiSuite[TSource SuiteSource](t *testing.T, setup func(*testing.T, string) (TSource, error))
- func RunPsql(t *testing.T, peer string, extraArgs ...string) (string, error)
- func Schema(s Suite) string
- func SetupCDCFlowStatusQuery(t *testing.T, env WorkflowRun, config *protos.FlowConnectionConfigs)
- func SetupClickHouseSuite[TSource SuiteSource](t *testing.T, cluster bool, ...) func(*testing.T) ClickHouseSuite
- func SetupGenericSuite[T GenericSuite](f func(t *testing.T) T) func(t *testing.T) Generic
- func SignalWorkflow[T any](ctx context.Context, env WorkflowRun, signal model.TypedSignal[T], value T)
- func SwitchboardDSN(peer string, options map[string]string) string
- func TableMappings(s GenericSuite, tables ...string) []*protos.TableMapping
- func TearDownPostgres(ctx context.Context, s PgSuite)
- type APITestSuite
- func (s APITestSuite) Connector() *connpostgres.PostgresConnector
- func (s APITestSuite) DestinationTable(table string) string
- func (s APITestSuite) Source() SuiteSource
- func (s APITestSuite) Suffix() string
- func (s APITestSuite) T() *testing.T
- func (s APITestSuite) Teardown(ctx context.Context)
- func (s APITestSuite) TestAlertConfig()
- func (s APITestSuite) TestCancelAddCancel()
- func (s APITestSuite) TestCancelRepairsPostgresZeroOIDs()
- func (s APITestSuite) TestCancelTableAdditionDuringSetupFlow()
- func (s APITestSuite) TestCancelTableAdditionRemoveAddRemove()
- func (s APITestSuite) TestCancelTableAddition_NoRemovalAssumed()
- func (s APITestSuite) TestCancelTableAddition_NoRemovalAssumedWithRemoval()
- func (s APITestSuite) TestCancelTableAddition_WithRemoval()
- func (s APITestSuite) TestClickHouseMirrorValidation_NoPrimaryKey()
- func (s APITestSuite) TestClickHouseMirrorValidation_NoPrimaryKey_ReplicaIdentityFull()
- func (s APITestSuite) TestClickHouseMirrorValidation_Pass()
- func (s APITestSuite) TestCreateCDCFlowAttachCanceledWorkflow()
- func (s APITestSuite) TestCreateCDCFlowAttachConcurrentRequests()
- func (s APITestSuite) TestCreateCDCFlowAttachConcurrentRequestsToxi()
- func (s APITestSuite) TestCreateCDCFlowAttachExternalFlowEntry()
- func (s APITestSuite) TestCreateCDCFlowAttachIdempotentAfterContinueAsNew()
- func (s APITestSuite) TestCreateCDCFlowAttachSequentialRequests()
- func (s APITestSuite) TestDoubleClickCancelTableAddition()
- func (s APITestSuite) TestDropCompleted()
- func (s APITestSuite) TestDropCompletedAndUnavailable()
- func (s APITestSuite) TestDropMissing()
- func (s APITestSuite) TestDropQRep()
- func (s APITestSuite) TestEditTablesBeforeResync()
- func (s APITestSuite) TestFlowStatusUpdate()
- func (s APITestSuite) TestGetColumnsNullability()
- func (s APITestSuite) TestGetMirrorRowCounts()
- func (s APITestSuite) TestGetMirrorRowCountsNonCDC()
- func (s APITestSuite) TestGetTablesExcludeViews()
- func (s APITestSuite) TestGetTablesUnloggedNotMirrorable()
- func (s APITestSuite) TestGetVersion()
- func (s APITestSuite) TestMirrorValidation_InvalidTableMappings()
- func (s APITestSuite) TestMirrorValidation_MissingSourceTable()
- func (s APITestSuite) TestMySQLFlavorSwap()
- func (s APITestSuite) TestPostgresDestinationValidation_ExtraColumnsOk()
- func (s APITestSuite) TestPostgresDestinationValidation_MissingColumns()
- func (s APITestSuite) TestPostgresDestinationValidation_MissingSchema()
- func (s APITestSuite) TestPostgresDestinationValidation_NonEmptyTableAllowedWithoutSnapshot()
- func (s APITestSuite) TestPostgresDestinationValidation_NonEmptyTableBlocksSnapshot()
- func (s APITestSuite) TestPostgresDestinationValidation_NumericPrecisionMismatch()
- func (s APITestSuite) TestPostgresDestinationValidation_NumericSuperset()
- func (s APITestSuite) TestPostgresDestinationValidation_TypeCompatible()
- func (s APITestSuite) TestPostgresDestinationValidation_UnboundedNumeric()
- func (s APITestSuite) TestPostgresDestinationValidation_UserDefinedTypeMatch()
- func (s APITestSuite) TestPostgresDestinationValidation_UserDefinedTypeMismatch()
- func (s APITestSuite) TestPostgresDestinationValidation_WithExcludedColumns()
- func (s APITestSuite) TestPostgresValidation_Pass()
- func (s APITestSuite) TestPostgresValidation_WrongPassword()
- func (s APITestSuite) TestQRep()
- func (s APITestSuite) TestResetMirrorSequences()
- func (s APITestSuite) TestResetMirrorSequences_EmptyTable()
- func (s APITestSuite) TestResetMirrorSequences_NonPgDestination()
- func (s APITestSuite) TestResyncCompleted()
- func (s APITestSuite) TestResyncFailed()
- func (s APITestSuite) TestResyncSourceTableMissing()
- func (s APITestSuite) TestResyncTablesNotInPublication()
- func (s APITestSuite) TestResyncWithSnapshotConfigDuringSnapshot()
- func (s APITestSuite) TestResyncWithSnapshotConfigOnCompletedSnapshotOnlyPipe()
- func (s APITestSuite) TestResyncWithSnapshotConfigOnPausedPipe()
- func (s APITestSuite) TestResyncWithSnapshotConfigOnRunningPipe()
- func (s APITestSuite) TestSchemaEndpoints()
- func (s APITestSuite) TestScripts()
- func (s APITestSuite) TestSettings()
- func (s APITestSuite) TestSnapshotNullPartitionKey()
- func (s APITestSuite) TestTableAdditionWithoutInitialLoad()
- func (s APITestSuite) TestTerminateDuringResyncDropFlow()
- func (s APITestSuite) TestTotalRowsSyncedByMirror()
- func (s APITestSuite) TestValidateCDCMirror_ServerIDInUseAcrossPeers()
- func (s APITestSuite) TestValidateCDCMirror_ServerIDPeerReuse()
- func (s APITestSuite) TestValidateCDCMirror_ServerIDResync()
- type BigQueryTestHelper
- func (b *BigQueryTestHelper) CheckNull(ctx context.Context, tableName string, colName []string) (bool, error)
- func (b *BigQueryTestHelper) CountObjectsInGCSPath(ctx context.Context, gcsPath string) (int, error)
- func (b *BigQueryTestHelper) CountRows(ctx context.Context, tableName string) (int, error)
- func (b *BigQueryTestHelper) CountRowsWithDataset(ctx context.Context, dataset, tableName string, nonNullCol string) (int, error)
- func (b *BigQueryTestHelper) DropDataset(ctx context.Context, datasetName string) error
- func (b *BigQueryTestHelper) ExecuteAndProcessQuery(ctx context.Context, query string) (*model.QRecordBatch, error)
- func (b *BigQueryTestHelper) RecreateDataset(ctx context.Context) error
- func (b *BigQueryTestHelper) RunInt64Query(ctx context.Context, query string) (int64, error)
- func (b *BigQueryTestHelper) SelectRow(ctx context.Context, tableName string, cols ...string) ([]bigquery.Value, error)
- type ClickHouseMVManager
- type ClickHouseSuite
- func (s ClickHouseSuite) Catalog() shared.CatalogPool
- func (s ClickHouseSuite) Connector() connectors.Connector
- func (s ClickHouseSuite) CountNonDeletedRows(table string) (int, error)
- func (s ClickHouseSuite) CreateRMTTable(tableName string, columns []TestClickHouseColumn, orderingKey string) error
- func (s ClickHouseSuite) CreateSlowInsertViaMV(tableName string, sleepSeconds int) (func(), error)
- func (s ClickHouseSuite) DestinationConnector() connectors.Connector
- func (s ClickHouseSuite) DestinationTable(table string) string
- func (s ClickHouseSuite) DropTable(tableName string) error
- func (s ClickHouseSuite) GetRows(table string, cols string) (*model.QRecordBatch, error)
- func (s ClickHouseSuite) IsCluster() bool
- func (s ClickHouseSuite) NewMVManager(tableName string, suffix string) *ClickHouseMVManager
- func (s ClickHouseSuite) Peer() *protos.Peer
- func (s ClickHouseSuite) PeerForDatabase(dbname string) *protos.Peer
- func (s ClickHouseSuite) S3Helper() *S3TestHelper
- func (s ClickHouseSuite) Source() SuiteSource
- func (s ClickHouseSuite) Suffix() string
- func (s ClickHouseSuite) T() *testing.T
- func (s ClickHouseSuite) Teardown(ctx context.Context)
- func (s ClickHouseSuite) Test_Addition_Removal()
- func (s ClickHouseSuite) Test_AvroNullableLax()
- func (s ClickHouseSuite) Test_Binary_Format_Base64()
- func (s ClickHouseSuite) Test_Binary_Format_Hex()
- func (s ClickHouseSuite) Test_Binary_Format_Raw()
- func (s ClickHouseSuite) Test_CTID_Inherited_Table()
- func (s ClickHouseSuite) Test_CTID_Inherited_Table_Extra_Columns()
- func (s ClickHouseSuite) Test_CTID_Multi_Level_Inherited_Table()
- func (s ClickHouseSuite) Test_CTID_Multi_Level_Partitioned_Table()
- func (s ClickHouseSuite) Test_CTID_Partitioned_Table()
- func (s ClickHouseSuite) Test_Chunking_Initial_Load_Parts_Per_Partition()
- func (s ClickHouseSuite) Test_CoalescingEngine()
- func (s ClickHouseSuite) Test_Column_Exclusion()
- func (s ClickHouseSuite) Test_Composite_PKey()
- func (s ClickHouseSuite) Test_Destination_Type_Conversion()
- func (s ClickHouseSuite) Test_Extra_CH_Columns()
- func (s ClickHouseSuite) Test_First_Row_Lag_Times_Recorded()
- func (s ClickHouseSuite) Test_Geometric_Types()
- func (s ClickHouseSuite) Test_InfiniteTimestamp()
- func (s ClickHouseSuite) Test_InitialLoadOnly_No_Primary_Key()
- func (s ClickHouseSuite) Test_JSON_CH()
- func (s ClickHouseSuite) Test_JSON_Null()
- func (s ClickHouseSuite) Test_Large_Numeric()
- func (s ClickHouseSuite) Test_Large_Text_CDC()
- func (s ClickHouseSuite) Test_Normalize_Metadata_With_Retry()
- func (s ClickHouseSuite) Test_NullEngine()
- func (s ClickHouseSuite) Test_NullableColumnSetting()
- func (s ClickHouseSuite) Test_NullableMirrorSetting()
- func (s ClickHouseSuite) Test_Nullable_Schema_Change()
- func (s ClickHouseSuite) Test_Nullable_Schema_Change_Replident_Full()
- func (s ClickHouseSuite) Test_Nullable_Schema_Change_Replident_Index()
- func (s ClickHouseSuite) Test_Numeric_Truncation_With_UnbNumAsString_FF()
- func (s ClickHouseSuite) Test_Numeric_Truncation_Without_UnbNumAsString_FF()
- func (s ClickHouseSuite) Test_Offload_Partition_Ranges()
- func (s ClickHouseSuite) Test_PG_AlterTableAddColumnDefault()
- func (s ClickHouseSuite) Test_PG_AlterTableAddColumnDefaultUntranslated()
- func (s ClickHouseSuite) Test_PG_Domain_Bytea()
- func (s ClickHouseSuite) Test_PartitionBy()
- func (s ClickHouseSuite) Test_PartitionByExpr()
- func (s ClickHouseSuite) Test_Partition_By_CTID_With_Num_Partitions_Override()
- func (s ClickHouseSuite) Test_Partition_Key_Empty()
- func (s ClickHouseSuite) Test_Partition_Key_Integer()
- func (s ClickHouseSuite) Test_Partition_Key_Null()
- func (s ClickHouseSuite) Test_Partition_Key_Timestamp()
- func (s ClickHouseSuite) Test_PgVector()
- func (s ClickHouseSuite) Test_PgVector_Version0()
- func (s ClickHouseSuite) Test_Removal_Shared_Destination()
- func (s ClickHouseSuite) Test_Replident_Full_Unchanged_TOAST_Updates()
- func (s ClickHouseSuite) Test_SchemaAsColumn()
- func (s ClickHouseSuite) Test_Schema_Change_After_Resync_Cluster()
- func (s ClickHouseSuite) Test_SkipSnapshotExport()
- func (s ClickHouseSuite) Test_Sync_Error_Cancels_InProgress_Normalize()
- func (s ClickHouseSuite) Test_Time64()
- func (s ClickHouseSuite) Test_Types_CH()
- func (s ClickHouseSuite) Test_Unbounded_Numeric_With_FF()
- func (s ClickHouseSuite) Test_Unbounded_Numeric_Without_FF()
- func (s ClickHouseSuite) Test_Unprivileged_Postgres_Columns()
- func (s ClickHouseSuite) Test_Update_PKey_Env_Disabled()
- func (s ClickHouseSuite) Test_Update_PKey_Env_Enabled()
- func (s ClickHouseSuite) Test_ValidatePartitionByExpression()
- func (s ClickHouseSuite) Test_WeirdTable_Dash()
- func (s ClickHouseSuite) Test_WeirdTable_Keyword()
- func (s ClickHouseSuite) Test_WeirdTable_MixedCase()
- func (s ClickHouseSuite) Test_WeirdTable_Question()
- func (s ClickHouseSuite) WeirdTable(tableName string)
- type CockroachDBSource
- func (s *CockroachDBSource) AdminConn() *pgx.Conn
- func (s *CockroachDBSource) CockroachDBConnector() *conncockroachdb.CockroachDBConnector
- func (s *CockroachDBSource) Config() *protos.CockroachDBConfig
- func (s *CockroachDBSource) Connector() connectors.Connector
- func (s *CockroachDBSource) Exec(ctx context.Context, sql string, args ...any) error
- func (s *CockroachDBSource) GeneratePeer(t *testing.T) *protos.Peer
- func (s *CockroachDBSource) GetRows(ctx context.Context, suffix, table, cols string) (*model.QRecordBatch, error)
- func (s *CockroachDBSource) Teardown(t *testing.T, ctx context.Context, suffix string)
- type FlowConnectionGenerationConfig
- type Generic
- func (s Generic) Test_Custom_Replication_Slot_Starting_With_Numbers_CDC_Only()
- func (s Generic) Test_Inheritance_Table_With_Dynamic_Setting()
- func (s Generic) Test_Inheritance_Table_Without_Dynamic_Setting()
- func (s Generic) Test_Initial_Custom_Partition()
- func (s Generic) Test_Partitioned_Table()
- func (s Generic) Test_Partitioned_Table_With_Different_Column_Ordering()
- func (s Generic) Test_Partitioned_Table_Without_Publish_Via_Partition_Root()
- func (s Generic) Test_Schema_Change_Drop_Consecutive_Columns()
- func (s Generic) Test_Schema_Change_Lost_Column_Bug()
- func (s Generic) Test_Schema_Changes_Cutoff_Bug()
- func (s Generic) Test_Simple_Flow()
- func (s Generic) Test_Simple_Schema_Changes()
- type GenericSuite
- type MongoSource
- func (s *MongoSource) AdminClient() *mongo.Client
- func (s *MongoSource) Config() *protos.MongoConfig
- func (s *MongoSource) Connector() connectors.Connector
- func (s *MongoSource) Exec(ctx context.Context, sql string, args ...any) error
- func (s *MongoSource) GeneratePeer(t *testing.T) *protos.Peer
- func (s *MongoSource) GetRows(ctx context.Context, suffix, table, cols string) (*model.QRecordBatch, error)
- func (s *MongoSource) Teardown(t *testing.T, ctx context.Context, suffix string)
- type MySQLTestContainerConfig
- type MySqlSource
- func (s *MySqlSource) Connector() connectors.Connector
- func (s *MySqlSource) Exec(ctx context.Context, sql string, args ...any) error
- func (s *MySqlSource) GeneratePeer(t *testing.T) *protos.Peer
- func (s *MySqlSource) GetRows(ctx context.Context, suffix string, table string, cols string) (*model.QRecordBatch, error)
- func (s *MySqlSource) Teardown(t *testing.T, ctx context.Context, suffix string)
- type PgSuite
- type PostgresSource
- func (s *PostgresSource) Connector() connectors.Connector
- func (s *PostgresSource) Exec(ctx context.Context, sql string, args ...any) error
- func (s *PostgresSource) GeneratePeer(t *testing.T) *protos.Peer
- func (s *PostgresSource) GetRows(ctx context.Context, suffix string, table string, cols string) (*model.QRecordBatch, error)
- func (s *PostgresSource) GetRowsOnly(ctx context.Context, suffix string, table string, cols string) (*model.QRecordBatch, error)
- func (s *PostgresSource) Query(ctx context.Context, query string) (*model.QRecordBatch, error)
- func (s *PostgresSource) Teardown(t *testing.T, ctx context.Context, suffix string)
- type RowSource
- type S3Environment
- type S3PeerCredentials
- type S3TestHelper
- type SnowflakeTestHelper
- func (s *SnowflakeTestHelper) CheckIsDeleted(ctx context.Context, query string) error
- func (s *SnowflakeTestHelper) CheckNull(ctx context.Context, tableName string, colNames []string) (bool, error)
- func (s *SnowflakeTestHelper) CheckSyncedAt(ctx context.Context, query string) error
- func (s *SnowflakeTestHelper) Cleanup(ctx context.Context) error
- func (s *SnowflakeTestHelper) CountNonNullRows(ctx context.Context, tableName string, columnName string) (int64, error)
- func (s *SnowflakeTestHelper) CountRows(ctx context.Context, tableName string) (int64, error)
- func (s *SnowflakeTestHelper) CountSRIDs(ctx context.Context, tableName string, columnName string) (int64, error)
- func (s *SnowflakeTestHelper) ExecuteAndProcessQuery(ctx context.Context, query string) (*model.QRecordBatch, error)
- func (s *SnowflakeTestHelper) RunCommand(ctx context.Context, command string) error
- func (s *SnowflakeTestHelper) RunIntQuery(ctx context.Context, query string) (int, error)
- type Suite
- type SuiteSource
- type TestClickHouseColumn
- type WorkflowRun
- func ExecuteDropFlow(ctx context.Context, tc client.Client, config *protos.FlowConnectionConfigs) WorkflowRun
- func ExecutePeerflow(t *testing.T, tc client.Client, config *protos.FlowConnectionConfigs) WorkflowRun
- func ExecuteWorkflow(ctx context.Context, tc client.Client, taskQueueID shared.TaskQueueID, wf any, ...) WorkflowRun
- func GetPeerflow(ctx context.Context, catalog shared.CatalogPool, tc client.Client, ...) (WorkflowRun, error)
- func RunQRepFlowWorkflow(t *testing.T, tc client.Client, config *protos.QRepConfig) WorkflowRun
- func (env WorkflowRun) Cancel(ctx context.Context)
- func (env WorkflowRun) Error(ctx context.Context) error
- func (env WorkflowRun) Finished(ctx context.Context) bool
- func (env WorkflowRun) GetFlowStatus(t *testing.T) protos.FlowStatus
- func (env WorkflowRun) Query(ctx context.Context, queryType string, args ...any) (converter.EncodedValue, error)
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func AttachSchema ¶
func CatalogTestAccessPool ¶
func CockroachDBTestContainerConfig ¶
func CockroachDBTestContainerConfig(t *testing.T) *protos.CockroachDBConfig
CockroachDBTestContainerConfig starts an isolated source for tests that change cluster settings. CI and Tilt pass COCKROACHDB_IMAGE so this test covers the same version as the shared source.
func ConsumeAllMessages ¶
func ConsumeAllMessages( ctx context.Context, namespaceName string, eventhubName string, expectedNum int, config *protos.EventHubGroupConfig, ) ([]string, error)
consume all messages from the eventhub with the given name“. returns as a list of strings.
func CreateTableForQRep ¶
func DeleteEventhub ¶
func EnvEqualRecordBatches ¶
func EnvEqualRecordBatches(t *testing.T, env WorkflowRun, q *model.QRecordBatch, other *model.QRecordBatch)
func EnvEqualTables ¶
func EnvEqualTables[TSource connectors.Connector](env WorkflowRun, suite RowSource, table string, cols string)
func EnvEqualTablesWithNames ¶
func EnvEqualTablesWithNames( env WorkflowRun, suite RowSource, srcTable string, dstTable string, cols string, )
func EnvGetRunID ¶
func EnvGetRunID(t *testing.T, env WorkflowRun) string
func EnvGetWorkflowState ¶
func EnvGetWorkflowState(t *testing.T, env WorkflowRun) cdc_state.CDCFlowWorkflowState
func EnvNoError ¶
func EnvNoError(t *testing.T, env WorkflowRun, err error)
Helper function to assert errors in go routines running concurrent to workflows This achieves two goals: 1. cancel workflow to avoid waiting on goroutine which has failed 2. get around t.FailNow being incorrect when called from non initial goroutine
func EnvWaitFor ¶
func EnvWaitForCount ¶
func EnvWaitForEqualTables ¶
func EnvWaitForEqualTables( env WorkflowRun, suite RowSource, reason string, table string, cols string, )
func EnvWaitForFinished ¶
func EnvWaitForFinished(t *testing.T, env WorkflowRun, timeout time.Duration)
func ExpectedDestinationIdentifier ¶
func ExpectedDestinationIdentifier(s GenericSuite, ident string) string
func ExpectedDestinationSchema ¶
func ExpectedDestinationSchema(s GenericSuite, dstTable string, userColumns []*protos.FieldDescription) *protos.TableSchema
func ExpectedDestinationTableName ¶
func ExpectedDestinationTableName(s GenericSuite, table string) string
func GetLogCount ¶
func GetOwnersSchema ¶
func GetOwnersSchema() *types.QRecordSchema
func GetOwnersSelectorStringsSF ¶
func GetOwnersSelectorStringsSF() [2]string
func GetPostgresToxicProxy ¶
GetPostgresToxicProxy gets or creates the PostgreSQL proxy
func GetTestDatabase ¶
func InitToxiproxy ¶
func InitToxiproxy() error
InitToxiproxy initializes the Toxiproxy client (singleton pattern)
func InsertScript ¶
InsertScript registers a transform script in the catalog's public.scripts table.
func NewApiClient ¶
func NewApiClient() (protos.FlowServiceClient, error)
func PopulateSourceTable ¶
func RequireEnvCanceled ¶
func RequireEnvCanceled(t *testing.T, env WorkflowRun)
func RequireEqualRecordBatches ¶
func RequireEqualRecordBatches(t *testing.T, q *model.QRecordBatch, other *model.QRecordBatch)
func RequireEqualTableSchemas ¶
func RequireEqualTableSchemas(t *testing.T, expected *protos.TableSchema, actual *protos.TableSchema) bool
func RequireEqualTables ¶
func RunApiSuite ¶
func SetupCDCFlowStatusQuery ¶
func SetupCDCFlowStatusQuery(t *testing.T, env WorkflowRun, config *protos.FlowConnectionConfigs)
func SetupClickHouseSuite ¶
func SetupGenericSuite ¶
func SignalWorkflow ¶
func SignalWorkflow[T any](ctx context.Context, env WorkflowRun, signal model.TypedSignal[T], value T)
func TableMappings ¶
func TableMappings(s GenericSuite, tables ...string) []*protos.TableMapping
func TearDownPostgres ¶
Types ¶
type APITestSuite ¶
type APITestSuite struct {
protos.FlowServiceClient
// contains filtered or unexported fields
}
func (APITestSuite) Connector ¶
func (s APITestSuite) Connector() *connpostgres.PostgresConnector
func (APITestSuite) DestinationTable ¶
func (s APITestSuite) DestinationTable(table string) string
func (APITestSuite) Source ¶
func (s APITestSuite) Source() SuiteSource
func (APITestSuite) Suffix ¶
func (s APITestSuite) Suffix() string
func (APITestSuite) T ¶
func (s APITestSuite) T() *testing.T
func (APITestSuite) Teardown ¶
func (s APITestSuite) Teardown(ctx context.Context)
func (APITestSuite) TestAlertConfig ¶
func (s APITestSuite) TestAlertConfig()
func (APITestSuite) TestCancelAddCancel ¶
func (s APITestSuite) TestCancelAddCancel()
func (APITestSuite) TestCancelRepairsPostgresZeroOIDs ¶
func (s APITestSuite) TestCancelRepairsPostgresZeroOIDs()
Zero table OIDs in the catalog (flows created before OIDs were stored there) are repaired from the live source during cancellation.
func (APITestSuite) TestCancelTableAdditionDuringSetupFlow ¶
func (s APITestSuite) TestCancelTableAdditionDuringSetupFlow()
func (APITestSuite) TestCancelTableAdditionRemoveAddRemove ¶
func (s APITestSuite) TestCancelTableAdditionRemoveAddRemove()
Tests that table addition cancellation doesn't get confused by the canceled table having a previous initial load and a previous successful addition then removal
func (APITestSuite) TestCancelTableAddition_NoRemovalAssumed ¶
func (s APITestSuite) TestCancelTableAddition_NoRemovalAssumed()
func (APITestSuite) TestCancelTableAddition_NoRemovalAssumedWithRemoval ¶
func (s APITestSuite) TestCancelTableAddition_NoRemovalAssumedWithRemoval()
func (APITestSuite) TestCancelTableAddition_WithRemoval ¶
func (s APITestSuite) TestCancelTableAddition_WithRemoval()
func (APITestSuite) TestClickHouseMirrorValidation_NoPrimaryKey ¶
func (s APITestSuite) TestClickHouseMirrorValidation_NoPrimaryKey()
func (APITestSuite) TestClickHouseMirrorValidation_NoPrimaryKey_ReplicaIdentityFull ¶
func (s APITestSuite) TestClickHouseMirrorValidation_NoPrimaryKey_ReplicaIdentityFull()
func (APITestSuite) TestClickHouseMirrorValidation_Pass ¶
func (s APITestSuite) TestClickHouseMirrorValidation_Pass()
func (APITestSuite) TestCreateCDCFlowAttachCanceledWorkflow ¶
func (s APITestSuite) TestCreateCDCFlowAttachCanceledWorkflow()
func (APITestSuite) TestCreateCDCFlowAttachConcurrentRequests ¶
func (s APITestSuite) TestCreateCDCFlowAttachConcurrentRequests()
func (APITestSuite) TestCreateCDCFlowAttachConcurrentRequestsToxi ¶
func (s APITestSuite) TestCreateCDCFlowAttachConcurrentRequestsToxi()
func (APITestSuite) TestCreateCDCFlowAttachExternalFlowEntry ¶
func (s APITestSuite) TestCreateCDCFlowAttachExternalFlowEntry()
func (APITestSuite) TestCreateCDCFlowAttachIdempotentAfterContinueAsNew ¶
func (s APITestSuite) TestCreateCDCFlowAttachIdempotentAfterContinueAsNew()
func (APITestSuite) TestCreateCDCFlowAttachSequentialRequests ¶
func (s APITestSuite) TestCreateCDCFlowAttachSequentialRequests()
func (APITestSuite) TestDoubleClickCancelTableAddition ¶
func (s APITestSuite) TestDoubleClickCancelTableAddition()
func (APITestSuite) TestDropCompleted ¶
func (s APITestSuite) TestDropCompleted()
func (APITestSuite) TestDropCompletedAndUnavailable ¶
func (s APITestSuite) TestDropCompletedAndUnavailable()
drop on completed mirror doesn't access peers, so should still drop immediately
func (APITestSuite) TestDropMissing ¶
func (s APITestSuite) TestDropMissing()
Simulate a mirror whose Temporal workflow no longer exists (e.g. completed or failed over 30 days ago and was already cleaned up).
func (APITestSuite) TestDropQRep ¶
func (s APITestSuite) TestDropQRep()
func (APITestSuite) TestEditTablesBeforeResync ¶
func (s APITestSuite) TestEditTablesBeforeResync()
func (APITestSuite) TestFlowStatusUpdate ¶
func (s APITestSuite) TestFlowStatusUpdate()
func (APITestSuite) TestGetColumnsNullability ¶
func (s APITestSuite) TestGetColumnsNullability()
func (APITestSuite) TestGetMirrorRowCounts ¶
func (s APITestSuite) TestGetMirrorRowCounts()
func (APITestSuite) TestGetMirrorRowCountsNonCDC ¶
func (s APITestSuite) TestGetMirrorRowCountsNonCDC()
func (APITestSuite) TestGetTablesExcludeViews ¶
func (s APITestSuite) TestGetTablesExcludeViews()
func (APITestSuite) TestGetTablesUnloggedNotMirrorable ¶
func (s APITestSuite) TestGetTablesUnloggedNotMirrorable()
func (APITestSuite) TestGetVersion ¶
func (s APITestSuite) TestGetVersion()
func (APITestSuite) TestMirrorValidation_InvalidTableMappings ¶
func (s APITestSuite) TestMirrorValidation_InvalidTableMappings()
func (APITestSuite) TestMirrorValidation_MissingSourceTable ¶
func (s APITestSuite) TestMirrorValidation_MissingSourceTable()
This is the canonical test that source validation is wired up. Specific validaton tests go as integration tests in connectors or flow/pkg.
func (APITestSuite) TestMySQLFlavorSwap ¶
func (s APITestSuite) TestMySQLFlavorSwap()
func (APITestSuite) TestPostgresDestinationValidation_ExtraColumnsOk ¶
func (s APITestSuite) TestPostgresDestinationValidation_ExtraColumnsOk()
func (APITestSuite) TestPostgresDestinationValidation_MissingColumns ¶
func (s APITestSuite) TestPostgresDestinationValidation_MissingColumns()
func (APITestSuite) TestPostgresDestinationValidation_MissingSchema ¶
func (s APITestSuite) TestPostgresDestinationValidation_MissingSchema()
func (APITestSuite) TestPostgresDestinationValidation_NonEmptyTableAllowedWithoutSnapshot ¶
func (s APITestSuite) TestPostgresDestinationValidation_NonEmptyTableAllowedWithoutSnapshot()
func (APITestSuite) TestPostgresDestinationValidation_NonEmptyTableBlocksSnapshot ¶
func (s APITestSuite) TestPostgresDestinationValidation_NonEmptyTableBlocksSnapshot()
func (APITestSuite) TestPostgresDestinationValidation_NumericPrecisionMismatch ¶
func (s APITestSuite) TestPostgresDestinationValidation_NumericPrecisionMismatch()
func (APITestSuite) TestPostgresDestinationValidation_NumericSuperset ¶
func (s APITestSuite) TestPostgresDestinationValidation_NumericSuperset()
func (APITestSuite) TestPostgresDestinationValidation_TypeCompatible ¶
func (s APITestSuite) TestPostgresDestinationValidation_TypeCompatible()
func (APITestSuite) TestPostgresDestinationValidation_UnboundedNumeric ¶
func (s APITestSuite) TestPostgresDestinationValidation_UnboundedNumeric()
func (APITestSuite) TestPostgresDestinationValidation_UserDefinedTypeMatch ¶
func (s APITestSuite) TestPostgresDestinationValidation_UserDefinedTypeMatch()
func (APITestSuite) TestPostgresDestinationValidation_UserDefinedTypeMismatch ¶
func (s APITestSuite) TestPostgresDestinationValidation_UserDefinedTypeMismatch()
func (APITestSuite) TestPostgresDestinationValidation_WithExcludedColumns ¶
func (s APITestSuite) TestPostgresDestinationValidation_WithExcludedColumns()
func (APITestSuite) TestPostgresValidation_Pass ¶
func (s APITestSuite) TestPostgresValidation_Pass()
func (APITestSuite) TestPostgresValidation_WrongPassword ¶
func (s APITestSuite) TestPostgresValidation_WrongPassword()
func (APITestSuite) TestQRep ¶
func (s APITestSuite) TestQRep()
func (APITestSuite) TestResetMirrorSequences ¶
func (s APITestSuite) TestResetMirrorSequences()
func (APITestSuite) TestResetMirrorSequences_EmptyTable ¶
func (s APITestSuite) TestResetMirrorSequences_EmptyTable()
func (APITestSuite) TestResetMirrorSequences_NonPgDestination ¶
func (s APITestSuite) TestResetMirrorSequences_NonPgDestination()
func (APITestSuite) TestResyncCompleted ¶
func (s APITestSuite) TestResyncCompleted()
func (APITestSuite) TestResyncFailed ¶
func (s APITestSuite) TestResyncFailed()
func (APITestSuite) TestResyncSourceTableMissing ¶
func (s APITestSuite) TestResyncSourceTableMissing()
func (APITestSuite) TestResyncTablesNotInPublication ¶
func (s APITestSuite) TestResyncTablesNotInPublication()
func (APITestSuite) TestResyncWithSnapshotConfigDuringSnapshot ¶
func (s APITestSuite) TestResyncWithSnapshotConfigDuringSnapshot()
TestResyncWithSnapshotConfigDuringSnapshot verifies that snapshot tuning parameters supplied via FlowConfigUpdate on a RESYNC state change issued while the flow is still in STATUS_SNAPSHOT (initial snapshot in-flight) are applied to the resynced flow and persisted in the catalog.
func (APITestSuite) TestResyncWithSnapshotConfigOnCompletedSnapshotOnlyPipe ¶
func (s APITestSuite) TestResyncWithSnapshotConfigOnCompletedSnapshotOnlyPipe()
TestResyncWithSnapshotConfigOnCompletedSnapshotOnlyPipe verifies that snapshot tuning parameters supplied via FlowConfigUpdate on a RESYNC state change are applied to an InitialSnapshotOnly flow that has already completed, and that the override is persisted in the catalog config.
func (APITestSuite) TestResyncWithSnapshotConfigOnPausedPipe ¶
func (s APITestSuite) TestResyncWithSnapshotConfigOnPausedPipe()
TestResyncWithSnapshotConfigOnPausedPipe verifies that snapshot tuning parameters supplied via FlowConfigUpdate on a RESYNC state change issued while the flow is PAUSED are applied to the resynced flow and persisted in the catalog.
func (APITestSuite) TestResyncWithSnapshotConfigOnRunningPipe ¶
func (s APITestSuite) TestResyncWithSnapshotConfigOnRunningPipe()
TestResyncWithSnapshotConfigOnRunningPipe verifies that snapshot tuning parameters (max_parallel_workers, num_tables_in_parallel, num_rows_per_partition, num_partitions_override) supplied via FlowConfigUpdate on a RESYNC state change are applied to the resynced flow and persisted in the catalog config.
func (APITestSuite) TestSchemaEndpoints ¶
func (s APITestSuite) TestSchemaEndpoints()
func (APITestSuite) TestScripts ¶
func (s APITestSuite) TestScripts()
func (APITestSuite) TestSettings ¶
func (s APITestSuite) TestSettings()
func (APITestSuite) TestSnapshotNullPartitionKey ¶
func (s APITestSuite) TestSnapshotNullPartitionKey()
func (APITestSuite) TestTableAdditionWithoutInitialLoad ¶
func (s APITestSuite) TestTableAdditionWithoutInitialLoad()
func (APITestSuite) TestTerminateDuringResyncDropFlow ¶
func (s APITestSuite) TestTerminateDuringResyncDropFlow()
TestTerminateDuringResyncDropFlow creates a pipe, resyncs it, then sends a TERMINATING signal while the DropFlowWorkflow is blocked dropping the replication slot. An idle open transaction on the source keeps the slot active (pg_drop_replication_slot blocks while the slot's xmin cannot advance). Once the TERMINATING signal is queued, the transaction is committed so the drop activity can finish — but the DropFlowWorkflow's signal goroutine has already set TerminateSignal, so it does NOT ContinueAsNew back to CDCFlowWorkflow and the pipe is fully torn down.
func (APITestSuite) TestTotalRowsSyncedByMirror ¶
func (s APITestSuite) TestTotalRowsSyncedByMirror()
func (APITestSuite) TestValidateCDCMirror_ServerIDInUseAcrossPeers ¶
func (s APITestSuite) TestValidateCDCMirror_ServerIDInUseAcrossPeers()
TestValidateCDCMirror_ServerIDInUseAcrossPeers verifies that a CDC mirror is rejected when its MySQL source peer pins a server_id that is already registered by another replica on the source database — here, a mirror streaming through a DIFFERENT peer pinning the same server_id. The cross-mirror peer-reuse check cannot relate the two peers, so only the replica registered on the source betrays the collision.
func (APITestSuite) TestValidateCDCMirror_ServerIDPeerReuse ¶
func (s APITestSuite) TestValidateCDCMirror_ServerIDPeerReuse()
TestValidateCDCMirror_ServerIDPeerReuse verifies that validating a CDC mirror whose MySQL source peer pins a fixed server_id is rejected when that peer already backs another streaming CDC mirror.
func (APITestSuite) TestValidateCDCMirror_ServerIDResync ¶
func (s APITestSuite) TestValidateCDCMirror_ServerIDResync()
TestValidateCDCMirror_ServerIDResync verifies that a running CDC mirror whose MySQL source peer pins a fixed server_id can be resynced. Resync validates the mirror before dropping it, while the mirror's own binlog connection is still registered on the source with the pinned server_id; validation must not mistake that registration for a foreign replica.
type BigQueryTestHelper ¶
type BigQueryTestHelper struct {
Config *protos.BigqueryConfig
ServiceAccount *utils.GcpServiceAccount
// contains filtered or unexported fields
}
func NewBigQueryTestHelper ¶
func NewBigQueryTestHelper(t *testing.T, datasetID string) (*BigQueryTestHelper, error)
NewBigQueryTestHelper creates a new BigQueryTestHelper.
func (*BigQueryTestHelper) CheckNull ¶
func (b *BigQueryTestHelper) CheckNull(ctx context.Context, tableName string, colName []string) (bool, error)
returns whether the function errors or there are no nulls
func (*BigQueryTestHelper) CountObjectsInGCSPath ¶
func (b *BigQueryTestHelper) CountObjectsInGCSPath(ctx context.Context, gcsPath string) (int, error)
CountObjectsInGCSPath counts the number of objects in a GCS path (gs://bucket/path/prefix)
func (*BigQueryTestHelper) CountRows ¶
CountRows(tableName) returns the number of rows in the given table.
func (*BigQueryTestHelper) CountRowsWithDataset ¶
func (*BigQueryTestHelper) DropDataset ¶
func (b *BigQueryTestHelper) DropDataset(ctx context.Context, datasetName string) error
DropDataset drops the dataset.
func (*BigQueryTestHelper) ExecuteAndProcessQuery ¶
func (b *BigQueryTestHelper) ExecuteAndProcessQuery(ctx context.Context, query string) (*model.QRecordBatch, error)
func (*BigQueryTestHelper) RecreateDataset ¶
func (b *BigQueryTestHelper) RecreateDataset(ctx context.Context) error
RecreateDataset recreates the dataset, i.e, deletes it if exists and creates it again.
func (*BigQueryTestHelper) RunInt64Query ¶
type ClickHouseMVManager ¶
type ClickHouseMVManager struct {
// contains filtered or unexported fields
}
ClickHouseMVManager manages materialized views for ClickHouse testing
func (*ClickHouseMVManager) CreateBadMV ¶
func (m *ClickHouseMVManager) CreateBadMV(ctx context.Context) error
CreateBadMV creates a materialized view on the configured table that will fail when data is inserted. Source schema is assumed to correspond to (id UInt64, val String).
type ClickHouseSuite ¶
type ClickHouseSuite struct {
// contains filtered or unexported fields
}
func (ClickHouseSuite) Catalog ¶
func (s ClickHouseSuite) Catalog() shared.CatalogPool
func (ClickHouseSuite) Connector ¶
func (s ClickHouseSuite) Connector() connectors.Connector
func (ClickHouseSuite) CountNonDeletedRows ¶
func (s ClickHouseSuite) CountNonDeletedRows(table string) (int, error)
CountNonDeletedRows returns the number of rows in table where _peerdb_is_deleted = 0. Unlike GetRows it does not use FINAL, so it works for engines that reject FINAL (e.g. MergeTree).
func (ClickHouseSuite) CreateRMTTable ¶
func (s ClickHouseSuite) CreateRMTTable(tableName string, columns []TestClickHouseColumn, orderingKey string) error
CreateRMTTable creates a ReplacingMergeTree table with the given name and columns.
func (ClickHouseSuite) CreateSlowInsertViaMV ¶
func (s ClickHouseSuite) CreateSlowInsertViaMV(tableName string, sleepSeconds int) (func(), error)
CreateSlowInsertViaMV attaches a materialized view to tableName that makes every insert into it takes sleepSeconds
func (ClickHouseSuite) DestinationConnector ¶
func (s ClickHouseSuite) DestinationConnector() connectors.Connector
func (ClickHouseSuite) DestinationTable ¶
func (s ClickHouseSuite) DestinationTable(table string) string
func (ClickHouseSuite) DropTable ¶
func (s ClickHouseSuite) DropTable(tableName string) error
func (ClickHouseSuite) GetRows ¶
func (s ClickHouseSuite) GetRows(table string, cols string) (*model.QRecordBatch, error)
func (ClickHouseSuite) IsCluster ¶
func (s ClickHouseSuite) IsCluster() bool
func (ClickHouseSuite) NewMVManager ¶
func (s ClickHouseSuite) NewMVManager(tableName string, suffix string) *ClickHouseMVManager
func (ClickHouseSuite) Peer ¶
func (s ClickHouseSuite) Peer() *protos.Peer
func (ClickHouseSuite) PeerForDatabase ¶
func (s ClickHouseSuite) PeerForDatabase(dbname string) *protos.Peer
func (ClickHouseSuite) S3Helper ¶
func (s ClickHouseSuite) S3Helper() *S3TestHelper
func (ClickHouseSuite) Source ¶
func (s ClickHouseSuite) Source() SuiteSource
func (ClickHouseSuite) Suffix ¶
func (s ClickHouseSuite) Suffix() string
func (ClickHouseSuite) T ¶
func (s ClickHouseSuite) T() *testing.T
func (ClickHouseSuite) Teardown ¶
func (s ClickHouseSuite) Teardown(ctx context.Context)
func (ClickHouseSuite) Test_Addition_Removal ¶
func (s ClickHouseSuite) Test_Addition_Removal()
func (ClickHouseSuite) Test_AvroNullableLax ¶
func (s ClickHouseSuite) Test_AvroNullableLax()
Test_AvroNullableLax tests PEERDB_AVRO_NULLABLE_LAX with multi-level inheritance and attnum gaps Need to modify code to trigger logging as the logging was added for the issue we were unable to reproduce
func (ClickHouseSuite) Test_Binary_Format_Base64 ¶
func (s ClickHouseSuite) Test_Binary_Format_Base64()
func (ClickHouseSuite) Test_Binary_Format_Hex ¶
func (s ClickHouseSuite) Test_Binary_Format_Hex()
func (ClickHouseSuite) Test_Binary_Format_Raw ¶
func (s ClickHouseSuite) Test_Binary_Format_Raw()
func (ClickHouseSuite) Test_CTID_Inherited_Table ¶
func (s ClickHouseSuite) Test_CTID_Inherited_Table()
func (ClickHouseSuite) Test_CTID_Inherited_Table_Extra_Columns ¶
func (s ClickHouseSuite) Test_CTID_Inherited_Table_Extra_Columns()
func (ClickHouseSuite) Test_CTID_Multi_Level_Inherited_Table ¶
func (s ClickHouseSuite) Test_CTID_Multi_Level_Inherited_Table()
func (ClickHouseSuite) Test_CTID_Multi_Level_Partitioned_Table ¶
func (s ClickHouseSuite) Test_CTID_Multi_Level_Partitioned_Table()
func (ClickHouseSuite) Test_CTID_Partitioned_Table ¶
func (s ClickHouseSuite) Test_CTID_Partitioned_Table()
func (ClickHouseSuite) Test_Chunking_Initial_Load_Parts_Per_Partition ¶
func (s ClickHouseSuite) Test_Chunking_Initial_Load_Parts_Per_Partition()
func (ClickHouseSuite) Test_CoalescingEngine ¶
func (s ClickHouseSuite) Test_CoalescingEngine()
func (ClickHouseSuite) Test_Column_Exclusion ¶
func (s ClickHouseSuite) Test_Column_Exclusion()
func (ClickHouseSuite) Test_Composite_PKey ¶
func (s ClickHouseSuite) Test_Composite_PKey()
func (ClickHouseSuite) Test_Destination_Type_Conversion ¶
func (s ClickHouseSuite) Test_Destination_Type_Conversion()
func (ClickHouseSuite) Test_Extra_CH_Columns ¶
func (s ClickHouseSuite) Test_Extra_CH_Columns()
func (ClickHouseSuite) Test_First_Row_Lag_Times_Recorded ¶
func (s ClickHouseSuite) Test_First_Row_Lag_Times_Recorded()
func (ClickHouseSuite) Test_Geometric_Types ¶
func (s ClickHouseSuite) Test_Geometric_Types()
func (ClickHouseSuite) Test_InfiniteTimestamp ¶
func (s ClickHouseSuite) Test_InfiniteTimestamp()
func (ClickHouseSuite) Test_InitialLoadOnly_No_Primary_Key ¶
func (s ClickHouseSuite) Test_InitialLoadOnly_No_Primary_Key()
func (ClickHouseSuite) Test_JSON_CH ¶
func (s ClickHouseSuite) Test_JSON_CH()
func (ClickHouseSuite) Test_JSON_Null ¶
func (s ClickHouseSuite) Test_JSON_Null()
func (ClickHouseSuite) Test_Large_Numeric ¶
func (s ClickHouseSuite) Test_Large_Numeric()
large NUMERICs (precision >76) are mapped to String on CH, test
func (ClickHouseSuite) Test_Large_Text_CDC ¶
func (s ClickHouseSuite) Test_Large_Text_CDC()
func (ClickHouseSuite) Test_Normalize_Metadata_With_Retry ¶
func (s ClickHouseSuite) Test_Normalize_Metadata_With_Retry()
Test_Normalize_Metadata_With_Retry tests the chunking normalization with a push to ClickHouse thrown in via renaming a target table.
func (ClickHouseSuite) Test_NullEngine ¶
func (s ClickHouseSuite) Test_NullEngine()
func (ClickHouseSuite) Test_NullableColumnSetting ¶
func (s ClickHouseSuite) Test_NullableColumnSetting()
func (ClickHouseSuite) Test_NullableMirrorSetting ¶
func (s ClickHouseSuite) Test_NullableMirrorSetting()
func (ClickHouseSuite) Test_Nullable_Schema_Change ¶
func (s ClickHouseSuite) Test_Nullable_Schema_Change()
func (ClickHouseSuite) Test_Nullable_Schema_Change_Replident_Full ¶
func (s ClickHouseSuite) Test_Nullable_Schema_Change_Replident_Full()
func (ClickHouseSuite) Test_Nullable_Schema_Change_Replident_Index ¶
func (s ClickHouseSuite) Test_Nullable_Schema_Change_Replident_Index()
old logic would mark pkey being nullable if replident index was used
func (ClickHouseSuite) Test_Numeric_Truncation_With_UnbNumAsString_FF ¶
func (s ClickHouseSuite) Test_Numeric_Truncation_With_UnbNumAsString_FF()
func (ClickHouseSuite) Test_Numeric_Truncation_Without_UnbNumAsString_FF ¶
func (s ClickHouseSuite) Test_Numeric_Truncation_Without_UnbNumAsString_FF()
func (ClickHouseSuite) Test_Offload_Partition_Ranges ¶
func (s ClickHouseSuite) Test_Offload_Partition_Ranges()
func (ClickHouseSuite) Test_PG_AlterTableAddColumnDefault ¶
func (s ClickHouseSuite) Test_PG_AlterTableAddColumnDefault()
Test_PG_AlterTableAddColumnDefault covers rows that existed before an ADD COLUMN with a DEFAULT.
func (ClickHouseSuite) Test_PG_AlterTableAddColumnDefaultUntranslated ¶
func (s ClickHouseSuite) Test_PG_AlterTableAddColumnDefaultUntranslated()
Test_PG_AlterTableAddColumnDefaultUntranslated covers defaults deliberately not carried over, plus one ClickHouse rejects outright. Neither may stall replication for the table.
func (ClickHouseSuite) Test_PG_Domain_Bytea ¶
func (s ClickHouseSuite) Test_PG_Domain_Bytea()
Test_PG_Domain_Bytea covers domains over bytea: table schema reports the base type, but pgoutput relation messages carry the domain's own oid
func (ClickHouseSuite) Test_PartitionBy ¶
func (s ClickHouseSuite) Test_PartitionBy()
func (ClickHouseSuite) Test_PartitionByExpr ¶
func (s ClickHouseSuite) Test_PartitionByExpr()
func (ClickHouseSuite) Test_Partition_By_CTID_With_Num_Partitions_Override ¶
func (s ClickHouseSuite) Test_Partition_By_CTID_With_Num_Partitions_Override()
func (ClickHouseSuite) Test_Partition_Key_Empty ¶
func (s ClickHouseSuite) Test_Partition_Key_Empty()
tests where panic happened when using custom partition key on empty table
func (ClickHouseSuite) Test_Partition_Key_Integer ¶
func (s ClickHouseSuite) Test_Partition_Key_Integer()
func (ClickHouseSuite) Test_Partition_Key_Null ¶
func (s ClickHouseSuite) Test_Partition_Key_Null()
edge case: min/max will be null, but null partition should still replicate all null rows
func (ClickHouseSuite) Test_Partition_Key_Timestamp ¶
func (s ClickHouseSuite) Test_Partition_Key_Timestamp()
func (ClickHouseSuite) Test_PgVector ¶
func (s ClickHouseSuite) Test_PgVector()
func (ClickHouseSuite) Test_PgVector_Version0 ¶
func (s ClickHouseSuite) Test_PgVector_Version0()
func (ClickHouseSuite) Test_Removal_Shared_Destination ¶
func (s ClickHouseSuite) Test_Removal_Shared_Destination()
Removing one of several source tables feeding a shared destination must keep that destination's raw rows & catalog schema mapping intact, otherwise the surviving source can no longer normalize into it. Disjoint id ranges keep the two sources' rows distinct under the destination's ReplacingMergeTree.
func (ClickHouseSuite) Test_Replident_Full_Unchanged_TOAST_Updates ¶
func (s ClickHouseSuite) Test_Replident_Full_Unchanged_TOAST_Updates()
func (ClickHouseSuite) Test_SchemaAsColumn ¶
func (s ClickHouseSuite) Test_SchemaAsColumn()
func (ClickHouseSuite) Test_Schema_Change_After_Resync_Cluster ¶
func (s ClickHouseSuite) Test_Schema_Change_After_Resync_Cluster()
func (ClickHouseSuite) Test_SkipSnapshotExport ¶
func (s ClickHouseSuite) Test_SkipSnapshotExport()
func (ClickHouseSuite) Test_Sync_Error_Cancels_InProgress_Normalize ¶
func (s ClickHouseSuite) Test_Sync_Error_Cancels_InProgress_Normalize()
func (ClickHouseSuite) Test_Time64 ¶
func (s ClickHouseSuite) Test_Time64()
func (ClickHouseSuite) Test_Types_CH ¶
func (s ClickHouseSuite) Test_Types_CH()
func (ClickHouseSuite) Test_Unbounded_Numeric_With_FF ¶
func (s ClickHouseSuite) Test_Unbounded_Numeric_With_FF()
func (ClickHouseSuite) Test_Unbounded_Numeric_Without_FF ¶
func (s ClickHouseSuite) Test_Unbounded_Numeric_Without_FF()
func (ClickHouseSuite) Test_Unprivileged_Postgres_Columns ¶
func (s ClickHouseSuite) Test_Unprivileged_Postgres_Columns()
func (ClickHouseSuite) Test_Update_PKey_Env_Disabled ¶
func (s ClickHouseSuite) Test_Update_PKey_Env_Disabled()
func (ClickHouseSuite) Test_Update_PKey_Env_Enabled ¶
func (s ClickHouseSuite) Test_Update_PKey_Env_Enabled()
func (ClickHouseSuite) Test_ValidatePartitionByExpression ¶
func (s ClickHouseSuite) Test_ValidatePartitionByExpression()
func (ClickHouseSuite) Test_WeirdTable_Dash ¶
func (s ClickHouseSuite) Test_WeirdTable_Dash()
func (ClickHouseSuite) Test_WeirdTable_Keyword ¶
func (s ClickHouseSuite) Test_WeirdTable_Keyword()
func (ClickHouseSuite) Test_WeirdTable_MixedCase ¶
func (s ClickHouseSuite) Test_WeirdTable_MixedCase()
func (ClickHouseSuite) Test_WeirdTable_Question ¶
func (s ClickHouseSuite) Test_WeirdTable_Question()
func (ClickHouseSuite) WeirdTable ¶
func (s ClickHouseSuite) WeirdTable(tableName string)
type CockroachDBSource ¶
type CockroachDBSource struct {
// contains filtered or unexported fields
}
func SetupCockroachDB ¶
func SetupCockroachDB(t *testing.T, suffix string) (*CockroachDBSource, error)
func SetupCockroachDBWithConfig ¶
func SetupCockroachDBWithConfig( t *testing.T, suffix string, config *protos.CockroachDBConfig, ) (*CockroachDBSource, error)
func (*CockroachDBSource) AdminConn ¶
func (s *CockroachDBSource) AdminConn() *pgx.Conn
func (*CockroachDBSource) CockroachDBConnector ¶
func (s *CockroachDBSource) CockroachDBConnector() *conncockroachdb.CockroachDBConnector
func (*CockroachDBSource) Config ¶
func (s *CockroachDBSource) Config() *protos.CockroachDBConfig
func (*CockroachDBSource) Connector ¶
func (s *CockroachDBSource) Connector() connectors.Connector
func (*CockroachDBSource) GeneratePeer ¶
func (s *CockroachDBSource) GeneratePeer(t *testing.T) *protos.Peer
func (*CockroachDBSource) GetRows ¶
func (s *CockroachDBSource) GetRows(ctx context.Context, suffix, table, cols string) (*model.QRecordBatch, error)
type FlowConnectionGenerationConfig ¶
type FlowConnectionGenerationConfig struct {
FlowJobName string
TableNameMapping map[string]string
Destination string
TableMappings []*protos.TableMapping
SoftDelete bool
}
func (*FlowConnectionGenerationConfig) GenerateFlowConnectionConfigs ¶
func (c *FlowConnectionGenerationConfig) GenerateFlowConnectionConfigs(s Suite) *protos.FlowConnectionConfigs
type Generic ¶
type Generic struct {
GenericSuite
}
func (Generic) Test_Custom_Replication_Slot_Starting_With_Numbers_CDC_Only ¶
func (s Generic) Test_Custom_Replication_Slot_Starting_With_Numbers_CDC_Only()
func (Generic) Test_Inheritance_Table_With_Dynamic_Setting ¶
func (s Generic) Test_Inheritance_Table_With_Dynamic_Setting()
func (Generic) Test_Inheritance_Table_Without_Dynamic_Setting ¶
func (s Generic) Test_Inheritance_Table_Without_Dynamic_Setting()
func (Generic) Test_Initial_Custom_Partition ¶
func (s Generic) Test_Initial_Custom_Partition()
func (Generic) Test_Partitioned_Table ¶
func (s Generic) Test_Partitioned_Table()
func (Generic) Test_Partitioned_Table_With_Different_Column_Ordering ¶
func (s Generic) Test_Partitioned_Table_With_Different_Column_Ordering()
func (Generic) Test_Partitioned_Table_Without_Publish_Via_Partition_Root ¶
func (s Generic) Test_Partitioned_Table_Without_Publish_Via_Partition_Root()
func (Generic) Test_Schema_Change_Drop_Consecutive_Columns ¶
func (s Generic) Test_Schema_Change_Drop_Consecutive_Columns()
func (Generic) Test_Schema_Change_Lost_Column_Bug ¶
func (s Generic) Test_Schema_Change_Lost_Column_Bug()
Test_Schema_Change_Lost_Column_Bug addresses a race condition where a column added to the source table without a subsequent DML operation can be "lost" during schema evolution. The scenario:
- CDC mirror is running with a small batch size
- ALTER TABLE adds good_column + INSERT (relation message sent for good_column)
- ALTER TABLE adds lost_column (NO subsequent DML, so no relation message yet)
- PeerDB syncs: adds good_column to destination, then applySchemaDelta updates catalog with latest source db schema (which includes lost_column)
- Next INSERT triggers relation message, adds lost_column to schemaDeltas
- PeerDB compares against catalog (which has lost_column) -> no delta detected
- Error: lost_column doesn't exist on destination but we try to insert data for it
func (Generic) Test_Schema_Changes_Cutoff_Bug ¶
func (s Generic) Test_Schema_Changes_Cutoff_Bug()
func (Generic) Test_Simple_Flow ¶
func (s Generic) Test_Simple_Flow()
func (Generic) Test_Simple_Schema_Changes ¶
func (s Generic) Test_Simple_Schema_Changes()
type GenericSuite ¶
type MongoSource ¶
type MongoSource struct {
// contains filtered or unexported fields
}
func SetupMongo ¶
func SetupMongo(t *testing.T, suffix string) (*MongoSource, error)
func (*MongoSource) AdminClient ¶
func (s *MongoSource) AdminClient() *mongo.Client
func (*MongoSource) Config ¶
func (s *MongoSource) Config() *protos.MongoConfig
func (*MongoSource) Connector ¶
func (s *MongoSource) Connector() connectors.Connector
func (*MongoSource) GeneratePeer ¶
func (s *MongoSource) GeneratePeer(t *testing.T) *protos.Peer
func (*MongoSource) GetRows ¶
func (s *MongoSource) GetRows(ctx context.Context, suffix, table, cols string) (*model.QRecordBatch, error)
type MySQLTestContainerConfig ¶
type MySQLTestContainerConfig struct {
Image string
// ExtraServerFlags are appended to a common small-footprint server flag base (which already
// includes flavor-aware low-resource flags), e.g. "--binlog-row-event-fragment-threshold=1024".
ExtraServerFlags []string
Flavor protos.MySqlFlavor
ReplicationMechanism protos.MySqlReplicationMechanism
}
MySQLTestContainerConfig parameterizes a throwaway MySQL/MariaDB testcontainer source.
type MySqlSource ¶
type MySqlSource struct {
*connmysql.MySqlConnector
Config *protos.MySqlConfig
// peer name, defaults to "mysql"
Name string
}
func SetupMySQL ¶
func SetupMySQL(t *testing.T, suffix string) (*MySqlSource, error)
func SetupMySQLTestContainerSource ¶
func SetupMySQLTestContainerSource( t *testing.T, namePrefix string, cfg MySQLTestContainerConfig, ) (*MySqlSource, string)
SetupMySQLTestContainerSource starts a throwaway MySQL/MariaDB server in a testcontainer and returns a MySqlSource pointed at it, plus the generated suffix used for its e2e database. It registers container and database cleanup on t. Use it to replace the shared CI source with an isolated server a test can reconfigure (custom image, extra server flags) - typically by assigning the returned source's peer to flowConnConfig.SourceName.
func (*MySqlSource) Connector ¶
func (s *MySqlSource) Connector() connectors.Connector
func (*MySqlSource) GeneratePeer ¶
func (s *MySqlSource) GeneratePeer(t *testing.T) *protos.Peer
type PgSuite ¶
type PgSuite interface {
Suite
Connector() *connpostgres.PostgresConnector
}
type PostgresSource ¶
type PostgresSource struct {
*connpostgres.PostgresConnector
}
func SetupPostgres ¶
func SetupPostgres(t *testing.T, suffix string) (*PostgresSource, error)
func SetupPostgresWithToxiproxy ¶
func SetupPostgresWithToxiproxy(t *testing.T, suffix string, port uint32) (*PostgresSource, *tp.Proxy, error)
SetupPostgresWithToxiproxy creates a PostgreSQL source that connects through Toxiproxy
func (*PostgresSource) Connector ¶
func (s *PostgresSource) Connector() connectors.Connector
func (*PostgresSource) GeneratePeer ¶
func (s *PostgresSource) GeneratePeer(t *testing.T) *protos.Peer
func (*PostgresSource) GetRows ¶
func (s *PostgresSource) GetRows(ctx context.Context, suffix string, table string, cols string) (*model.QRecordBatch, error)
func (*PostgresSource) GetRowsOnly ¶
func (s *PostgresSource) GetRowsOnly(ctx context.Context, suffix string, table string, cols string) (*model.QRecordBatch, error)
to avoid fetching rows from "child" tables ala Postgres table inheritance
func (*PostgresSource) Query ¶
func (s *PostgresSource) Query(ctx context.Context, query string) (*model.QRecordBatch, error)
type RowSource ¶
type RowSource interface {
Suite
GetRows(table, cols string) (*model.QRecordBatch, error)
}
type S3PeerCredentials ¶
type S3TestHelper ¶
type S3TestHelper struct {
S3Config *protos.S3Config
BucketName string
// contains filtered or unexported fields
}
func NewS3TestHelper ¶
func NewS3TestHelper(ctx context.Context, s3environment S3Environment) (*S3TestHelper, error)
func (*S3TestHelper) CleanUp ¶
func (h *S3TestHelper) CleanUp(ctx context.Context) error
Delete all generated objects during the test
func (*S3TestHelper) ListAllFiles ¶
func (h *S3TestHelper) ListAllFiles( ctx context.Context, jobName string, ) ([]s3types.Object, error)
List all files from the S3 bucket. returns as a list of S3Objects.
type SnowflakeTestHelper ¶
type SnowflakeTestHelper struct {
// config is the Snowflake config.
Config *protos.SnowflakeConfig
// TestSchemaName is the schema to use for testing.
TestSchemaName string
// dbName is the database used for testing.
TestDatabaseName string
// contains filtered or unexported fields
}
func NewSnowflakeTestHelper ¶
func NewSnowflakeTestHelper(t *testing.T) (*SnowflakeTestHelper, error)
func (*SnowflakeTestHelper) CheckIsDeleted ¶
func (s *SnowflakeTestHelper) CheckIsDeleted(ctx context.Context, query string) error
func (*SnowflakeTestHelper) CheckSyncedAt ¶
func (s *SnowflakeTestHelper) CheckSyncedAt(ctx context.Context, query string) error
func (*SnowflakeTestHelper) Cleanup ¶
func (s *SnowflakeTestHelper) Cleanup(ctx context.Context) error
Cleanup drops the database.
func (*SnowflakeTestHelper) CountNonNullRows ¶
func (s *SnowflakeTestHelper) CountNonNullRows(ctx context.Context, tableName string, columnName string) (int64, error)
CountRows(tableName) returns the non-null number of rows in the given table.
func (*SnowflakeTestHelper) CountRows ¶
CountRows(tableName) returns the number of rows in the given table.
func (*SnowflakeTestHelper) CountSRIDs ¶
func (*SnowflakeTestHelper) ExecuteAndProcessQuery ¶
func (s *SnowflakeTestHelper) ExecuteAndProcessQuery(ctx context.Context, query string) (*model.QRecordBatch, error)
func (*SnowflakeTestHelper) RunCommand ¶
func (s *SnowflakeTestHelper) RunCommand(ctx context.Context, command string) error
RunCommand runs the given command.
func (*SnowflakeTestHelper) RunIntQuery ¶
runs a query that returns an int result
type SuiteSource ¶
type SuiteSource interface {
Teardown(t *testing.T, ctx context.Context, suffix string)
GeneratePeer(t *testing.T) *protos.Peer
Connector() connectors.Connector
Exec(ctx context.Context, sql string, args ...any) error
GetRows(ctx context.Context, suffix, table, cols string) (*model.QRecordBatch, error)
}
type TestClickHouseColumn ¶
type WorkflowRun ¶
type WorkflowRun struct {
client.WorkflowRun
// contains filtered or unexported fields
}
func ExecuteDropFlow ¶
func ExecuteDropFlow(ctx context.Context, tc client.Client, config *protos.FlowConnectionConfigs) WorkflowRun
func ExecutePeerflow ¶
func ExecutePeerflow(t *testing.T, tc client.Client, config *protos.FlowConnectionConfigs) WorkflowRun
func ExecuteWorkflow ¶
func ExecuteWorkflow(ctx context.Context, tc client.Client, taskQueueID shared.TaskQueueID, wf any, args ...any) WorkflowRun
func GetPeerflow ¶
func GetPeerflow(ctx context.Context, catalog shared.CatalogPool, tc client.Client, flowName string) (WorkflowRun, error)
func RunQRepFlowWorkflow ¶
func RunQRepFlowWorkflow(t *testing.T, tc client.Client, config *protos.QRepConfig) WorkflowRun
func (WorkflowRun) Cancel ¶
func (env WorkflowRun) Cancel(ctx context.Context)
func (WorkflowRun) GetFlowStatus ¶
func (env WorkflowRun) GetFlowStatus(t *testing.T) protos.FlowStatus
func (WorkflowRun) Query ¶
func (env WorkflowRun) Query(ctx context.Context, queryType string, args ...any) (converter.EncodedValue, error)
Source Files
¶
- api_suite.go
- api_suite_cancel_table_addition.go
- api_suite_drop_flow.go
- api_suite_flow_status.go
- api_suite_mysql_server_id_collisions.go
- bigquery_helper.go
- clickhouse.go
- clickhouse_mv.go
- clickhouse_suite.go
- cockroachdb.go
- congen.go
- eventhub.go
- generic_suite.go
- mongo.go
- mysql.go
- pg.go
- s3_helper.go
- snowflake_helper.go
- switchboard.go
- test_utils.go