290
inclusiveLowWaterMarkTime = highPriorityWaterMarkInfo.Timestamp
291
}
293
return 0, NewStreamError("InclusiveLowWaterMark is not set", serviceerror.NewInternal("Invalid inclusive low watermark"))
294
}
295
296
>
if err := stream.Send(&adminservice.StreamWorkflowReplicationMessagesRequest{
stream_receiver.go
297
>
Attributes: &adminservice.StreamWorkflowReplicationMessagesRequest_SyncReplicationState{
298
>
SyncReplicationState: &replicationspb.SyncReplicationState{
299
>
InclusiveLowWatermark: inclusiveLowWaterMark,
300
>
InclusiveLowWatermarkTime: timestamppb.New(inclusiveLowWaterMarkTime),
301
>
HighPriorityState: highPriorityWatermark,
302
>
LowPriorityState: lowPriorityWatermark,
303
>
},
304
>
},
305
>
}); err != nil {
306
return 0, NewStreamError("stream_receiver failed to send", err)
307
}
308
>
metrics.ReplicationTasksRecvBacklog.With(r.MetricsHandler).Record(
stream_receiver.go
309
>
int64(size),
310
>
metrics.FromClusterIDTag(r.serverShardKey.ClusterID),
311
>
metrics.ToClusterIDTag(r.clientShardKey.ClusterID),
312
>
)
313
>
metrics.ReplicationTasksSend.With(r.MetricsHandler).Record(
314
>
int64(1),
315
>
metrics.FromClusterIDTag(r.clientShardKey.ClusterID),
316
>
metrics.ToClusterIDTag(r.serverShardKey.ClusterID),
317
>
metrics.OperationTag(metrics.SyncWatermarkScope),
318
>
)
319
>
return inclusiveLowWaterMark, nil
320
}
321