78
serverShardKey ClusterShardKey,
79
config *configs.Config,
81
>
logger := log.With(
82
>
shardContext.GetLogger(),
83
>
tag.TargetCluster(clientClusterName), // client is the target cluster (passive cluster)
84
>
tag.TargetShardID(clientShardKey.ShardID),
85
>
tag.ShardID(serverShardKey.ShardID), // server is the source cluster (active cluster)
86
>
tag.Operation("replication-stream-sender"),
87
>
)
88
>
return &StreamSenderImpl{
89
>
server: server,
90
>
shardContext: shardContext,
91
>
historyEngine: historyEngine,
92
>
taskConverter: taskConverter,
93
>
metrics: shardContext.GetMetricsHandler(),
94
>
logger: logger,
95
>
status: common.DaemonStatusInitialized,
96
>
clientClusterName: clientClusterName,
97
>
clientShardKey: clientShardKey,
98
>
serverShardKey: serverShardKey,
99
>
clientClusterShardCount: clientClusterShardCount,
100
>
recvSignalChan: make(chan struct{}, 1),
101
>
shutdownChan: channel.NewShutdownOnce(),
102
>
config: config,
103
>
isTieredStackEnabled: config.EnableReplicationTaskTieredProcessing(),
104
>
flowController: NewSenderFlowController(config, logger),
105
>
ssRateLimiter: ssRateLimiter,
106
>
}
107
>
}
108
109
func (s *StreamSenderImpl) Start() {