e2e

package
v0.0.0-...-eab1039 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 10, 2026 License: AGPL-3.0 Imports: 89 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func AddSuffix

func AddSuffix(s Suite, str string) string

func AttachSchema

func AttachSchema(s Suite, table string) string

func CatalogTestAccessPool

func CatalogTestAccessPool() (*pgxpool.Pool, error)

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 CreatePeer

func CreatePeer(t *testing.T, peer *protos.Peer)

func CreateQRepWorkflowConfig

func CreateQRepWorkflowConfig(
	t *testing.T,
	flowJobName string,
	sourceTable string,
	dstTable string,
	query string,
	dest string,
	stagingPath string,
	setupDst bool,
	syncedAtCol string,
	isDeletedCol string,
) *protos.QRepConfig

func CreateTableForQRep

func CreateTableForQRep(ctx context.Context, conn *pgx.Conn, suffix string, tableName string) error

func DeleteEventhub

func DeleteEventhub(
	ctx context.Context,
	eventhubName string,
	eventhubConfig *protos.EventHubConfig,
) error

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 EnvTrue

func EnvTrue(t *testing.T, env WorkflowRun, val bool)

func EnvWaitFor

func EnvWaitFor(t *testing.T, env WorkflowRun, timeout time.Duration, reason string, f func() bool)

func EnvWaitForCount

func EnvWaitForCount(
	env WorkflowRun,
	suite RowSource,
	reason string,
	dstTable string,
	cols string,
	expectedCount int,
)

func EnvWaitForEqualTables

func EnvWaitForEqualTables(
	env WorkflowRun,
	suite RowSource,
	reason string,
	table string,
	cols string,
)

func EnvWaitForEqualTablesWithNames

func EnvWaitForEqualTablesWithNames(
	env WorkflowRun,
	suite RowSource,
	reason string,
	srcTable string,
	dstTable string,
	cols string,
)

func EnvWaitForEqualTablesWithNames_Only

func EnvWaitForEqualTablesWithNames_Only(
	env WorkflowRun,
	suite RowSource,
	reason string,
	srcTable string,
	dstTable 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 GeneratePostgresPeer

func GeneratePostgresPeer(t *testing.T) *protos.Peer

func GetLogCount

func GetLogCount(ctx context.Context, catalog shared.CatalogPool, flowJobName, errorType, pattern string) (int, error)

func GetOwnersSchema

func GetOwnersSchema() *types.QRecordSchema

func GetOwnersSelectorStringsSF

func GetOwnersSelectorStringsSF() [2]string

func GetPostgresToxicProxy

func GetPostgresToxicProxy(t *testing.T, suffix string, port uint32) (*tp.Proxy, error)

GetPostgresToxicProxy gets or creates the PostgreSQL proxy

func GetTestDatabase

func GetTestDatabase(suffix string) string

func InitToxiproxy

func InitToxiproxy() error

InitToxiproxy initializes the Toxiproxy client (singleton pattern)

func InsertScript

func InsertScript(t *testing.T, name string, lang string, source string)

InsertScript registers a transform script in the catalog's public.scripts table.

func NewApiClient

func NewApiClient() (protos.FlowServiceClient, error)

func NewTemporalClient

func NewTemporalClient(t *testing.T) client.Client

func PopulateSourceTable

func PopulateSourceTable(ctx context.Context, conn *pgx.Conn, suffix string, tableName string, rowCount int) error

func RequireEmptyDestinationTable

func RequireEmptyDestinationTable(suite RowSource, dstTable string, cols string)

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 RequireEqualTables(suite RowSource, table string, cols string)

func RequireEqualTablesWithNames

func RequireEqualTablesWithNames(suite RowSource, srcTable string, dstTable string, cols string)

func RevokePermissionForTableColumns

func RevokePermissionForTableColumns(ctx context.Context, conn *pgx.Conn, tableIdentifier string, selectedColumns []string) error

func RunApiSuite

func RunApiSuite[TSource SuiteSource](
	t *testing.T,
	setup func(*testing.T, string) (TSource, error),
)

func RunPsql

func RunPsql(t *testing.T, peer string, extraArgs ...string) (string, error)

func Schema

func Schema(s Suite) string

func SetupCDCFlowStatusQuery

func SetupCDCFlowStatusQuery(t *testing.T, env WorkflowRun, config *protos.FlowConnectionConfigs)

func SetupClickHouseSuite

func SetupClickHouseSuite[TSource SuiteSource](
	t *testing.T,
	cluster bool,
	setupSource func(*testing.T) (TSource, string, error),
) func(*testing.T) ClickHouseSuite

func SetupGenericSuite

func SetupGenericSuite[T GenericSuite](f func(t *testing.T) T) func(t *testing.T) Generic

func SignalWorkflow

func SignalWorkflow[T any](ctx context.Context, env WorkflowRun, signal model.TypedSignal[T], value T)

func SwitchboardDSN

func SwitchboardDSN(peer string, options map[string]string) string

func TableMappings

func TableMappings(s GenericSuite, tables ...string) []*protos.TableMapping

func TearDownPostgres

func TearDownPostgres(ctx context.Context, s PgSuite)

Types

type APITestSuite

type APITestSuite struct {
	protos.FlowServiceClient
	// contains filtered or unexported fields
}

func (APITestSuite) Connector

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

func (b *BigQueryTestHelper) CountRows(ctx context.Context, tableName string) (int, error)

CountRows(tableName) returns the number of rows in the given table.

func (*BigQueryTestHelper) CountRowsWithDataset

func (b *BigQueryTestHelper) CountRowsWithDataset(ctx context.Context, dataset, tableName string, nonNullCol string) (int, error)

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

func (b *BigQueryTestHelper) RunInt64Query(ctx context.Context, query string) (int64, error)

func (*BigQueryTestHelper) SelectRow

func (b *BigQueryTestHelper) SelectRow(ctx context.Context, tableName string, cols ...string) ([]bigquery.Value, error)

check if NaN, Inf double values are null

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).

func (*ClickHouseMVManager) DropBadMV

func (m *ClickHouseMVManager) DropBadMV(ctx context.Context) error

DropBadMV removes the materialized view and its target table

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 (*CockroachDBSource) Connector

func (s *CockroachDBSource) Connector() connectors.Connector

func (*CockroachDBSource) Exec

func (s *CockroachDBSource) Exec(ctx context.Context, sql string, args ...any) error

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)

func (*CockroachDBSource) Teardown

func (s *CockroachDBSource) Teardown(t *testing.T, ctx context.Context, suffix string)

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:

  1. CDC mirror is running with a small batch size
  2. ALTER TABLE adds good_column + INSERT (relation message sent for good_column)
  3. ALTER TABLE adds lost_column (NO subsequent DML, so no relation message yet)
  4. PeerDB syncs: adds good_column to destination, then applySchemaDelta updates catalog with latest source db schema (which includes lost_column)
  5. Next INSERT triggers relation message, adds lost_column to schemaDeltas
  6. PeerDB compares against catalog (which has lost_column) -> no delta detected
  7. 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 GenericSuite interface {
	RowSource
	Peer() *protos.Peer
	DestinationConnector() connectors.Connector
	DestinationTable(table string) string
}

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) Exec

func (s *MongoSource) Exec(ctx context.Context, sql string, args ...any) error

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)

func (*MongoSource) Teardown

func (s *MongoSource) Teardown(t *testing.T, ctx context.Context, suffix string)

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) Exec

func (s *MySqlSource) Exec(ctx context.Context, sql string, args ...any) error

func (*MySqlSource) GeneratePeer

func (s *MySqlSource) GeneratePeer(t *testing.T) *protos.Peer

func (*MySqlSource) GetRows

func (s *MySqlSource) GetRows(ctx context.Context, suffix string, table string, cols string) (*model.QRecordBatch, error)

func (*MySqlSource) Teardown

func (s *MySqlSource) Teardown(t *testing.T, ctx context.Context, suffix string)

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) Exec

func (s *PostgresSource) Exec(ctx context.Context, sql string, args ...any) error

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)

func (*PostgresSource) Teardown

func (s *PostgresSource) Teardown(t *testing.T, ctx context.Context, suffix string)

type RowSource

type RowSource interface {
	Suite
	GetRows(table, cols string) (*model.QRecordBatch, error)
}

type S3Environment

type S3Environment int
const (
	Aws S3Environment = iota
	Gcs
	Minio
	MinioTls
)

type S3PeerCredentials

type S3PeerCredentials struct {
	AccessKeyID     string `json:"accessKeyId"`
	SecretAccessKey string `json:"secretAccessKey"`
	AwsRoleArn      string `json:"awsRoleArn"`
	SessionToken    string `json:"sessionToken"`
	Region          string `json:"region"`
	Endpoint        string `json:"endpoint"`
}

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) CheckNull

func (s *SnowflakeTestHelper) CheckNull(ctx context.Context, tableName string, colNames []string) (bool, 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

func (s *SnowflakeTestHelper) CountRows(ctx context.Context, tableName string) (int64, error)

CountRows(tableName) returns the number of rows in the given table.

func (*SnowflakeTestHelper) CountSRIDs

func (s *SnowflakeTestHelper) CountSRIDs(ctx context.Context, tableName string, columnName string) (int64, error)

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

func (s *SnowflakeTestHelper) RunIntQuery(ctx context.Context, query string) (int, error)

runs a query that returns an int result

type Suite

type Suite interface {
	e2eshared.Suite
	T() *testing.T
	Suffix() string
	Source() SuiteSource
}

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 TestClickHouseColumn struct {
	Name string
	Type string
}

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) Error

func (env WorkflowRun) Error(ctx context.Context) error

func (WorkflowRun) Finished

func (env WorkflowRun) Finished(ctx context.Context) bool

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)

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL