83
clientShardKey ClusterShardKey,
84
serverShardKey ClusterShardKey,
86
>
logger := log.With(
87
>
processToolBox.Logger,
88
>
tag.SourceCluster(processToolBox.ClusterMetadata.ClusterNameForFailoverVersion(true, int64(serverShardKey.ClusterID))),
89
>
tag.SourceShardID(serverShardKey.ShardID),
90
>
tag.ShardID(clientShardKey.ShardID), // client is the local cluster (target cluster, passive cluster)
91
>
tag.Operation("replication-stream-receiver"),
92
>
)
93
>
highPriorityTaskTracker := NewExecutableTaskTracker(logger, processToolBox.MetricsHandler)
94
>
lowPriorityTaskTracker := NewExecutableTaskTracker(logger, processToolBox.MetricsHandler)
95
>
receiver := &StreamReceiverImpl{
96
>
ProcessToolBox: processToolBox,
97
>
98
>
status: common.DaemonStatusInitialized,
99
>
clientShardKey: clientShardKey,
100
>
serverShardKey: serverShardKey,
101
>
highPriorityTaskTracker: highPriorityTaskTracker,
102
>
lowPriorityTaskTracker: lowPriorityTaskTracker,
103
>
shutdownChan: channel.NewShutdownOnce(),
104
>
logger: logger,
105
>
stream: newStream(
106
>
processToolBox,
107
>
clientShardKey,
108
>
serverShardKey,
109
>
),
110
>
taskConverter: taskConverter,
111
>
receiverMode: ReceiverModeUnset,
112
>
slowSubmissionTimestamps: make(map[enumsspb.TaskPriority]time.Time),
113
>
recvSignalChan: make(chan struct{}, 1),
114
>
}
115
>
taskTrackerMap := make(map[enumsspb.TaskPriority]FlowControlSignalProvider)
116
>
taskTrackerMap[enumsspb.TASK_PRIORITY_HIGH] = func() *FlowControlSignal {
117
return &FlowControlSignal{
118
taskTrackingCount: highPriorityTaskTracker.Size(),