Documentation ¶
Overview ¶
Package replication is a generated GoMock package.
Package replication is a generated GoMock package.
Package replication is a generated GoMock package.
Package replication is a generated GoMock package.
Package replication is a generated GoMock package.
Package replication is a generated GoMock package.
Package replication is a generated GoMock package.
Package replication is a generated GoMock package.
Package replication is a generated GoMock package.
Package replication is a generated GoMock package.
Package replication is a generated GoMock package.
Package replication is a generated GoMock package.
Package replication is a generated GoMock package.
Package replication is a generated GoMock package.
Index ¶
- Constants
- Variables
- func ClusterIDToClusterNameShardCount(allClusterInfo map[string]cluster.ClusterInformation, clusterID int32) (string, int32, error)
- func IsRetryableError(err error) bool
- func IsStreamError(err error) bool
- func NewExecutableTaskConverter(processToolBox ProcessToolBox) *executableTaskConverterImpl
- func NewExecutionManagerDLQWriter(executionManager ExecutionManager) *executionManagerDLQWriter
- func NewReceiverFlowControl(signals map[enums.TaskPriority]FlowControlSignalProvider, ...) *streamReceiverFlowControllerImpl
- func NewSequentialBatchableTaskQueue(task TrackableExecutableTask, ...) ctasks.SequentialTaskQueue[TrackableExecutableTask]
- func NewSequentialTaskQueue(task TrackableExecutableTask) ctasks.SequentialTaskQueue[TrackableExecutableTask]
- func NewTaskProcessorManager(config *configs.Config, shard shard.Context, engine shard.Engine, ...) *taskProcessorManagerImpl
- func SequentialTaskQueueCompareLess(this TrackableExecutableTask, that TrackableExecutableTask) bool
- func TaskOperationTag(replicationTask *replicationspb.ReplicationTask) string
- func ValidateTasksHaveSamePriority(messageBatchPriority enums.TaskPriority, ...) error
- func WorkflowKeyHashFn(item interface{}) uint32
- func WrapEventLoop(originalEventLoop func() error, streamStopper func(), logger log.Logger, ...)
- type AckManager
- type BatchableTask
- type BiDirectionStream
- type BiDirectionStreamClient
- type BiDirectionStreamClientProvider
- type BiDirectionStreamImpl
- type ClusterChannelKey
- type ClusterShardKey
- type ClusterShardKeyPair
- type DLQHandler
- type DLQWriteRequest
- type DLQWriter
- type DLQWriterAdapter
- type EagerNamespaceRefresher
- type ExecutableActivityStateTask
- func (e *ExecutableActivityStateTask) BatchWith(incomingTask BatchableTask) (TrackableExecutableTask, bool)
- func (e *ExecutableActivityStateTask) CanBatch() bool
- func (e *ExecutableActivityStateTask) Execute() error
- func (e *ExecutableActivityStateTask) HandleErr(err error) error
- func (e *ExecutableActivityStateTask) MarkPoisonPill() error
- func (e *ExecutableActivityStateTask) MarkUnbatchable()
- func (e *ExecutableActivityStateTask) QueueID() interface{}
- type ExecutableHistoryTask
- func (e *ExecutableHistoryTask) BatchWith(incomingTask BatchableTask) (TrackableExecutableTask, bool)
- func (e *ExecutableHistoryTask) CanBatch() bool
- func (e *ExecutableHistoryTask) Execute() error
- func (e *ExecutableHistoryTask) HandleErr(err error) error
- func (e *ExecutableHistoryTask) MarkPoisonPill() error
- func (e *ExecutableHistoryTask) MarkUnbatchable()
- func (e *ExecutableHistoryTask) QueueID() interface{}
- type ExecutableNoopTask
- type ExecutableTask
- type ExecutableTaskConverter
- type ExecutableTaskImpl
- func (e *ExecutableTaskImpl) Abort()
- func (e *ExecutableTaskImpl) Ack()
- func (e *ExecutableTaskImpl) Attempt() int
- func (e *ExecutableTaskImpl) Cancel()
- func (e *ExecutableTaskImpl) DeleteWorkflow(ctx context.Context, workflowKey definition.WorkflowKey) (retError error)
- func (e *ExecutableTaskImpl) GetNamespaceInfo(ctx context.Context, namespaceID string) (string, bool, error)
- func (e *ExecutableTaskImpl) IsRetryableError(err error) bool
- func (e *ExecutableTaskImpl) Nack(err error)
- func (e *ExecutableTaskImpl) Reschedule()
- func (e *ExecutableTaskImpl) Resend(ctx context.Context, remoteCluster string, ...) (bool, error)
- func (e *ExecutableTaskImpl) RetryPolicy() backoff.RetryPolicy
- func (e *ExecutableTaskImpl) SourceClusterName() string
- func (e *ExecutableTaskImpl) State() ctasks.State
- func (e *ExecutableTaskImpl) TaskCreationTime() time.Time
- func (e *ExecutableTaskImpl) TaskID() int64
- func (e *ExecutableTaskImpl) TerminalState() bool
- type ExecutableTaskTracker
- type ExecutableTaskTrackerImpl
- func (t *ExecutableTaskTrackerImpl) Cancel()
- func (t *ExecutableTaskTrackerImpl) LowWatermark() *WatermarkInfo
- func (t *ExecutableTaskTrackerImpl) Size() int
- func (t *ExecutableTaskTrackerImpl) TrackTasks(exclusiveHighWatermarkInfo WatermarkInfo, tasks ...TrackableExecutableTask) []TrackableExecutableTask
- type ExecutableUnknownTask
- type ExecutableWorkflowStateTask
- type ExecutionManager
- type FlowControlSignal
- type FlowControlSignalProvider
- type MockAckManager
- func (m *MockAckManager) ConvertTask(ctx context.Context, task tasks.Task) (*repication.ReplicationTask, error)
- func (m *MockAckManager) EXPECT() *MockAckManagerMockRecorder
- func (m *MockAckManager) GetMaxTaskInfo() (int64, time.Time)
- func (m *MockAckManager) GetReplicationTasksIter(ctx context.Context, pollingCluster string, ...) (collection.Iterator[tasks.Task], error)
- func (m *MockAckManager) GetTask(ctx context.Context, taskInfo *repication.ReplicationTaskInfo) (*repication.ReplicationTask, error)
- func (m *MockAckManager) GetTasks(ctx context.Context, pollingCluster string, queryMessageID int64) (*repication.ReplicationMessages, error)
- func (m *MockAckManager) NotifyNewTasks(tasks []tasks.Task)
- func (m *MockAckManager) SubscribeNotification() (<-chan struct{}, string)
- func (m *MockAckManager) UnsubscribeNotification(arg0 string)
- type MockAckManagerMockRecorder
- func (mr *MockAckManagerMockRecorder) ConvertTask(ctx, task interface{}) *gomock.Call
- func (mr *MockAckManagerMockRecorder) GetMaxTaskInfo() *gomock.Call
- func (mr *MockAckManagerMockRecorder) GetReplicationTasksIter(ctx, pollingCluster, minInclusiveTaskID, maxExclusiveTaskID interface{}) *gomock.Call
- func (mr *MockAckManagerMockRecorder) GetTask(ctx, taskInfo interface{}) *gomock.Call
- func (mr *MockAckManagerMockRecorder) GetTasks(ctx, pollingCluster, queryMessageID interface{}) *gomock.Call
- func (mr *MockAckManagerMockRecorder) NotifyNewTasks(tasks interface{}) *gomock.Call
- func (mr *MockAckManagerMockRecorder) SubscribeNotification() *gomock.Call
- func (mr *MockAckManagerMockRecorder) UnsubscribeNotification(arg0 interface{}) *gomock.Call
- type MockBatchableTask
- func (m *MockBatchableTask) Abort()
- func (m *MockBatchableTask) Ack()
- func (m *MockBatchableTask) BatchWith(task BatchableTask) (TrackableExecutableTask, bool)
- func (m *MockBatchableTask) CanBatch() bool
- func (m *MockBatchableTask) Cancel()
- func (m *MockBatchableTask) EXPECT() *MockBatchableTaskMockRecorder
- func (m *MockBatchableTask) Execute() error
- func (m *MockBatchableTask) HandleErr(err error) error
- func (m *MockBatchableTask) IsRetryableError(err error) bool
- func (m *MockBatchableTask) MarkPoisonPill() error
- func (m *MockBatchableTask) MarkUnbatchable()
- func (m *MockBatchableTask) Nack(err error)
- func (m *MockBatchableTask) QueueID() interface{}
- func (m *MockBatchableTask) Reschedule()
- func (m *MockBatchableTask) RetryPolicy() backoff.RetryPolicy
- func (m *MockBatchableTask) SourceClusterName() string
- func (m *MockBatchableTask) State() tasks.State
- func (m *MockBatchableTask) TaskCreationTime() time.Time
- func (m *MockBatchableTask) TaskID() int64
- type MockBatchableTaskMockRecorder
- func (mr *MockBatchableTaskMockRecorder) Abort() *gomock.Call
- func (mr *MockBatchableTaskMockRecorder) Ack() *gomock.Call
- func (mr *MockBatchableTaskMockRecorder) BatchWith(task interface{}) *gomock.Call
- func (mr *MockBatchableTaskMockRecorder) CanBatch() *gomock.Call
- func (mr *MockBatchableTaskMockRecorder) Cancel() *gomock.Call
- func (mr *MockBatchableTaskMockRecorder) Execute() *gomock.Call
- func (mr *MockBatchableTaskMockRecorder) HandleErr(err interface{}) *gomock.Call
- func (mr *MockBatchableTaskMockRecorder) IsRetryableError(err interface{}) *gomock.Call
- func (mr *MockBatchableTaskMockRecorder) MarkPoisonPill() *gomock.Call
- func (mr *MockBatchableTaskMockRecorder) MarkUnbatchable() *gomock.Call
- func (mr *MockBatchableTaskMockRecorder) Nack(err interface{}) *gomock.Call
- func (mr *MockBatchableTaskMockRecorder) QueueID() *gomock.Call
- func (mr *MockBatchableTaskMockRecorder) Reschedule() *gomock.Call
- func (mr *MockBatchableTaskMockRecorder) RetryPolicy() *gomock.Call
- func (mr *MockBatchableTaskMockRecorder) SourceClusterName() *gomock.Call
- func (mr *MockBatchableTaskMockRecorder) State() *gomock.Call
- func (mr *MockBatchableTaskMockRecorder) TaskCreationTime() *gomock.Call
- func (mr *MockBatchableTaskMockRecorder) TaskID() *gomock.Call
- type MockDLQHandler
- func (m *MockDLQHandler) EXPECT() *MockDLQHandlerMockRecorder
- func (m *MockDLQHandler) GetMessages(ctx context.Context, sourceCluster string, lastMessageID int64, pageSize int, ...) ([]*repication.ReplicationTask, []*repication.ReplicationTaskInfo, []byte, ...)
- func (m *MockDLQHandler) MergeMessages(ctx context.Context, sourceCluster string, lastMessageID int64, pageSize int, ...) ([]byte, error)
- func (m *MockDLQHandler) PurgeMessages(ctx context.Context, sourceCluster string, lastMessageID int64) error
- type MockDLQHandlerMockRecorder
- func (mr *MockDLQHandlerMockRecorder) GetMessages(ctx, sourceCluster, lastMessageID, pageSize, pageToken interface{}) *gomock.Call
- func (mr *MockDLQHandlerMockRecorder) MergeMessages(ctx, sourceCluster, lastMessageID, pageSize, pageToken interface{}) *gomock.Call
- func (mr *MockDLQHandlerMockRecorder) PurgeMessages(ctx, sourceCluster, lastMessageID interface{}) *gomock.Call
- type MockEagerNamespaceRefresher
- func (m *MockEagerNamespaceRefresher) EXPECT() *MockEagerNamespaceRefresherMockRecorder
- func (m *MockEagerNamespaceRefresher) SyncNamespaceFromSourceCluster(ctx context.Context, namespaceId namespace.ID, sourceCluster string) (*namespace.Namespace, error)
- func (m *MockEagerNamespaceRefresher) UpdateNamespaceFailoverVersion(namespaceId namespace.ID, targetFailoverVersion int64) error
- type MockEagerNamespaceRefresherMockRecorder
- type MockExecutableTask
- func (m *MockExecutableTask) Abort()
- func (m *MockExecutableTask) Ack()
- func (m *MockExecutableTask) Attempt() int
- func (m *MockExecutableTask) Cancel()
- func (m *MockExecutableTask) DeleteWorkflow(ctx context.Context, workflowKey definition.WorkflowKey) error
- func (m *MockExecutableTask) EXPECT() *MockExecutableTaskMockRecorder
- func (m *MockExecutableTask) GetNamespaceInfo(ctx context.Context, namespaceID string) (string, bool, error)
- func (m *MockExecutableTask) IsRetryableError(err error) bool
- func (m *MockExecutableTask) Nack(err error)
- func (m *MockExecutableTask) Reschedule()
- func (m *MockExecutableTask) Resend(ctx context.Context, remoteCluster string, ...) (bool, error)
- func (m *MockExecutableTask) RetryPolicy() backoff.RetryPolicy
- func (m *MockExecutableTask) SourceClusterName() string
- func (m *MockExecutableTask) State() tasks.State
- func (m *MockExecutableTask) TaskCreationTime() time.Time
- func (m *MockExecutableTask) TaskID() int64
- func (m *MockExecutableTask) TerminalState() bool
- type MockExecutableTaskMockRecorder
- func (mr *MockExecutableTaskMockRecorder) Abort() *gomock.Call
- func (mr *MockExecutableTaskMockRecorder) Ack() *gomock.Call
- func (mr *MockExecutableTaskMockRecorder) Attempt() *gomock.Call
- func (mr *MockExecutableTaskMockRecorder) Cancel() *gomock.Call
- func (mr *MockExecutableTaskMockRecorder) DeleteWorkflow(ctx, workflowKey interface{}) *gomock.Call
- func (mr *MockExecutableTaskMockRecorder) GetNamespaceInfo(ctx, namespaceID interface{}) *gomock.Call
- func (mr *MockExecutableTaskMockRecorder) IsRetryableError(err interface{}) *gomock.Call
- func (mr *MockExecutableTaskMockRecorder) Nack(err interface{}) *gomock.Call
- func (mr *MockExecutableTaskMockRecorder) Reschedule() *gomock.Call
- func (mr *MockExecutableTaskMockRecorder) Resend(ctx, remoteCluster, retryErr, remainingAttempt interface{}) *gomock.Call
- func (mr *MockExecutableTaskMockRecorder) RetryPolicy() *gomock.Call
- func (mr *MockExecutableTaskMockRecorder) SourceClusterName() *gomock.Call
- func (mr *MockExecutableTaskMockRecorder) State() *gomock.Call
- func (mr *MockExecutableTaskMockRecorder) TaskCreationTime() *gomock.Call
- func (mr *MockExecutableTaskMockRecorder) TaskID() *gomock.Call
- func (mr *MockExecutableTaskMockRecorder) TerminalState() *gomock.Call
- type MockExecutableTaskTracker
- func (m *MockExecutableTaskTracker) Cancel()
- func (m *MockExecutableTaskTracker) EXPECT() *MockExecutableTaskTrackerMockRecorder
- func (m *MockExecutableTaskTracker) LowWatermark() *WatermarkInfo
- func (m *MockExecutableTaskTracker) Size() int
- func (m *MockExecutableTaskTracker) TrackTasks(exclusiveHighWatermarkInfo WatermarkInfo, tasks ...TrackableExecutableTask) []TrackableExecutableTask
- type MockExecutableTaskTrackerMockRecorder
- func (mr *MockExecutableTaskTrackerMockRecorder) Cancel() *gomock.Call
- func (mr *MockExecutableTaskTrackerMockRecorder) LowWatermark() *gomock.Call
- func (mr *MockExecutableTaskTrackerMockRecorder) Size() *gomock.Call
- func (mr *MockExecutableTaskTrackerMockRecorder) TrackTasks(exclusiveHighWatermarkInfo interface{}, tasks ...interface{}) *gomock.Call
- type MockReceiverFlowController
- type MockReceiverFlowControllerMockRecorder
- type MockSenderFlowController
- type MockSenderFlowControllerMockRecorder
- type MockSourceTaskConverter
- type MockSourceTaskConverterMockRecorder
- type MockStreamReceiver
- type MockStreamReceiverMockRecorder
- type MockStreamSender
- type MockStreamSenderMockRecorder
- type MockTaskExecutor
- type MockTaskExecutorMockRecorder
- type MockTaskFetcherFactory
- type MockTaskFetcherFactoryMockRecorder
- type MockTaskProcessor
- type MockTaskProcessorMockRecorder
- type MockTrackableExecutableTask
- func (m *MockTrackableExecutableTask) Abort()
- func (m *MockTrackableExecutableTask) Ack()
- func (m *MockTrackableExecutableTask) Cancel()
- func (m *MockTrackableExecutableTask) EXPECT() *MockTrackableExecutableTaskMockRecorder
- func (m *MockTrackableExecutableTask) Execute() error
- func (m *MockTrackableExecutableTask) HandleErr(err error) error
- func (m *MockTrackableExecutableTask) IsRetryableError(err error) bool
- func (m *MockTrackableExecutableTask) MarkPoisonPill() error
- func (m *MockTrackableExecutableTask) Nack(err error)
- func (m *MockTrackableExecutableTask) QueueID() interface{}
- func (m *MockTrackableExecutableTask) Reschedule()
- func (m *MockTrackableExecutableTask) RetryPolicy() backoff.RetryPolicy
- func (m *MockTrackableExecutableTask) SourceClusterName() string
- func (m *MockTrackableExecutableTask) State() tasks.State
- func (m *MockTrackableExecutableTask) TaskCreationTime() time.Time
- func (m *MockTrackableExecutableTask) TaskID() int64
- type MockTrackableExecutableTaskMockRecorder
- func (mr *MockTrackableExecutableTaskMockRecorder) Abort() *gomock.Call
- func (mr *MockTrackableExecutableTaskMockRecorder) Ack() *gomock.Call
- func (mr *MockTrackableExecutableTaskMockRecorder) Cancel() *gomock.Call
- func (mr *MockTrackableExecutableTaskMockRecorder) Execute() *gomock.Call
- func (mr *MockTrackableExecutableTaskMockRecorder) HandleErr(err interface{}) *gomock.Call
- func (mr *MockTrackableExecutableTaskMockRecorder) IsRetryableError(err interface{}) *gomock.Call
- func (mr *MockTrackableExecutableTaskMockRecorder) MarkPoisonPill() *gomock.Call
- func (mr *MockTrackableExecutableTaskMockRecorder) Nack(err interface{}) *gomock.Call
- func (mr *MockTrackableExecutableTaskMockRecorder) QueueID() *gomock.Call
- func (mr *MockTrackableExecutableTaskMockRecorder) Reschedule() *gomock.Call
- func (mr *MockTrackableExecutableTaskMockRecorder) RetryPolicy() *gomock.Call
- func (mr *MockTrackableExecutableTaskMockRecorder) SourceClusterName() *gomock.Call
- func (mr *MockTrackableExecutableTaskMockRecorder) State() *gomock.Call
- func (mr *MockTrackableExecutableTaskMockRecorder) TaskCreationTime() *gomock.Call
- func (mr *MockTrackableExecutableTaskMockRecorder) TaskID() *gomock.Call
- type MocktaskFetcher
- type MocktaskFetcherMockRecorder
- type NoopDLQWriter
- type ProcessToolBox
- type ReceiverFlowController
- type ReceiverMode
- type SenderFlowController
- type SenderFlowControllerImpl
- type SequentialBatchableTaskQueue
- func (q *SequentialBatchableTaskQueue) Add(task TrackableExecutableTask)
- func (q *SequentialBatchableTaskQueue) ID() interface{}
- func (q *SequentialBatchableTaskQueue) IsEmpty() bool
- func (q *SequentialBatchableTaskQueue) Len() int
- func (q *SequentialBatchableTaskQueue) Peek() TrackableExecutableTask
- func (q *SequentialBatchableTaskQueue) Remove() (task TrackableExecutableTask)
- type SequentialTaskQueue
- func (q *SequentialTaskQueue) Add(task TrackableExecutableTask)
- func (q *SequentialTaskQueue) ID() interface{}
- func (q *SequentialTaskQueue) IsEmpty() bool
- func (q *SequentialTaskQueue) Len() int
- func (q *SequentialTaskQueue) Peek() TrackableExecutableTask
- func (q *SequentialTaskQueue) Remove() TrackableExecutableTask
- type SourceTaskConverter
- type SourceTaskConverterImpl
- type SourceTaskConverterProvider
- type Stream
- type StreamBiDirectionStreamClientProvider
- type StreamError
- type StreamReceiver
- type StreamReceiverImpl
- type StreamReceiverMonitor
- type StreamReceiverMonitorImpl
- type StreamResp
- type StreamSender
- type StreamSenderImpl
- type TaskExecutor
- type TaskExecutorParams
- type TaskExecutorProvider
- type TaskFetcherFactory
- type TaskParser
- type TaskProcessor
- type TrackableExecutableTask
- type WatermarkInfo
Constants ¶
const MarkPoisonPillMaxAttempts = 3
const (
ResendAttempt = 2
)
const TaskMaxSkipCount int = 1000
Variables ¶
var ( TaskRetryPolicy = backoff.NewExponentialRetryPolicy(1 * time.Second). WithBackoffCoefficient(1.2). WithMaximumInterval(5 * time.Second). WithMaximumAttempts(80). WithExpirationInterval(5 * time.Minute) ErrResendAttemptExceeded = serviceerror.NewInternal("resend history attempts exceeded") )
var ( // ErrClosed indicates stream closed before a read/write operation ErrClosed = serviceerror.NewUnavailable("stream closed") )
var ( // ErrUnknownReplicationTask is the error to indicate unknown replication task type ErrUnknownReplicationTask = serviceerror.NewInvalidArgument("unknown replication task") )
var Module = fx.Provide( NewTaskFetcherFactory, func(m persistence.ExecutionManager) ExecutionManager { return m }, NewExecutionManagerDLQWriter, replicationTaskConverterFactoryProvider, replicationTaskExecutorProvider, fx.Annotated{ Name: "HighPriorityTaskScheduler", Target: replicationStreamHighPrioritySchedulerProvider, }, fx.Annotated{ Name: "LowPriorityTaskScheduler", Target: replicationStreamLowPrioritySchedulerProvider, }, executableTaskConverterProvider, streamReceiverMonitorProvider, ndcHistoryResenderProvider, eagerNamespaceRefresherProvider, sequentialTaskQueueFactoryProvider, dlqWriterAdapterProvider, newDLQWriterToggle, historyPaginatedFetcherProvider, remoteEventHandlerProvider, localEventHandlerProvider, historyEventsHandlerProvider, )
Functions ¶
func ClusterIDToClusterNameShardCount ¶ added in v1.21.0
func IsRetryableError ¶ added in v1.24.0
func IsStreamError ¶ added in v1.24.0
func NewExecutableTaskConverter ¶ added in v1.23.0
func NewExecutableTaskConverter( processToolBox ProcessToolBox, ) *executableTaskConverterImpl
func NewExecutionManagerDLQWriter ¶ added in v1.23.0
func NewExecutionManagerDLQWriter(executionManager ExecutionManager) *executionManagerDLQWriter
NewExecutionManagerDLQWriter creates a new DLQWriter that uses the ExecutionManager.
func NewReceiverFlowControl ¶ added in v1.25.0
func NewReceiverFlowControl(signals map[enums.TaskPriority]FlowControlSignalProvider, config *configs.Config) *streamReceiverFlowControllerImpl
func NewSequentialBatchableTaskQueue ¶ added in v1.23.0
func NewSequentialBatchableTaskQueue( task TrackableExecutableTask, batchedIndividualTaskHandler func(task TrackableExecutableTask), logger log.Logger, metricsHandler metrics.Handler, ) ctasks.SequentialTaskQueue[TrackableExecutableTask]
func NewSequentialTaskQueue ¶ added in v1.21.0
func NewSequentialTaskQueue(task TrackableExecutableTask) ctasks.SequentialTaskQueue[TrackableExecutableTask]
func NewTaskProcessorManager ¶
func NewTaskProcessorManager( config *configs.Config, shard shard.Context, engine shard.Engine, workflowCache wcache.Cache, workflowDeleteManager deletemanager.DeleteManager, clientBean client.Bean, eventSerializer serialization.Serializer, replicationTaskFetcherFactory TaskFetcherFactory, taskExecutorProvider TaskExecutorProvider, dlqWriter DLQWriter, ) *taskProcessorManagerImpl
func SequentialTaskQueueCompareLess ¶ added in v1.21.0
func SequentialTaskQueueCompareLess(this TrackableExecutableTask, that TrackableExecutableTask) bool
func TaskOperationTag ¶ added in v1.21.0
func TaskOperationTag( replicationTask *replicationspb.ReplicationTask, ) string
func ValidateTasksHaveSamePriority ¶ added in v1.25.0
func ValidateTasksHaveSamePriority(messageBatchPriority enums.TaskPriority, tasks ...*replicationpb.ReplicationTask) error
func WorkflowKeyHashFn ¶ added in v1.23.0
func WorkflowKeyHashFn( item interface{}, ) uint32
func WrapEventLoop ¶ added in v1.24.0
func WrapEventLoop( originalEventLoop func() error, streamStopper func(), logger log.Logger, metricsHandler metrics.Handler, fromClusterKey ClusterShardKey, toClusterKey ClusterShardKey, retryInterval time.Duration, )
Types ¶
type AckManager ¶
type AckManager interface { NotifyNewTasks(tasks []tasks.Task) GetMaxTaskInfo() (int64, time.Time) GetTasks(ctx context.Context, pollingCluster string, queryMessageID int64) (*replicationspb.ReplicationMessages, error) GetTask(ctx context.Context, taskInfo *replicationspb.ReplicationTaskInfo) (*replicationspb.ReplicationTask, error) SubscribeNotification() (<-chan struct{}, string) UnsubscribeNotification(string) ConvertTask( ctx context.Context, task tasks.Task, ) (*replicationspb.ReplicationTask, error) GetReplicationTasksIter( ctx context.Context, pollingCluster string, minInclusiveTaskID int64, maxExclusiveTaskID int64, ) (collection.Iterator[tasks.Task], error) }
func NewAckManager ¶
func NewAckManager( shardContext shard.Context, workflowCache wcache.Cache, eventBlobCache persistence.XDCCache, executionMgr persistence.ExecutionManager, logger log.Logger, ) AckManager
type BatchableTask ¶ added in v1.23.0
type BatchableTask interface { TrackableExecutableTask // BatchWith task and return a new TrackableExecutableTask BatchWith(task BatchableTask) (TrackableExecutableTask, bool) CanBatch() bool // MarkUnbatchable will mark current task not batchable, so CanBatch() will return false MarkUnbatchable() }
type BiDirectionStream ¶ added in v1.21.0
type BiDirectionStreamClient ¶ added in v1.21.0
type BiDirectionStreamClientProvider ¶ added in v1.21.0
type BiDirectionStreamImpl ¶ added in v1.21.0
type BiDirectionStreamImpl[Req any, Resp any] struct { sync.Mutex // contains filtered or unexported fields }
func NewBiDirectionStream ¶ added in v1.21.0
func (*BiDirectionStreamImpl[Req, Resp]) Close ¶ added in v1.21.0
func (s *BiDirectionStreamImpl[Req, Resp]) Close()
func (*BiDirectionStreamImpl[Req, Resp]) IsValid ¶ added in v1.21.0
func (s *BiDirectionStreamImpl[Req, Resp]) IsValid() bool
type ClusterChannelKey ¶ added in v1.25.0
type ClusterChannelKey struct {
ClusterName string
}
type ClusterShardKey ¶ added in v1.21.0
func NewClusterShardKey ¶ added in v1.21.0
func NewClusterShardKey( ClusterID int32, ClusterShardID int32, ) ClusterShardKey
type ClusterShardKeyPair ¶ added in v1.21.0
type ClusterShardKeyPair struct { Client ClusterShardKey Server ClusterShardKey }
type DLQHandler ¶
type DLQHandler interface { GetMessages( ctx context.Context, sourceCluster string, lastMessageID int64, pageSize int, pageToken []byte, ) ([]*replicationspb.ReplicationTask, []*replicationspb.ReplicationTaskInfo, []byte, error) PurgeMessages( ctx context.Context, sourceCluster string, lastMessageID int64, ) error MergeMessages( ctx context.Context, sourceCluster string, lastMessageID int64, pageSize int, pageToken []byte, ) ([]byte, error) }
DLQHandler is the interface handles replication DLQ messages
func NewLazyDLQHandler ¶
func NewLazyDLQHandler( shard shard.Context, deleteManager deletemanager.DeleteManager, workflowCache wcache.Cache, clientBean client.Bean, taskExecutorProvider TaskExecutorProvider, ) DLQHandler
type DLQWriteRequest ¶ added in v1.23.0
type DLQWriteRequest struct { ShardID int32 SourceCluster string ReplicationTaskInfo *persistencespb.ReplicationTaskInfo }
DLQWriteRequest is a request to write a task to the DLQ.
type DLQWriter ¶ added in v1.23.0
type DLQWriter interface {
WriteTaskToDLQ(ctx context.Context, request DLQWriteRequest) error
}
DLQWriter is an interface that can be implemented easily by the two different queue solutions that we have. - Queue V1 implements this interface via persistence.ExecutionManager. - Queue V2 will implement this interface via go.temporal.io/server/service/history/queues.DLQWriter.
We want this interface to make the migration referenced by persistence.QueueV2 easier.
type DLQWriterAdapter ¶ added in v1.23.0
type DLQWriterAdapter struct {
// contains filtered or unexported fields
}
DLQWriterAdapter is a DLQWriter that uses the QueueV2 queues.DLQWriter object.
func NewDLQWriterAdapter ¶ added in v1.23.0
func NewDLQWriterAdapter( dlqWriter *queues.DLQWriter, taskParser TaskParser, currentClusterName string, ) *DLQWriterAdapter
NewDLQWriterAdapter creates a new DLQWriter from a QueueV2 queues.DLQWriter.
func (*DLQWriterAdapter) WriteTaskToDLQ ¶ added in v1.23.0
func (d *DLQWriterAdapter) WriteTaskToDLQ( ctx context.Context, request DLQWriteRequest, ) error
WriteTaskToDLQ implements [DLQWriter.WriteTaskToDLQ] by calling queues.DLQWriter.Write.
type EagerNamespaceRefresher ¶ added in v1.22.0
type EagerNamespaceRefresher interface { UpdateNamespaceFailoverVersion(namespaceId namespace.ID, targetFailoverVersion int64) error SyncNamespaceFromSourceCluster(ctx context.Context, namespaceId namespace.ID, sourceCluster string) (*namespace.Namespace, error) }
func NewEagerNamespaceRefresher ¶ added in v1.22.0
func NewEagerNamespaceRefresher( metadataManager persistence.MetadataManager, namespaceRegistry namespace.Registry, logger log.Logger, clientBean client.Bean, replicationTaskExecutor namespace.ReplicationTaskExecutor, currentCluster string, metricsHandler metrics.Handler) EagerNamespaceRefresher
type ExecutableActivityStateTask ¶ added in v1.21.0
type ExecutableActivityStateTask struct { ProcessToolBox definition.WorkflowKey ExecutableTask // contains filtered or unexported fields }
func NewExecutableActivityStateTask ¶ added in v1.21.0
func NewExecutableActivityStateTask( processToolBox ProcessToolBox, taskID int64, taskCreationTime time.Time, task *replicationspb.SyncActivityTaskAttributes, sourceClusterName string, priority enumsspb.TaskPriority, ) *ExecutableActivityStateTask
func (*ExecutableActivityStateTask) BatchWith ¶ added in v1.24.0
func (e *ExecutableActivityStateTask) BatchWith(incomingTask BatchableTask) (TrackableExecutableTask, bool)
func (*ExecutableActivityStateTask) CanBatch ¶ added in v1.24.0
func (e *ExecutableActivityStateTask) CanBatch() bool
func (*ExecutableActivityStateTask) Execute ¶ added in v1.21.0
func (e *ExecutableActivityStateTask) Execute() error
func (*ExecutableActivityStateTask) HandleErr ¶ added in v1.21.0
func (e *ExecutableActivityStateTask) HandleErr(err error) error
func (*ExecutableActivityStateTask) MarkPoisonPill ¶ added in v1.21.0
func (e *ExecutableActivityStateTask) MarkPoisonPill() error
func (*ExecutableActivityStateTask) MarkUnbatchable ¶ added in v1.24.0
func (e *ExecutableActivityStateTask) MarkUnbatchable()
func (*ExecutableActivityStateTask) QueueID ¶ added in v1.21.0
func (e *ExecutableActivityStateTask) QueueID() interface{}
type ExecutableHistoryTask ¶ added in v1.21.0
type ExecutableHistoryTask struct { ProcessToolBox definition.WorkflowKey ExecutableTask // contains filtered or unexported fields }
func NewExecutableHistoryTask ¶ added in v1.21.0
func NewExecutableHistoryTask( processToolBox ProcessToolBox, taskID int64, taskCreationTime time.Time, task *replicationspb.HistoryTaskAttributes, sourceClusterName string, priority enumsspb.TaskPriority, ) *ExecutableHistoryTask
func (*ExecutableHistoryTask) BatchWith ¶ added in v1.23.0
func (e *ExecutableHistoryTask) BatchWith(incomingTask BatchableTask) (TrackableExecutableTask, bool)
func (*ExecutableHistoryTask) CanBatch ¶ added in v1.23.0
func (e *ExecutableHistoryTask) CanBatch() bool
func (*ExecutableHistoryTask) Execute ¶ added in v1.21.0
func (e *ExecutableHistoryTask) Execute() error
func (*ExecutableHistoryTask) HandleErr ¶ added in v1.21.0
func (e *ExecutableHistoryTask) HandleErr(err error) error
func (*ExecutableHistoryTask) MarkPoisonPill ¶ added in v1.21.0
func (e *ExecutableHistoryTask) MarkPoisonPill() error
func (*ExecutableHistoryTask) MarkUnbatchable ¶ added in v1.23.0
func (e *ExecutableHistoryTask) MarkUnbatchable()
func (*ExecutableHistoryTask) QueueID ¶ added in v1.21.0
func (e *ExecutableHistoryTask) QueueID() interface{}
type ExecutableNoopTask ¶ added in v1.21.0
type ExecutableNoopTask struct {
ExecutableTask
}
func NewExecutableNoopTask ¶ added in v1.21.0
func NewExecutableNoopTask( processToolBox ProcessToolBox, taskID int64, taskCreationTime time.Time, sourceClusterName string, ) *ExecutableNoopTask
func (*ExecutableNoopTask) Execute ¶ added in v1.21.0
func (e *ExecutableNoopTask) Execute() error
func (*ExecutableNoopTask) HandleErr ¶ added in v1.21.0
func (e *ExecutableNoopTask) HandleErr(err error) error
func (*ExecutableNoopTask) MarkPoisonPill ¶ added in v1.21.0
func (e *ExecutableNoopTask) MarkPoisonPill() error
func (*ExecutableNoopTask) QueueID ¶ added in v1.21.0
func (e *ExecutableNoopTask) QueueID() interface{}
type ExecutableTask ¶ added in v1.21.0
type ExecutableTask interface { TaskID() int64 TaskCreationTime() time.Time SourceClusterName() string Ack() Nack(err error) Abort() Cancel() Reschedule() IsRetryableError(err error) bool RetryPolicy() backoff.RetryPolicy State() ctasks.State TerminalState() bool Attempt() int Resend( ctx context.Context, remoteCluster string, retryErr *serviceerrors.RetryReplication, remainingAttempt int, ) (bool, error) DeleteWorkflow( ctx context.Context, workflowKey definition.WorkflowKey, ) (retError error) GetNamespaceInfo( ctx context.Context, namespaceID string, ) (string, bool, error) }
type ExecutableTaskConverter ¶ added in v1.23.0
type ExecutableTaskConverter interface { Convert( taskClusterName string, clientShardKey ClusterShardKey, serverShardKey ClusterShardKey, replicationTasks ...*replicationspb.ReplicationTask, ) []TrackableExecutableTask }
type ExecutableTaskImpl ¶ added in v1.21.0
type ExecutableTaskImpl struct { ProcessToolBox // contains filtered or unexported fields }
func NewExecutableTask ¶ added in v1.21.0
func NewExecutableTask( processToolBox ProcessToolBox, taskID int64, metricsTag string, taskCreationTime time.Time, taskReceivedTime time.Time, sourceClusterName string, priority enumsspb.TaskPriority, ) *ExecutableTaskImpl
func (*ExecutableTaskImpl) Abort ¶ added in v1.21.0
func (e *ExecutableTaskImpl) Abort()
func (*ExecutableTaskImpl) Ack ¶ added in v1.21.0
func (e *ExecutableTaskImpl) Ack()
func (*ExecutableTaskImpl) Attempt ¶ added in v1.21.0
func (e *ExecutableTaskImpl) Attempt() int
func (*ExecutableTaskImpl) Cancel ¶ added in v1.21.0
func (e *ExecutableTaskImpl) Cancel()
func (*ExecutableTaskImpl) DeleteWorkflow ¶ added in v1.21.0
func (e *ExecutableTaskImpl) DeleteWorkflow( ctx context.Context, workflowKey definition.WorkflowKey, ) (retError error)
func (*ExecutableTaskImpl) GetNamespaceInfo ¶ added in v1.21.0
func (*ExecutableTaskImpl) IsRetryableError ¶ added in v1.21.0
func (e *ExecutableTaskImpl) IsRetryableError(err error) bool
func (*ExecutableTaskImpl) Nack ¶ added in v1.21.0
func (e *ExecutableTaskImpl) Nack(err error)
func (*ExecutableTaskImpl) Reschedule ¶ added in v1.21.0
func (e *ExecutableTaskImpl) Reschedule()
func (*ExecutableTaskImpl) Resend ¶ added in v1.21.0
func (e *ExecutableTaskImpl) Resend( ctx context.Context, remoteCluster string, retryErr *serviceerrors.RetryReplication, remainingAttempt int, ) (bool, error)
func (*ExecutableTaskImpl) RetryPolicy ¶ added in v1.21.0
func (e *ExecutableTaskImpl) RetryPolicy() backoff.RetryPolicy
func (*ExecutableTaskImpl) SourceClusterName ¶ added in v1.23.0
func (e *ExecutableTaskImpl) SourceClusterName() string
func (*ExecutableTaskImpl) State ¶ added in v1.21.0
func (e *ExecutableTaskImpl) State() ctasks.State
func (*ExecutableTaskImpl) TaskCreationTime ¶ added in v1.21.0
func (e *ExecutableTaskImpl) TaskCreationTime() time.Time
func (*ExecutableTaskImpl) TaskID ¶ added in v1.21.0
func (e *ExecutableTaskImpl) TaskID() int64
func (*ExecutableTaskImpl) TerminalState ¶ added in v1.21.0
func (e *ExecutableTaskImpl) TerminalState() bool
type ExecutableTaskTracker ¶ added in v1.21.0
type ExecutableTaskTracker interface { TrackTasks(exclusiveHighWatermarkInfo WatermarkInfo, tasks ...TrackableExecutableTask) []TrackableExecutableTask LowWatermark() *WatermarkInfo Size() int Cancel() }
type ExecutableTaskTrackerImpl ¶ added in v1.21.0
func NewExecutableTaskTracker ¶ added in v1.21.0
func NewExecutableTaskTracker( logger log.Logger, metricsHandler metrics.Handler, ) *ExecutableTaskTrackerImpl
func (*ExecutableTaskTrackerImpl) Cancel ¶ added in v1.21.0
func (t *ExecutableTaskTrackerImpl) Cancel()
func (*ExecutableTaskTrackerImpl) LowWatermark ¶ added in v1.21.0
func (t *ExecutableTaskTrackerImpl) LowWatermark() *WatermarkInfo
func (*ExecutableTaskTrackerImpl) Size ¶ added in v1.21.0
func (t *ExecutableTaskTrackerImpl) Size() int
func (*ExecutableTaskTrackerImpl) TrackTasks ¶ added in v1.21.0
func (t *ExecutableTaskTrackerImpl) TrackTasks( exclusiveHighWatermarkInfo WatermarkInfo, tasks ...TrackableExecutableTask, ) []TrackableExecutableTask
TrackTasks add tasks for tracking, return valid tasks (dedup) if task tracker is cancelled, then newly added tasks will also be cancelled tasks should be sorted by task ID, all task IDs < exclusiveHighWatermarkInfo
type ExecutableUnknownTask ¶ added in v1.21.0
type ExecutableUnknownTask struct { ProcessToolBox ExecutableTask // contains filtered or unexported fields }
func NewExecutableUnknownTask ¶ added in v1.21.0
func NewExecutableUnknownTask( processToolBox ProcessToolBox, taskID int64, taskCreationTime time.Time, task any, ) *ExecutableUnknownTask
func (*ExecutableUnknownTask) Execute ¶ added in v1.21.0
func (e *ExecutableUnknownTask) Execute() error
func (*ExecutableUnknownTask) HandleErr ¶ added in v1.21.0
func (e *ExecutableUnknownTask) HandleErr(err error) error
func (*ExecutableUnknownTask) IsRetryableError ¶ added in v1.21.0
func (e *ExecutableUnknownTask) IsRetryableError(err error) bool
func (*ExecutableUnknownTask) MarkPoisonPill ¶ added in v1.21.0
func (e *ExecutableUnknownTask) MarkPoisonPill() error
func (*ExecutableUnknownTask) QueueID ¶ added in v1.21.0
func (e *ExecutableUnknownTask) QueueID() interface{}
type ExecutableWorkflowStateTask ¶ added in v1.21.0
type ExecutableWorkflowStateTask struct { ProcessToolBox definition.WorkflowKey ExecutableTask // contains filtered or unexported fields }
func NewExecutableWorkflowStateTask ¶ added in v1.21.0
func NewExecutableWorkflowStateTask( processToolBox ProcessToolBox, taskID int64, taskCreationTime time.Time, task *replicationspb.SyncWorkflowStateTaskAttributes, sourceClusterName string, priority enumsspb.TaskPriority, ) *ExecutableWorkflowStateTask
func (*ExecutableWorkflowStateTask) Execute ¶ added in v1.21.0
func (e *ExecutableWorkflowStateTask) Execute() error
func (*ExecutableWorkflowStateTask) HandleErr ¶ added in v1.21.0
func (e *ExecutableWorkflowStateTask) HandleErr(err error) error
func (*ExecutableWorkflowStateTask) MarkPoisonPill ¶ added in v1.21.0
func (e *ExecutableWorkflowStateTask) MarkPoisonPill() error
func (*ExecutableWorkflowStateTask) QueueID ¶ added in v1.21.0
func (e *ExecutableWorkflowStateTask) QueueID() interface{}
type ExecutionManager ¶ added in v1.23.0
type ExecutionManager interface { PutReplicationTaskToDLQ( ctx context.Context, request *persistence.PutReplicationTaskToDLQRequest, ) error }
ExecutionManager is a trimmed version of go.temporal.io/server/common/persistence.ExecutionManager that only provides the methods we need.
type FlowControlSignal ¶ added in v1.25.0
type FlowControlSignal struct {
// contains filtered or unexported fields
}
FlowControlSignal holds signals to make flow control decision, more signalsProvider can be added here i.e. total persistence rps, cpu usage etc.
type FlowControlSignalProvider ¶ added in v1.25.0
type FlowControlSignalProvider func() *FlowControlSignal
type MockAckManager ¶
type MockAckManager struct {
// contains filtered or unexported fields
}
MockAckManager is a mock of AckManager interface.
func NewMockAckManager ¶
func NewMockAckManager(ctrl *gomock.Controller) *MockAckManager
NewMockAckManager creates a new mock instance.
func (*MockAckManager) ConvertTask ¶ added in v1.21.0
func (m *MockAckManager) ConvertTask(ctx context.Context, task tasks.Task) (*repication.ReplicationTask, error)
ConvertTask mocks base method.
func (*MockAckManager) EXPECT ¶
func (m *MockAckManager) EXPECT() *MockAckManagerMockRecorder
EXPECT returns an object that allows the caller to indicate expected use.
func (*MockAckManager) GetMaxTaskInfo ¶ added in v1.17.2
func (m *MockAckManager) GetMaxTaskInfo() (int64, time.Time)
GetMaxTaskInfo mocks base method.
func (*MockAckManager) GetReplicationTasksIter ¶ added in v1.21.0
func (m *MockAckManager) GetReplicationTasksIter(ctx context.Context, pollingCluster string, minInclusiveTaskID, maxExclusiveTaskID int64) (collection.Iterator[tasks.Task], error)
GetReplicationTasksIter mocks base method.
func (*MockAckManager) GetTask ¶
func (m *MockAckManager) GetTask(ctx context.Context, taskInfo *repication.ReplicationTaskInfo) (*repication.ReplicationTask, error)
GetTask mocks base method.
func (*MockAckManager) GetTasks ¶
func (m *MockAckManager) GetTasks(ctx context.Context, pollingCluster string, queryMessageID int64) (*repication.ReplicationMessages, error)
GetTasks mocks base method.
func (*MockAckManager) NotifyNewTasks ¶
func (m *MockAckManager) NotifyNewTasks(tasks []tasks.Task)
NotifyNewTasks mocks base method.
func (*MockAckManager) SubscribeNotification ¶ added in v1.21.0
func (m *MockAckManager) SubscribeNotification() (<-chan struct{}, string)
SubscribeNotification mocks base method.
func (*MockAckManager) UnsubscribeNotification ¶ added in v1.21.0
func (m *MockAckManager) UnsubscribeNotification(arg0 string)
UnsubscribeNotification mocks base method.
type MockAckManagerMockRecorder ¶
type MockAckManagerMockRecorder struct {
// contains filtered or unexported fields
}
MockAckManagerMockRecorder is the mock recorder for MockAckManager.
func (*MockAckManagerMockRecorder) ConvertTask ¶ added in v1.21.0
func (mr *MockAckManagerMockRecorder) ConvertTask(ctx, task interface{}) *gomock.Call
ConvertTask indicates an expected call of ConvertTask.
func (*MockAckManagerMockRecorder) GetMaxTaskInfo ¶ added in v1.17.2
func (mr *MockAckManagerMockRecorder) GetMaxTaskInfo() *gomock.Call
GetMaxTaskInfo indicates an expected call of GetMaxTaskInfo.
func (*MockAckManagerMockRecorder) GetReplicationTasksIter ¶ added in v1.21.0
func (mr *MockAckManagerMockRecorder) GetReplicationTasksIter(ctx, pollingCluster, minInclusiveTaskID, maxExclusiveTaskID interface{}) *gomock.Call
GetReplicationTasksIter indicates an expected call of GetReplicationTasksIter.
func (*MockAckManagerMockRecorder) GetTask ¶
func (mr *MockAckManagerMockRecorder) GetTask(ctx, taskInfo interface{}) *gomock.Call
GetTask indicates an expected call of GetTask.
func (*MockAckManagerMockRecorder) GetTasks ¶
func (mr *MockAckManagerMockRecorder) GetTasks(ctx, pollingCluster, queryMessageID interface{}) *gomock.Call
GetTasks indicates an expected call of GetTasks.
func (*MockAckManagerMockRecorder) NotifyNewTasks ¶
func (mr *MockAckManagerMockRecorder) NotifyNewTasks(tasks interface{}) *gomock.Call
NotifyNewTasks indicates an expected call of NotifyNewTasks.
func (*MockAckManagerMockRecorder) SubscribeNotification ¶ added in v1.21.0
func (mr *MockAckManagerMockRecorder) SubscribeNotification() *gomock.Call
SubscribeNotification indicates an expected call of SubscribeNotification.
func (*MockAckManagerMockRecorder) UnsubscribeNotification ¶ added in v1.21.0
func (mr *MockAckManagerMockRecorder) UnsubscribeNotification(arg0 interface{}) *gomock.Call
UnsubscribeNotification indicates an expected call of UnsubscribeNotification.
type MockBatchableTask ¶ added in v1.23.0
type MockBatchableTask struct {
// contains filtered or unexported fields
}
MockBatchableTask is a mock of BatchableTask interface.
func NewMockBatchableTask ¶ added in v1.23.0
func NewMockBatchableTask(ctrl *gomock.Controller) *MockBatchableTask
NewMockBatchableTask creates a new mock instance.
func (*MockBatchableTask) Abort ¶ added in v1.23.0
func (m *MockBatchableTask) Abort()
Abort mocks base method.
func (*MockBatchableTask) Ack ¶ added in v1.23.0
func (m *MockBatchableTask) Ack()
Ack mocks base method.
func (*MockBatchableTask) BatchWith ¶ added in v1.23.0
func (m *MockBatchableTask) BatchWith(task BatchableTask) (TrackableExecutableTask, bool)
BatchWith mocks base method.
func (*MockBatchableTask) CanBatch ¶ added in v1.23.0
func (m *MockBatchableTask) CanBatch() bool
CanBatch mocks base method.
func (*MockBatchableTask) Cancel ¶ added in v1.23.0
func (m *MockBatchableTask) Cancel()
Cancel mocks base method.
func (*MockBatchableTask) EXPECT ¶ added in v1.23.0
func (m *MockBatchableTask) EXPECT() *MockBatchableTaskMockRecorder
EXPECT returns an object that allows the caller to indicate expected use.
func (*MockBatchableTask) Execute ¶ added in v1.23.0
func (m *MockBatchableTask) Execute() error
Execute mocks base method.
func (*MockBatchableTask) HandleErr ¶ added in v1.23.0
func (m *MockBatchableTask) HandleErr(err error) error
HandleErr mocks base method.
func (*MockBatchableTask) IsRetryableError ¶ added in v1.23.0
func (m *MockBatchableTask) IsRetryableError(err error) bool
IsRetryableError mocks base method.
func (*MockBatchableTask) MarkPoisonPill ¶ added in v1.23.0
func (m *MockBatchableTask) MarkPoisonPill() error
MarkPoisonPill mocks base method.
func (*MockBatchableTask) MarkUnbatchable ¶ added in v1.23.0
func (m *MockBatchableTask) MarkUnbatchable()
MarkUnbatchable mocks base method.
func (*MockBatchableTask) Nack ¶ added in v1.23.0
func (m *MockBatchableTask) Nack(err error)
Nack mocks base method.
func (*MockBatchableTask) QueueID ¶ added in v1.23.0
func (m *MockBatchableTask) QueueID() interface{}
QueueID mocks base method.
func (*MockBatchableTask) Reschedule ¶ added in v1.23.0
func (m *MockBatchableTask) Reschedule()
Reschedule mocks base method.
func (*MockBatchableTask) RetryPolicy ¶ added in v1.23.0
func (m *MockBatchableTask) RetryPolicy() backoff.RetryPolicy
RetryPolicy mocks base method.
func (*MockBatchableTask) SourceClusterName ¶ added in v1.25.0
func (m *MockBatchableTask) SourceClusterName() string
SourceClusterName mocks base method.
func (*MockBatchableTask) State ¶ added in v1.23.0
func (m *MockBatchableTask) State() tasks.State
State mocks base method.
func (*MockBatchableTask) TaskCreationTime ¶ added in v1.23.0
func (m *MockBatchableTask) TaskCreationTime() time.Time
TaskCreationTime mocks base method.
func (*MockBatchableTask) TaskID ¶ added in v1.23.0
func (m *MockBatchableTask) TaskID() int64
TaskID mocks base method.
type MockBatchableTaskMockRecorder ¶ added in v1.23.0
type MockBatchableTaskMockRecorder struct {
// contains filtered or unexported fields
}
MockBatchableTaskMockRecorder is the mock recorder for MockBatchableTask.
func (*MockBatchableTaskMockRecorder) Abort ¶ added in v1.23.0
func (mr *MockBatchableTaskMockRecorder) Abort() *gomock.Call
Abort indicates an expected call of Abort.
func (*MockBatchableTaskMockRecorder) Ack ¶ added in v1.23.0
func (mr *MockBatchableTaskMockRecorder) Ack() *gomock.Call
Ack indicates an expected call of Ack.
func (*MockBatchableTaskMockRecorder) BatchWith ¶ added in v1.23.0
func (mr *MockBatchableTaskMockRecorder) BatchWith(task interface{}) *gomock.Call
BatchWith indicates an expected call of BatchWith.
func (*MockBatchableTaskMockRecorder) CanBatch ¶ added in v1.23.0
func (mr *MockBatchableTaskMockRecorder) CanBatch() *gomock.Call
CanBatch indicates an expected call of CanBatch.
func (*MockBatchableTaskMockRecorder) Cancel ¶ added in v1.23.0
func (mr *MockBatchableTaskMockRecorder) Cancel() *gomock.Call
Cancel indicates an expected call of Cancel.
func (*MockBatchableTaskMockRecorder) Execute ¶ added in v1.23.0
func (mr *MockBatchableTaskMockRecorder) Execute() *gomock.Call
Execute indicates an expected call of Execute.
func (*MockBatchableTaskMockRecorder) HandleErr ¶ added in v1.23.0
func (mr *MockBatchableTaskMockRecorder) HandleErr(err interface{}) *gomock.Call
HandleErr indicates an expected call of HandleErr.
func (*MockBatchableTaskMockRecorder) IsRetryableError ¶ added in v1.23.0
func (mr *MockBatchableTaskMockRecorder) IsRetryableError(err interface{}) *gomock.Call
IsRetryableError indicates an expected call of IsRetryableError.
func (*MockBatchableTaskMockRecorder) MarkPoisonPill ¶ added in v1.23.0
func (mr *MockBatchableTaskMockRecorder) MarkPoisonPill() *gomock.Call
MarkPoisonPill indicates an expected call of MarkPoisonPill.
func (*MockBatchableTaskMockRecorder) MarkUnbatchable ¶ added in v1.23.0
func (mr *MockBatchableTaskMockRecorder) MarkUnbatchable() *gomock.Call
MarkUnbatchable indicates an expected call of MarkUnbatchable.
func (*MockBatchableTaskMockRecorder) Nack ¶ added in v1.23.0
func (mr *MockBatchableTaskMockRecorder) Nack(err interface{}) *gomock.Call
Nack indicates an expected call of Nack.
func (*MockBatchableTaskMockRecorder) QueueID ¶ added in v1.23.0
func (mr *MockBatchableTaskMockRecorder) QueueID() *gomock.Call
QueueID indicates an expected call of QueueID.
func (*MockBatchableTaskMockRecorder) Reschedule ¶ added in v1.23.0
func (mr *MockBatchableTaskMockRecorder) Reschedule() *gomock.Call
Reschedule indicates an expected call of Reschedule.
func (*MockBatchableTaskMockRecorder) RetryPolicy ¶ added in v1.23.0
func (mr *MockBatchableTaskMockRecorder) RetryPolicy() *gomock.Call
RetryPolicy indicates an expected call of RetryPolicy.
func (*MockBatchableTaskMockRecorder) SourceClusterName ¶ added in v1.25.0
func (mr *MockBatchableTaskMockRecorder) SourceClusterName() *gomock.Call
SourceClusterName indicates an expected call of SourceClusterName.
func (*MockBatchableTaskMockRecorder) State ¶ added in v1.23.0
func (mr *MockBatchableTaskMockRecorder) State() *gomock.Call
State indicates an expected call of State.
func (*MockBatchableTaskMockRecorder) TaskCreationTime ¶ added in v1.23.0
func (mr *MockBatchableTaskMockRecorder) TaskCreationTime() *gomock.Call
TaskCreationTime indicates an expected call of TaskCreationTime.
func (*MockBatchableTaskMockRecorder) TaskID ¶ added in v1.23.0
func (mr *MockBatchableTaskMockRecorder) TaskID() *gomock.Call
TaskID indicates an expected call of TaskID.
type MockDLQHandler ¶
type MockDLQHandler struct {
// contains filtered or unexported fields
}
MockDLQHandler is a mock of DLQHandler interface.
func NewMockDLQHandler ¶
func NewMockDLQHandler(ctrl *gomock.Controller) *MockDLQHandler
NewMockDLQHandler creates a new mock instance.
func (*MockDLQHandler) EXPECT ¶
func (m *MockDLQHandler) EXPECT() *MockDLQHandlerMockRecorder
EXPECT returns an object that allows the caller to indicate expected use.
func (*MockDLQHandler) GetMessages ¶
func (m *MockDLQHandler) GetMessages(ctx context.Context, sourceCluster string, lastMessageID int64, pageSize int, pageToken []byte) ([]*repication.ReplicationTask, []*repication.ReplicationTaskInfo, []byte, error)
GetMessages mocks base method.
func (*MockDLQHandler) MergeMessages ¶
func (m *MockDLQHandler) MergeMessages(ctx context.Context, sourceCluster string, lastMessageID int64, pageSize int, pageToken []byte) ([]byte, error)
MergeMessages mocks base method.
func (*MockDLQHandler) PurgeMessages ¶
func (m *MockDLQHandler) PurgeMessages(ctx context.Context, sourceCluster string, lastMessageID int64) error
PurgeMessages mocks base method.
type MockDLQHandlerMockRecorder ¶
type MockDLQHandlerMockRecorder struct {
// contains filtered or unexported fields
}
MockDLQHandlerMockRecorder is the mock recorder for MockDLQHandler.
func (*MockDLQHandlerMockRecorder) GetMessages ¶
func (mr *MockDLQHandlerMockRecorder) GetMessages(ctx, sourceCluster, lastMessageID, pageSize, pageToken interface{}) *gomock.Call
GetMessages indicates an expected call of GetMessages.
func (*MockDLQHandlerMockRecorder) MergeMessages ¶
func (mr *MockDLQHandlerMockRecorder) MergeMessages(ctx, sourceCluster, lastMessageID, pageSize, pageToken interface{}) *gomock.Call
MergeMessages indicates an expected call of MergeMessages.
func (*MockDLQHandlerMockRecorder) PurgeMessages ¶
func (mr *MockDLQHandlerMockRecorder) PurgeMessages(ctx, sourceCluster, lastMessageID interface{}) *gomock.Call
PurgeMessages indicates an expected call of PurgeMessages.
type MockEagerNamespaceRefresher ¶ added in v1.23.0
type MockEagerNamespaceRefresher struct {
// contains filtered or unexported fields
}
MockEagerNamespaceRefresher is a mock of EagerNamespaceRefresher interface.
func NewMockEagerNamespaceRefresher ¶ added in v1.23.0
func NewMockEagerNamespaceRefresher(ctrl *gomock.Controller) *MockEagerNamespaceRefresher
NewMockEagerNamespaceRefresher creates a new mock instance.
func (*MockEagerNamespaceRefresher) EXPECT ¶ added in v1.23.0
func (m *MockEagerNamespaceRefresher) EXPECT() *MockEagerNamespaceRefresherMockRecorder
EXPECT returns an object that allows the caller to indicate expected use.
func (*MockEagerNamespaceRefresher) SyncNamespaceFromSourceCluster ¶ added in v1.23.0
func (m *MockEagerNamespaceRefresher) SyncNamespaceFromSourceCluster(ctx context.Context, namespaceId namespace.ID, sourceCluster string) (*namespace.Namespace, error)
SyncNamespaceFromSourceCluster mocks base method.
func (*MockEagerNamespaceRefresher) UpdateNamespaceFailoverVersion ¶ added in v1.23.0
func (m *MockEagerNamespaceRefresher) UpdateNamespaceFailoverVersion(namespaceId namespace.ID, targetFailoverVersion int64) error
UpdateNamespaceFailoverVersion mocks base method.
type MockEagerNamespaceRefresherMockRecorder ¶ added in v1.23.0
type MockEagerNamespaceRefresherMockRecorder struct {
// contains filtered or unexported fields
}
MockEagerNamespaceRefresherMockRecorder is the mock recorder for MockEagerNamespaceRefresher.
func (*MockEagerNamespaceRefresherMockRecorder) SyncNamespaceFromSourceCluster ¶ added in v1.23.0
func (mr *MockEagerNamespaceRefresherMockRecorder) SyncNamespaceFromSourceCluster(ctx, namespaceId, sourceCluster interface{}) *gomock.Call
SyncNamespaceFromSourceCluster indicates an expected call of SyncNamespaceFromSourceCluster.
func (*MockEagerNamespaceRefresherMockRecorder) UpdateNamespaceFailoverVersion ¶ added in v1.23.0
func (mr *MockEagerNamespaceRefresherMockRecorder) UpdateNamespaceFailoverVersion(namespaceId, targetFailoverVersion interface{}) *gomock.Call
UpdateNamespaceFailoverVersion indicates an expected call of UpdateNamespaceFailoverVersion.
type MockExecutableTask ¶ added in v1.21.0
type MockExecutableTask struct {
// contains filtered or unexported fields
}
MockExecutableTask is a mock of ExecutableTask interface.
func NewMockExecutableTask ¶ added in v1.21.0
func NewMockExecutableTask(ctrl *gomock.Controller) *MockExecutableTask
NewMockExecutableTask creates a new mock instance.
func (*MockExecutableTask) Abort ¶ added in v1.21.0
func (m *MockExecutableTask) Abort()
Abort mocks base method.
func (*MockExecutableTask) Ack ¶ added in v1.21.0
func (m *MockExecutableTask) Ack()
Ack mocks base method.
func (*MockExecutableTask) Attempt ¶ added in v1.21.0
func (m *MockExecutableTask) Attempt() int
Attempt mocks base method.
func (*MockExecutableTask) Cancel ¶ added in v1.21.0
func (m *MockExecutableTask) Cancel()
Cancel mocks base method.
func (*MockExecutableTask) DeleteWorkflow ¶ added in v1.21.0
func (m *MockExecutableTask) DeleteWorkflow(ctx context.Context, workflowKey definition.WorkflowKey) error
DeleteWorkflow mocks base method.
func (*MockExecutableTask) EXPECT ¶ added in v1.21.0
func (m *MockExecutableTask) EXPECT() *MockExecutableTaskMockRecorder
EXPECT returns an object that allows the caller to indicate expected use.
func (*MockExecutableTask) GetNamespaceInfo ¶ added in v1.21.0
func (m *MockExecutableTask) GetNamespaceInfo(ctx context.Context, namespaceID string) (string, bool, error)
GetNamespaceInfo mocks base method.
func (*MockExecutableTask) IsRetryableError ¶ added in v1.21.0
func (m *MockExecutableTask) IsRetryableError(err error) bool
IsRetryableError mocks base method.
func (*MockExecutableTask) Nack ¶ added in v1.21.0
func (m *MockExecutableTask) Nack(err error)
Nack mocks base method.
func (*MockExecutableTask) Reschedule ¶ added in v1.21.0
func (m *MockExecutableTask) Reschedule()
Reschedule mocks base method.
func (*MockExecutableTask) Resend ¶ added in v1.21.0
func (m *MockExecutableTask) Resend(ctx context.Context, remoteCluster string, retryErr *serviceerror.RetryReplication, remainingAttempt int) (bool, error)
Resend mocks base method.
func (*MockExecutableTask) RetryPolicy ¶ added in v1.21.0
func (m *MockExecutableTask) RetryPolicy() backoff.RetryPolicy
RetryPolicy mocks base method.
func (*MockExecutableTask) SourceClusterName ¶ added in v1.23.0
func (m *MockExecutableTask) SourceClusterName() string
SourceClusterName mocks base method.
func (*MockExecutableTask) State ¶ added in v1.21.0
func (m *MockExecutableTask) State() tasks.State
State mocks base method.
func (*MockExecutableTask) TaskCreationTime ¶ added in v1.21.0
func (m *MockExecutableTask) TaskCreationTime() time.Time
TaskCreationTime mocks base method.
func (*MockExecutableTask) TaskID ¶ added in v1.21.0
func (m *MockExecutableTask) TaskID() int64
TaskID mocks base method.
func (*MockExecutableTask) TerminalState ¶ added in v1.21.0
func (m *MockExecutableTask) TerminalState() bool
TerminalState mocks base method.
type MockExecutableTaskMockRecorder ¶ added in v1.21.0
type MockExecutableTaskMockRecorder struct {
// contains filtered or unexported fields
}
MockExecutableTaskMockRecorder is the mock recorder for MockExecutableTask.
func (*MockExecutableTaskMockRecorder) Abort ¶ added in v1.21.0
func (mr *MockExecutableTaskMockRecorder) Abort() *gomock.Call
Abort indicates an expected call of Abort.
func (*MockExecutableTaskMockRecorder) Ack ¶ added in v1.21.0
func (mr *MockExecutableTaskMockRecorder) Ack() *gomock.Call
Ack indicates an expected call of Ack.
func (*MockExecutableTaskMockRecorder) Attempt ¶ added in v1.21.0
func (mr *MockExecutableTaskMockRecorder) Attempt() *gomock.Call
Attempt indicates an expected call of Attempt.
func (*MockExecutableTaskMockRecorder) Cancel ¶ added in v1.21.0
func (mr *MockExecutableTaskMockRecorder) Cancel() *gomock.Call
Cancel indicates an expected call of Cancel.
func (*MockExecutableTaskMockRecorder) DeleteWorkflow ¶ added in v1.21.0
func (mr *MockExecutableTaskMockRecorder) DeleteWorkflow(ctx, workflowKey interface{}) *gomock.Call
DeleteWorkflow indicates an expected call of DeleteWorkflow.
func (*MockExecutableTaskMockRecorder) GetNamespaceInfo ¶ added in v1.21.0
func (mr *MockExecutableTaskMockRecorder) GetNamespaceInfo(ctx, namespaceID interface{}) *gomock.Call
GetNamespaceInfo indicates an expected call of GetNamespaceInfo.
func (*MockExecutableTaskMockRecorder) IsRetryableError ¶ added in v1.21.0
func (mr *MockExecutableTaskMockRecorder) IsRetryableError(err interface{}) *gomock.Call
IsRetryableError indicates an expected call of IsRetryableError.
func (*MockExecutableTaskMockRecorder) Nack ¶ added in v1.21.0
func (mr *MockExecutableTaskMockRecorder) Nack(err interface{}) *gomock.Call
Nack indicates an expected call of Nack.
func (*MockExecutableTaskMockRecorder) Reschedule ¶ added in v1.21.0
func (mr *MockExecutableTaskMockRecorder) Reschedule() *gomock.Call
Reschedule indicates an expected call of Reschedule.
func (*MockExecutableTaskMockRecorder) Resend ¶ added in v1.21.0
func (mr *MockExecutableTaskMockRecorder) Resend(ctx, remoteCluster, retryErr, remainingAttempt interface{}) *gomock.Call
Resend indicates an expected call of Resend.
func (*MockExecutableTaskMockRecorder) RetryPolicy ¶ added in v1.21.0
func (mr *MockExecutableTaskMockRecorder) RetryPolicy() *gomock.Call
RetryPolicy indicates an expected call of RetryPolicy.
func (*MockExecutableTaskMockRecorder) SourceClusterName ¶ added in v1.23.0
func (mr *MockExecutableTaskMockRecorder) SourceClusterName() *gomock.Call
SourceClusterName indicates an expected call of SourceClusterName.
func (*MockExecutableTaskMockRecorder) State ¶ added in v1.21.0
func (mr *MockExecutableTaskMockRecorder) State() *gomock.Call
State indicates an expected call of State.
func (*MockExecutableTaskMockRecorder) TaskCreationTime ¶ added in v1.21.0
func (mr *MockExecutableTaskMockRecorder) TaskCreationTime() *gomock.Call
TaskCreationTime indicates an expected call of TaskCreationTime.
func (*MockExecutableTaskMockRecorder) TaskID ¶ added in v1.21.0
func (mr *MockExecutableTaskMockRecorder) TaskID() *gomock.Call
TaskID indicates an expected call of TaskID.
func (*MockExecutableTaskMockRecorder) TerminalState ¶ added in v1.21.0
func (mr *MockExecutableTaskMockRecorder) TerminalState() *gomock.Call
TerminalState indicates an expected call of TerminalState.
type MockExecutableTaskTracker ¶ added in v1.21.0
type MockExecutableTaskTracker struct {
// contains filtered or unexported fields
}
MockExecutableTaskTracker is a mock of ExecutableTaskTracker interface.
func NewMockExecutableTaskTracker ¶ added in v1.21.0
func NewMockExecutableTaskTracker(ctrl *gomock.Controller) *MockExecutableTaskTracker
NewMockExecutableTaskTracker creates a new mock instance.
func (*MockExecutableTaskTracker) Cancel ¶ added in v1.21.0
func (m *MockExecutableTaskTracker) Cancel()
Cancel mocks base method.
func (*MockExecutableTaskTracker) EXPECT ¶ added in v1.21.0
func (m *MockExecutableTaskTracker) EXPECT() *MockExecutableTaskTrackerMockRecorder
EXPECT returns an object that allows the caller to indicate expected use.
func (*MockExecutableTaskTracker) LowWatermark ¶ added in v1.21.0
func (m *MockExecutableTaskTracker) LowWatermark() *WatermarkInfo
LowWatermark mocks base method.
func (*MockExecutableTaskTracker) Size ¶ added in v1.21.0
func (m *MockExecutableTaskTracker) Size() int
Size mocks base method.
func (*MockExecutableTaskTracker) TrackTasks ¶ added in v1.21.0
func (m *MockExecutableTaskTracker) TrackTasks(exclusiveHighWatermarkInfo WatermarkInfo, tasks ...TrackableExecutableTask) []TrackableExecutableTask
TrackTasks mocks base method.
type MockExecutableTaskTrackerMockRecorder ¶ added in v1.21.0
type MockExecutableTaskTrackerMockRecorder struct {
// contains filtered or unexported fields
}
MockExecutableTaskTrackerMockRecorder is the mock recorder for MockExecutableTaskTracker.
func (*MockExecutableTaskTrackerMockRecorder) Cancel ¶ added in v1.21.0
func (mr *MockExecutableTaskTrackerMockRecorder) Cancel() *gomock.Call
Cancel indicates an expected call of Cancel.
func (*MockExecutableTaskTrackerMockRecorder) LowWatermark ¶ added in v1.21.0
func (mr *MockExecutableTaskTrackerMockRecorder) LowWatermark() *gomock.Call
LowWatermark indicates an expected call of LowWatermark.
func (*MockExecutableTaskTrackerMockRecorder) Size ¶ added in v1.21.0
func (mr *MockExecutableTaskTrackerMockRecorder) Size() *gomock.Call
Size indicates an expected call of Size.
func (*MockExecutableTaskTrackerMockRecorder) TrackTasks ¶ added in v1.21.0
func (mr *MockExecutableTaskTrackerMockRecorder) TrackTasks(exclusiveHighWatermarkInfo interface{}, tasks ...interface{}) *gomock.Call
TrackTasks indicates an expected call of TrackTasks.
type MockReceiverFlowController ¶ added in v1.25.0
type MockReceiverFlowController struct {
// contains filtered or unexported fields
}
MockReceiverFlowController is a mock of ReceiverFlowController interface.
func NewMockReceiverFlowController ¶ added in v1.25.0
func NewMockReceiverFlowController(ctrl *gomock.Controller) *MockReceiverFlowController
NewMockReceiverFlowController creates a new mock instance.
func (*MockReceiverFlowController) EXPECT ¶ added in v1.25.0
func (m *MockReceiverFlowController) EXPECT() *MockReceiverFlowControllerMockRecorder
EXPECT returns an object that allows the caller to indicate expected use.
func (*MockReceiverFlowController) GetFlowControlInfo ¶ added in v1.25.0
func (m *MockReceiverFlowController) GetFlowControlInfo(priority enums.TaskPriority) enums.ReplicationFlowControlCommand
GetFlowControlInfo mocks base method.
type MockReceiverFlowControllerMockRecorder ¶ added in v1.25.0
type MockReceiverFlowControllerMockRecorder struct {
// contains filtered or unexported fields
}
MockReceiverFlowControllerMockRecorder is the mock recorder for MockReceiverFlowController.
func (*MockReceiverFlowControllerMockRecorder) GetFlowControlInfo ¶ added in v1.25.0
func (mr *MockReceiverFlowControllerMockRecorder) GetFlowControlInfo(priority interface{}) *gomock.Call
GetFlowControlInfo indicates an expected call of GetFlowControlInfo.
type MockSenderFlowController ¶ added in v1.25.0
type MockSenderFlowController struct {
// contains filtered or unexported fields
}
MockSenderFlowController is a mock of SenderFlowController interface.
func NewMockSenderFlowController ¶ added in v1.25.0
func NewMockSenderFlowController(ctrl *gomock.Controller) *MockSenderFlowController
NewMockSenderFlowController creates a new mock instance.
func (*MockSenderFlowController) EXPECT ¶ added in v1.25.0
func (m *MockSenderFlowController) EXPECT() *MockSenderFlowControllerMockRecorder
EXPECT returns an object that allows the caller to indicate expected use.
func (*MockSenderFlowController) RefreshReceiverFlowControlInfo ¶ added in v1.25.0
func (m *MockSenderFlowController) RefreshReceiverFlowControlInfo(syncState *repication.SyncReplicationState)
RefreshReceiverFlowControlInfo mocks base method.
func (*MockSenderFlowController) Wait ¶ added in v1.25.0
func (m *MockSenderFlowController) Wait(priority enums.TaskPriority)
Wait mocks base method.
type MockSenderFlowControllerMockRecorder ¶ added in v1.25.0
type MockSenderFlowControllerMockRecorder struct {
// contains filtered or unexported fields
}
MockSenderFlowControllerMockRecorder is the mock recorder for MockSenderFlowController.
func (*MockSenderFlowControllerMockRecorder) RefreshReceiverFlowControlInfo ¶ added in v1.25.0
func (mr *MockSenderFlowControllerMockRecorder) RefreshReceiverFlowControlInfo(syncState interface{}) *gomock.Call
RefreshReceiverFlowControlInfo indicates an expected call of RefreshReceiverFlowControlInfo.
func (*MockSenderFlowControllerMockRecorder) Wait ¶ added in v1.25.0
func (mr *MockSenderFlowControllerMockRecorder) Wait(priority interface{}) *gomock.Call
Wait indicates an expected call of Wait.
type MockSourceTaskConverter ¶ added in v1.22.0
type MockSourceTaskConverter struct {
// contains filtered or unexported fields
}
MockSourceTaskConverter is a mock of SourceTaskConverter interface.
func NewMockSourceTaskConverter ¶ added in v1.22.0
func NewMockSourceTaskConverter(ctrl *gomock.Controller) *MockSourceTaskConverter
NewMockSourceTaskConverter creates a new mock instance.
func (*MockSourceTaskConverter) Convert ¶ added in v1.22.0
func (m *MockSourceTaskConverter) Convert(task tasks.Task) (*repication.ReplicationTask, error)
Convert mocks base method.
func (*MockSourceTaskConverter) EXPECT ¶ added in v1.22.0
func (m *MockSourceTaskConverter) EXPECT() *MockSourceTaskConverterMockRecorder
EXPECT returns an object that allows the caller to indicate expected use.
type MockSourceTaskConverterMockRecorder ¶ added in v1.22.0
type MockSourceTaskConverterMockRecorder struct {
// contains filtered or unexported fields
}
MockSourceTaskConverterMockRecorder is the mock recorder for MockSourceTaskConverter.
func (*MockSourceTaskConverterMockRecorder) Convert ¶ added in v1.22.0
func (mr *MockSourceTaskConverterMockRecorder) Convert(task interface{}) *gomock.Call
Convert indicates an expected call of Convert.
type MockStreamReceiver ¶ added in v1.21.0
type MockStreamReceiver struct {
// contains filtered or unexported fields
}
MockStreamReceiver is a mock of StreamReceiver interface.
func NewMockStreamReceiver ¶ added in v1.21.0
func NewMockStreamReceiver(ctrl *gomock.Controller) *MockStreamReceiver
NewMockStreamReceiver creates a new mock instance.
func (*MockStreamReceiver) EXPECT ¶ added in v1.21.0
func (m *MockStreamReceiver) EXPECT() *MockStreamReceiverMockRecorder
EXPECT returns an object that allows the caller to indicate expected use.
func (*MockStreamReceiver) IsValid ¶ added in v1.21.0
func (m *MockStreamReceiver) IsValid() bool
IsValid mocks base method.
func (*MockStreamReceiver) Key ¶ added in v1.21.0
func (m *MockStreamReceiver) Key() ClusterShardKeyPair
Key mocks base method.
func (*MockStreamReceiver) Stop ¶ added in v1.21.0
func (m *MockStreamReceiver) Stop()
Stop mocks base method.
type MockStreamReceiverMockRecorder ¶ added in v1.21.0
type MockStreamReceiverMockRecorder struct {
// contains filtered or unexported fields
}
MockStreamReceiverMockRecorder is the mock recorder for MockStreamReceiver.
func (*MockStreamReceiverMockRecorder) IsValid ¶ added in v1.21.0
func (mr *MockStreamReceiverMockRecorder) IsValid() *gomock.Call
IsValid indicates an expected call of IsValid.
func (*MockStreamReceiverMockRecorder) Key ¶ added in v1.21.0
func (mr *MockStreamReceiverMockRecorder) Key() *gomock.Call
Key indicates an expected call of Key.
func (*MockStreamReceiverMockRecorder) Stop ¶ added in v1.21.0
func (mr *MockStreamReceiverMockRecorder) Stop() *gomock.Call
Stop indicates an expected call of Stop.
type MockStreamSender ¶ added in v1.21.0
type MockStreamSender struct {
// contains filtered or unexported fields
}
MockStreamSender is a mock of StreamSender interface.
func NewMockStreamSender ¶ added in v1.21.0
func NewMockStreamSender(ctrl *gomock.Controller) *MockStreamSender
NewMockStreamSender creates a new mock instance.
func (*MockStreamSender) EXPECT ¶ added in v1.21.0
func (m *MockStreamSender) EXPECT() *MockStreamSenderMockRecorder
EXPECT returns an object that allows the caller to indicate expected use.
func (*MockStreamSender) IsValid ¶ added in v1.21.0
func (m *MockStreamSender) IsValid() bool
IsValid mocks base method.
func (*MockStreamSender) Key ¶ added in v1.21.0
func (m *MockStreamSender) Key() ClusterShardKeyPair
Key mocks base method.
func (*MockStreamSender) Stop ¶ added in v1.21.0
func (m *MockStreamSender) Stop()
Stop mocks base method.
type MockStreamSenderMockRecorder ¶ added in v1.21.0
type MockStreamSenderMockRecorder struct {
// contains filtered or unexported fields
}
MockStreamSenderMockRecorder is the mock recorder for MockStreamSender.
func (*MockStreamSenderMockRecorder) IsValid ¶ added in v1.21.0
func (mr *MockStreamSenderMockRecorder) IsValid() *gomock.Call
IsValid indicates an expected call of IsValid.
func (*MockStreamSenderMockRecorder) Key ¶ added in v1.21.0
func (mr *MockStreamSenderMockRecorder) Key() *gomock.Call
Key indicates an expected call of Key.
func (*MockStreamSenderMockRecorder) Stop ¶ added in v1.21.0
func (mr *MockStreamSenderMockRecorder) Stop() *gomock.Call
Stop indicates an expected call of Stop.
type MockTaskExecutor ¶
type MockTaskExecutor struct {
// contains filtered or unexported fields
}
MockTaskExecutor is a mock of TaskExecutor interface.
func NewMockTaskExecutor ¶
func NewMockTaskExecutor(ctrl *gomock.Controller) *MockTaskExecutor
NewMockTaskExecutor creates a new mock instance.
func (*MockTaskExecutor) EXPECT ¶
func (m *MockTaskExecutor) EXPECT() *MockTaskExecutorMockRecorder
EXPECT returns an object that allows the caller to indicate expected use.
func (*MockTaskExecutor) Execute ¶
func (m *MockTaskExecutor) Execute(ctx context.Context, replicationTask *repication.ReplicationTask, forceApply bool) error
Execute mocks base method.
type MockTaskExecutorMockRecorder ¶
type MockTaskExecutorMockRecorder struct {
// contains filtered or unexported fields
}
MockTaskExecutorMockRecorder is the mock recorder for MockTaskExecutor.
func (*MockTaskExecutorMockRecorder) Execute ¶
func (mr *MockTaskExecutorMockRecorder) Execute(ctx, replicationTask, forceApply interface{}) *gomock.Call
Execute indicates an expected call of Execute.
type MockTaskFetcherFactory ¶
type MockTaskFetcherFactory struct {
// contains filtered or unexported fields
}
MockTaskFetcherFactory is a mock of TaskFetcherFactory interface.
func NewMockTaskFetcherFactory ¶
func NewMockTaskFetcherFactory(ctrl *gomock.Controller) *MockTaskFetcherFactory
NewMockTaskFetcherFactory creates a new mock instance.
func (*MockTaskFetcherFactory) EXPECT ¶
func (m *MockTaskFetcherFactory) EXPECT() *MockTaskFetcherFactoryMockRecorder
EXPECT returns an object that allows the caller to indicate expected use.
func (*MockTaskFetcherFactory) GetOrCreateFetcher ¶
func (m *MockTaskFetcherFactory) GetOrCreateFetcher(clusterName string) taskFetcher
GetOrCreateFetcher mocks base method.
func (*MockTaskFetcherFactory) Start ¶
func (m *MockTaskFetcherFactory) Start()
Start mocks base method.
func (*MockTaskFetcherFactory) Stop ¶
func (m *MockTaskFetcherFactory) Stop()
Stop mocks base method.
type MockTaskFetcherFactoryMockRecorder ¶
type MockTaskFetcherFactoryMockRecorder struct {
// contains filtered or unexported fields
}
MockTaskFetcherFactoryMockRecorder is the mock recorder for MockTaskFetcherFactory.
func (*MockTaskFetcherFactoryMockRecorder) GetOrCreateFetcher ¶
func (mr *MockTaskFetcherFactoryMockRecorder) GetOrCreateFetcher(clusterName interface{}) *gomock.Call
GetOrCreateFetcher indicates an expected call of GetOrCreateFetcher.
func (*MockTaskFetcherFactoryMockRecorder) Start ¶
func (mr *MockTaskFetcherFactoryMockRecorder) Start() *gomock.Call
Start indicates an expected call of Start.
func (*MockTaskFetcherFactoryMockRecorder) Stop ¶
func (mr *MockTaskFetcherFactoryMockRecorder) Stop() *gomock.Call
Stop indicates an expected call of Stop.
type MockTaskProcessor ¶
type MockTaskProcessor struct {
// contains filtered or unexported fields
}
MockTaskProcessor is a mock of TaskProcessor interface.
func NewMockTaskProcessor ¶
func NewMockTaskProcessor(ctrl *gomock.Controller) *MockTaskProcessor
NewMockTaskProcessor creates a new mock instance.
func (*MockTaskProcessor) EXPECT ¶
func (m *MockTaskProcessor) EXPECT() *MockTaskProcessorMockRecorder
EXPECT returns an object that allows the caller to indicate expected use.
type MockTaskProcessorMockRecorder ¶
type MockTaskProcessorMockRecorder struct {
// contains filtered or unexported fields
}
MockTaskProcessorMockRecorder is the mock recorder for MockTaskProcessor.
func (*MockTaskProcessorMockRecorder) Start ¶
func (mr *MockTaskProcessorMockRecorder) Start() *gomock.Call
Start indicates an expected call of Start.
func (*MockTaskProcessorMockRecorder) Stop ¶
func (mr *MockTaskProcessorMockRecorder) Stop() *gomock.Call
Stop indicates an expected call of Stop.
type MockTrackableExecutableTask ¶ added in v1.21.0
type MockTrackableExecutableTask struct {
// contains filtered or unexported fields
}
MockTrackableExecutableTask is a mock of TrackableExecutableTask interface.
func NewMockTrackableExecutableTask ¶ added in v1.21.0
func NewMockTrackableExecutableTask(ctrl *gomock.Controller) *MockTrackableExecutableTask
NewMockTrackableExecutableTask creates a new mock instance.
func (*MockTrackableExecutableTask) Abort ¶ added in v1.21.0
func (m *MockTrackableExecutableTask) Abort()
Abort mocks base method.
func (*MockTrackableExecutableTask) Ack ¶ added in v1.21.0
func (m *MockTrackableExecutableTask) Ack()
Ack mocks base method.
func (*MockTrackableExecutableTask) Cancel ¶ added in v1.21.0
func (m *MockTrackableExecutableTask) Cancel()
Cancel mocks base method.
func (*MockTrackableExecutableTask) EXPECT ¶ added in v1.21.0
func (m *MockTrackableExecutableTask) EXPECT() *MockTrackableExecutableTaskMockRecorder
EXPECT returns an object that allows the caller to indicate expected use.
func (*MockTrackableExecutableTask) Execute ¶ added in v1.21.0
func (m *MockTrackableExecutableTask) Execute() error
Execute mocks base method.
func (*MockTrackableExecutableTask) HandleErr ¶ added in v1.21.0
func (m *MockTrackableExecutableTask) HandleErr(err error) error
HandleErr mocks base method.
func (*MockTrackableExecutableTask) IsRetryableError ¶ added in v1.21.0
func (m *MockTrackableExecutableTask) IsRetryableError(err error) bool
IsRetryableError mocks base method.
func (*MockTrackableExecutableTask) MarkPoisonPill ¶ added in v1.21.0
func (m *MockTrackableExecutableTask) MarkPoisonPill() error
MarkPoisonPill mocks base method.
func (*MockTrackableExecutableTask) Nack ¶ added in v1.21.0
func (m *MockTrackableExecutableTask) Nack(err error)
Nack mocks base method.
func (*MockTrackableExecutableTask) QueueID ¶ added in v1.21.0
func (m *MockTrackableExecutableTask) QueueID() interface{}
QueueID mocks base method.
func (*MockTrackableExecutableTask) Reschedule ¶ added in v1.21.0
func (m *MockTrackableExecutableTask) Reschedule()
Reschedule mocks base method.
func (*MockTrackableExecutableTask) RetryPolicy ¶ added in v1.21.0
func (m *MockTrackableExecutableTask) RetryPolicy() backoff.RetryPolicy
RetryPolicy mocks base method.
func (*MockTrackableExecutableTask) SourceClusterName ¶ added in v1.25.0
func (m *MockTrackableExecutableTask) SourceClusterName() string
SourceClusterName mocks base method.
func (*MockTrackableExecutableTask) State ¶ added in v1.21.0
func (m *MockTrackableExecutableTask) State() tasks.State
State mocks base method.
func (*MockTrackableExecutableTask) TaskCreationTime ¶ added in v1.21.0
func (m *MockTrackableExecutableTask) TaskCreationTime() time.Time
TaskCreationTime mocks base method.
func (*MockTrackableExecutableTask) TaskID ¶ added in v1.21.0
func (m *MockTrackableExecutableTask) TaskID() int64
TaskID mocks base method.
type MockTrackableExecutableTaskMockRecorder ¶ added in v1.21.0
type MockTrackableExecutableTaskMockRecorder struct {
// contains filtered or unexported fields
}
MockTrackableExecutableTaskMockRecorder is the mock recorder for MockTrackableExecutableTask.
func (*MockTrackableExecutableTaskMockRecorder) Abort ¶ added in v1.21.0
func (mr *MockTrackableExecutableTaskMockRecorder) Abort() *gomock.Call
Abort indicates an expected call of Abort.
func (*MockTrackableExecutableTaskMockRecorder) Ack ¶ added in v1.21.0
func (mr *MockTrackableExecutableTaskMockRecorder) Ack() *gomock.Call
Ack indicates an expected call of Ack.
func (*MockTrackableExecutableTaskMockRecorder) Cancel ¶ added in v1.21.0
func (mr *MockTrackableExecutableTaskMockRecorder) Cancel() *gomock.Call
Cancel indicates an expected call of Cancel.
func (*MockTrackableExecutableTaskMockRecorder) Execute ¶ added in v1.21.0
func (mr *MockTrackableExecutableTaskMockRecorder) Execute() *gomock.Call
Execute indicates an expected call of Execute.
func (*MockTrackableExecutableTaskMockRecorder) HandleErr ¶ added in v1.21.0
func (mr *MockTrackableExecutableTaskMockRecorder) HandleErr(err interface{}) *gomock.Call
HandleErr indicates an expected call of HandleErr.
func (*MockTrackableExecutableTaskMockRecorder) IsRetryableError ¶ added in v1.21.0
func (mr *MockTrackableExecutableTaskMockRecorder) IsRetryableError(err interface{}) *gomock.Call
IsRetryableError indicates an expected call of IsRetryableError.
func (*MockTrackableExecutableTaskMockRecorder) MarkPoisonPill ¶ added in v1.21.0
func (mr *MockTrackableExecutableTaskMockRecorder) MarkPoisonPill() *gomock.Call
MarkPoisonPill indicates an expected call of MarkPoisonPill.
func (*MockTrackableExecutableTaskMockRecorder) Nack ¶ added in v1.21.0
func (mr *MockTrackableExecutableTaskMockRecorder) Nack(err interface{}) *gomock.Call
Nack indicates an expected call of Nack.
func (*MockTrackableExecutableTaskMockRecorder) QueueID ¶ added in v1.21.0
func (mr *MockTrackableExecutableTaskMockRecorder) QueueID() *gomock.Call
QueueID indicates an expected call of QueueID.
func (*MockTrackableExecutableTaskMockRecorder) Reschedule ¶ added in v1.21.0
func (mr *MockTrackableExecutableTaskMockRecorder) Reschedule() *gomock.Call
Reschedule indicates an expected call of Reschedule.
func (*MockTrackableExecutableTaskMockRecorder) RetryPolicy ¶ added in v1.21.0
func (mr *MockTrackableExecutableTaskMockRecorder) RetryPolicy() *gomock.Call
RetryPolicy indicates an expected call of RetryPolicy.
func (*MockTrackableExecutableTaskMockRecorder) SourceClusterName ¶ added in v1.25.0
func (mr *MockTrackableExecutableTaskMockRecorder) SourceClusterName() *gomock.Call
SourceClusterName indicates an expected call of SourceClusterName.
func (*MockTrackableExecutableTaskMockRecorder) State ¶ added in v1.21.0
func (mr *MockTrackableExecutableTaskMockRecorder) State() *gomock.Call
State indicates an expected call of State.
func (*MockTrackableExecutableTaskMockRecorder) TaskCreationTime ¶ added in v1.21.0
func (mr *MockTrackableExecutableTaskMockRecorder) TaskCreationTime() *gomock.Call
TaskCreationTime indicates an expected call of TaskCreationTime.
func (*MockTrackableExecutableTaskMockRecorder) TaskID ¶ added in v1.21.0
func (mr *MockTrackableExecutableTaskMockRecorder) TaskID() *gomock.Call
TaskID indicates an expected call of TaskID.
type MocktaskFetcher ¶
type MocktaskFetcher struct {
// contains filtered or unexported fields
}
MocktaskFetcher is a mock of taskFetcher interface.
func NewMocktaskFetcher ¶
func NewMocktaskFetcher(ctrl *gomock.Controller) *MocktaskFetcher
NewMocktaskFetcher creates a new mock instance.
func (*MocktaskFetcher) EXPECT ¶
func (m *MocktaskFetcher) EXPECT() *MocktaskFetcherMockRecorder
EXPECT returns an object that allows the caller to indicate expected use.
type MocktaskFetcherMockRecorder ¶
type MocktaskFetcherMockRecorder struct {
// contains filtered or unexported fields
}
MocktaskFetcherMockRecorder is the mock recorder for MocktaskFetcher.
func (*MocktaskFetcherMockRecorder) Stop ¶
func (mr *MocktaskFetcherMockRecorder) Stop() *gomock.Call
Stop indicates an expected call of Stop.
type NoopDLQWriter ¶ added in v1.23.0
type NoopDLQWriter struct { }
NoopDLQWriter is a DLQWriter that does nothing. The zero value is a valid instance.
func (NoopDLQWriter) WriteTaskToDLQ ¶ added in v1.23.0
func (w NoopDLQWriter) WriteTaskToDLQ(context.Context, DLQWriteRequest) error
type ProcessToolBox ¶ added in v1.21.0
type ProcessToolBox struct { fx.In Config *configs.Config ClusterMetadata cluster.Metadata ClientBean client.Bean ShardController shard.Controller NamespaceCache namespace.Registry EagerNamespaceRefresher EagerNamespaceRefresher NDCHistoryResender xdc.NDCHistoryResender HighPriorityTaskScheduler ctasks.Scheduler[TrackableExecutableTask] `name:"HighPriorityTaskScheduler"` // consider using a single TaskScheduler i.e. InterleavedWeightedRoundRobinScheduler instead of two LowPriorityTaskScheduler ctasks.Scheduler[TrackableExecutableTask] `name:"LowPriorityTaskScheduler"` MetricsHandler metrics.Handler Logger log.Logger EventSerializer serialization.Serializer DLQWriter DLQWriter HistoryEventsHandler eventhandler.HistoryEventsHandler }
type ReceiverFlowController ¶ added in v1.25.0
type ReceiverFlowController interface {
GetFlowControlInfo(priority enums.TaskPriority) enums.ReplicationFlowControlCommand
}
type ReceiverMode ¶ added in v1.25.0
type ReceiverMode int32
const ( ReceiverModeUnset ReceiverMode = 0 ReceiverModeSingleStack ReceiverMode = 1 // default mode. It only uses High Priority Task Tracker for processing tasks. ReceiverModeTieredStack ReceiverMode = 2 )
type SenderFlowController ¶ added in v1.25.0
type SenderFlowController interface { // Wait will block go routine until the sender is allowed to send a task Wait(priority enums.TaskPriority) RefreshReceiverFlowControlInfo(syncState *replicationpb.SyncReplicationState) }
type SenderFlowControllerImpl ¶ added in v1.25.0
type SenderFlowControllerImpl struct {
// contains filtered or unexported fields
}
func NewSenderFlowController ¶ added in v1.25.0
func NewSenderFlowController(config *configs.Config, logger log.Logger) *SenderFlowControllerImpl
func (*SenderFlowControllerImpl) RefreshReceiverFlowControlInfo ¶ added in v1.25.0
func (s *SenderFlowControllerImpl) RefreshReceiverFlowControlInfo(syncState *replicationpb.SyncReplicationState)
func (*SenderFlowControllerImpl) Wait ¶ added in v1.25.0
func (s *SenderFlowControllerImpl) Wait(priority enums.TaskPriority)
type SequentialBatchableTaskQueue ¶ added in v1.23.0
func (*SequentialBatchableTaskQueue) Add ¶ added in v1.23.0
func (q *SequentialBatchableTaskQueue) Add(task TrackableExecutableTask)
Add will try to batch input task with the last task in the queue. Since most likely incoming task are ordered by task ID, we only try to batch incoming task with last task in the queue.
func (*SequentialBatchableTaskQueue) ID ¶ added in v1.23.0
func (q *SequentialBatchableTaskQueue) ID() interface{}
func (*SequentialBatchableTaskQueue) IsEmpty ¶ added in v1.23.0
func (q *SequentialBatchableTaskQueue) IsEmpty() bool
func (*SequentialBatchableTaskQueue) Len ¶ added in v1.23.0
func (q *SequentialBatchableTaskQueue) Len() int
func (*SequentialBatchableTaskQueue) Peek ¶ added in v1.23.0
func (q *SequentialBatchableTaskQueue) Peek() TrackableExecutableTask
func (*SequentialBatchableTaskQueue) Remove ¶ added in v1.23.0
func (q *SequentialBatchableTaskQueue) Remove() (task TrackableExecutableTask)
type SequentialTaskQueue ¶ added in v1.21.0
func (*SequentialTaskQueue) Add ¶ added in v1.21.0
func (q *SequentialTaskQueue) Add(task TrackableExecutableTask)
func (*SequentialTaskQueue) ID ¶ added in v1.21.0
func (q *SequentialTaskQueue) ID() interface{}
func (*SequentialTaskQueue) IsEmpty ¶ added in v1.21.0
func (q *SequentialTaskQueue) IsEmpty() bool
func (*SequentialTaskQueue) Len ¶ added in v1.21.0
func (q *SequentialTaskQueue) Len() int
func (*SequentialTaskQueue) Peek ¶ added in v1.21.0
func (q *SequentialTaskQueue) Peek() TrackableExecutableTask
func (*SequentialTaskQueue) Remove ¶ added in v1.21.0
func (q *SequentialTaskQueue) Remove() TrackableExecutableTask
type SourceTaskConverter ¶ added in v1.22.0
type SourceTaskConverter interface {
Convert(task tasks.Task) (*replicationspb.ReplicationTask, error)
}
type SourceTaskConverterImpl ¶ added in v1.22.0
type SourceTaskConverterImpl struct {
// contains filtered or unexported fields
}
func NewSourceTaskConverter ¶ added in v1.22.0
func (*SourceTaskConverterImpl) Convert ¶ added in v1.22.0
func (c *SourceTaskConverterImpl) Convert( task tasks.Task, ) (*replicationspb.ReplicationTask, error)
type SourceTaskConverterProvider ¶ added in v1.22.0
type Stream ¶ added in v1.21.0
type Stream BiDirectionStream[*adminservice.StreamWorkflowReplicationMessagesRequest, *adminservice.StreamWorkflowReplicationMessagesResponse]
type StreamBiDirectionStreamClientProvider ¶ added in v1.21.0
type StreamBiDirectionStreamClientProvider struct {
// contains filtered or unexported fields
}
func NewStreamBiDirectionStreamClientProvider ¶ added in v1.21.0
func NewStreamBiDirectionStreamClientProvider( clusterMetadata cluster.Metadata, clientBean client.Bean, ) *StreamBiDirectionStreamClientProvider
func (*StreamBiDirectionStreamClientProvider) Get ¶ added in v1.21.0
func (p *StreamBiDirectionStreamClientProvider) Get( ctx context.Context, clientShardKey ClusterShardKey, serverShardKey ClusterShardKey, ) (BiDirectionStreamClient[*adminservice.StreamWorkflowReplicationMessagesRequest, *adminservice.StreamWorkflowReplicationMessagesResponse], error)
type StreamError ¶ added in v1.24.0
type StreamError struct { Message string // contains filtered or unexported fields }
func NewStreamError ¶ added in v1.24.0
type StreamReceiver ¶ added in v1.21.0
type StreamReceiver interface { IsValid() bool Key() ClusterShardKeyPair Stop() }
type StreamReceiverImpl ¶ added in v1.21.0
type StreamReceiverImpl struct { ProcessToolBox // contains filtered or unexported fields }
func NewStreamReceiver ¶ added in v1.21.0
func NewStreamReceiver( processToolBox ProcessToolBox, taskConverter ExecutableTaskConverter, clientShardKey ClusterShardKey, serverShardKey ClusterShardKey, ) *StreamReceiverImpl
func (*StreamReceiverImpl) IsValid ¶ added in v1.21.0
func (r *StreamReceiverImpl) IsValid() bool
func (*StreamReceiverImpl) Key ¶ added in v1.21.0
func (r *StreamReceiverImpl) Key() ClusterShardKeyPair
func (*StreamReceiverImpl) Start ¶ added in v1.21.0
func (r *StreamReceiverImpl) Start()
Start starts the processor
func (*StreamReceiverImpl) Stop ¶ added in v1.21.0
func (r *StreamReceiverImpl) Stop()
Stop stops the processor
type StreamReceiverMonitor ¶ added in v1.21.0
type StreamReceiverMonitor interface { RegisterInboundStream(streamSender StreamSender) Start() Stop() }
type StreamReceiverMonitorImpl ¶ added in v1.21.0
type StreamReceiverMonitorImpl struct { ProcessToolBox sync.Mutex // contains filtered or unexported fields }
func NewStreamReceiverMonitor ¶ added in v1.21.0
func NewStreamReceiverMonitor( processToolBox ProcessToolBox, executableTaskConverter ExecutableTaskConverter, enableStreaming bool, ) *StreamReceiverMonitorImpl
func (*StreamReceiverMonitorImpl) RegisterInboundStream ¶ added in v1.21.0
func (m *StreamReceiverMonitorImpl) RegisterInboundStream( streamSender StreamSender, )
func (*StreamReceiverMonitorImpl) Start ¶ added in v1.21.0
func (m *StreamReceiverMonitorImpl) Start()
func (*StreamReceiverMonitorImpl) Stop ¶ added in v1.21.0
func (m *StreamReceiverMonitorImpl) Stop()
type StreamResp ¶ added in v1.21.0
type StreamSender ¶ added in v1.21.0
type StreamSender interface { IsValid() bool Key() ClusterShardKeyPair Stop() }
type StreamSenderImpl ¶ added in v1.21.0
type StreamSenderImpl struct {
// contains filtered or unexported fields
}
func NewStreamSender ¶ added in v1.21.0
func NewStreamSender( server historyservice.HistoryService_StreamWorkflowReplicationMessagesServer, shardContext shard.Context, historyEngine shard.Engine, taskConverter SourceTaskConverter, clientClusterName string, clientClusterShardCount int32, clientShardKey ClusterShardKey, serverShardKey ClusterShardKey, config *configs.Config, ) *StreamSenderImpl
func (*StreamSenderImpl) IsValid ¶ added in v1.21.0
func (s *StreamSenderImpl) IsValid() bool
func (*StreamSenderImpl) Key ¶ added in v1.21.0
func (s *StreamSenderImpl) Key() ClusterShardKeyPair
func (*StreamSenderImpl) Start ¶ added in v1.21.0
func (s *StreamSenderImpl) Start()
func (*StreamSenderImpl) Stop ¶ added in v1.21.0
func (s *StreamSenderImpl) Stop()
func (*StreamSenderImpl) Wait ¶ added in v1.21.0
func (s *StreamSenderImpl) Wait()
type TaskExecutor ¶
type TaskExecutor interface {
Execute(ctx context.Context, replicationTask *replicationspb.ReplicationTask, forceApply bool) error
}
func NewTaskExecutor ¶
func NewTaskExecutor( remoteCluster string, shardContext shard.Context, nDCHistoryResender xdc.NDCHistoryResender, deleteManager deletemanager.DeleteManager, workflowCache wcache.Cache, ) TaskExecutor
NewTaskExecutor creates a replication task executor The executor uses by 1) DLQ replication task handler 2) history replication task processor
type TaskExecutorParams ¶
type TaskExecutorParams struct { RemoteCluster string // TODO: Remove this remote cluster from executor then it can use singleton. Shard shard.Context HistoryResender xdc.NDCHistoryResender DeleteManager deletemanager.DeleteManager WorkflowCache wcache.Cache }
type TaskExecutorProvider ¶
type TaskExecutorProvider func(params TaskExecutorParams) TaskExecutor
type TaskFetcherFactory ¶
type TaskFetcherFactory interface { GetOrCreateFetcher(clusterName string) taskFetcher Start() Stop() }
TaskFetcherFactory is a group of fetchers, one per source DC.
type TaskParser ¶ added in v1.23.0
type TaskParser interface {
ParseReplicationTask(replicationTask *persistencespb.ReplicationTaskInfo) (tasks.Task, error)
}
TaskParser is a trimmed version of go.temporal.io/server/common/persistence/serialization.Serializer that only provides the methods we need.
type TaskProcessor ¶
type TaskProcessor interface { Start() Stop() }
TaskProcessor is the interface for task processor
func NewTaskProcessor ¶
func NewTaskProcessor( sourceShardID int32, shard shard.Context, historyEngine shard.Engine, config *configs.Config, metricsHandler metrics.Handler, replicationTaskFetcher taskFetcher, replicationTaskExecutor TaskExecutor, eventSerializer serialization.Serializer, dlqWriter DLQWriter, ) TaskProcessor
NewTaskProcessor creates a new replication task processor.
type TrackableExecutableTask ¶ added in v1.21.0
type WatermarkInfo ¶ added in v1.21.0
Source Files ¶
- ack_manager.go
- ack_manager_mock.go
- batchable_task.go
- batchable_task_mock.go
- bi_direction_stream.go
- dlq_handler.go
- dlq_handler_mock.go
- dlq_writer.go
- eager_namespace_refresher.go
- eager_namespace_refresher_mock.go
- executable_activity_state_task.go
- executable_history_task.go
- executable_noop_task.go
- executable_task.go
- executable_task_converter.go
- executable_task_mock.go
- executable_task_tool_box.go
- executable_task_tracker.go
- executable_task_tracker_mock.go
- executable_unknown_task.go
- executable_workflow_state_task.go
- fx.go
- grpc_stream_client.go
- metrics.go
- noop_dlq_writer.go
- poller_manager.go
- raw_task_converter.go
- raw_task_converter_mock.go
- sequential_batch_queue.go
- sequential_queue.go
- stream.go
- stream_receiver.go
- stream_receiver_flow_controller.go
- stream_receiver_flow_controller_mock.go
- stream_receiver_mock.go
- stream_receiver_monitor.go
- stream_sender.go
- stream_sender_flow_controller.go
- stream_sender_flow_controller_mock.go
- stream_sender_mock.go
- task_executor.go
- task_executor_mock.go
- task_fetcher.go
- task_fetcher_mock.go
- task_processor.go
- task_processor_manager.go
- task_processor_mock.go
Directories ¶
Path | Synopsis |
---|---|
Package eventhandler is a generated GoMock package.
|
Package eventhandler is a generated GoMock package. |