60
}
61
62
>
func NamespaceHandoverWorkflow(ctx workflow.Context, params NamespaceHandoverParams) (retErr error) {
handover_workflow.go
63
>
if err := validateAndSetNamespaceHandoverParams(¶ms); err != nil {
64
return err
65
}
67
>
InitialInterval: time.Second,
68
>
MaximumInterval: time.Second,
69
>
BackoffCoefficient: 1,
70
>
}
71
>
ao := workflow.ActivityOptions{
72
>
StartToCloseTimeout: time.Second * 10,
73
>
RetryPolicy: retryPolicy,
74
>
}
75
>
ctx = workflow.WithActivityOptions(ctx, ao)
76
>
77
>
// ** Step 1: Get Cluster Metadata **
78
>
var metadataResp MetadataResponse
79
>
MetadataRequest := MetadataRequest{Namespace: params.Namespace}
80
>
var a *activities
81
>
err := workflow.ExecuteActivity(ctx, a.GetMetadata, MetadataRequest).Get(ctx, &metadataResp)
82
>
if err != nil {
83
return err
84
}
85
86
// ** Step 2: Get current replication status **
88
>
err = workflow.ExecuteActivity(ctx, a.GetMaxReplicationTaskIDs).Get(ctx, &repStatus)
89
>
if err != nil {
90
return err
91
}
92
93
// ** Step 3: Wait for Remote Cluster to catch-up on Replication Tasks
95
>
StartToCloseTimeout: time.Hour,
96
>
HeartbeatTimeout: time.Second * 10,
97
>
RetryPolicy: retryPolicy,
98
>
}
99
>
ctx2 := workflow.WithActivityOptions(ctx, ao2)
100
>
waitRequest := WaitReplicationRequest{
101
>
Namespace: params.Namespace,
102
>
ShardCount: metadataResp.ShardCount,
103
>
RemoteCluster: params.RemoteCluster,
104
>
AllowedLagging: time.Duration(params.AllowedLaggingSeconds) * time.Second,
105
>
WaitForTaskIds: repStatus.MaxReplicationTaskIds,
106
>
AllowedLaggingTasks: params.AllowedLaggingTasks,
107
>
}
108
>
err = workflow.ExecuteActivity(ctx2, a.WaitReplication, waitRequest).Get(ctx2, nil)
109
>
if err != nil {
110
return err
111
}