Sitelet https://github.com/kubernetes/kops/commit/2aa9c38d670dd3bd78e470a9364c0fe36d1792b6
Skip to content

Commit 2aa9c38

Browse files
committed
Restore pkg/assets/assetcopy/ package from master
1 parent 901b471 commit 2aa9c38

4 files changed

Lines changed: 526 additions & 0 deletions

File tree

‎pkg/assets/assetcopy/copy.go‎

Lines changed: 126 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,126 @@
1+
/*
2+
Copyright 2021 The Kubernetes Authors.
3+
4+
Licensed under the Apache License, Version 2.0 (the "License");
5+
you may not use this file except in compliance with the License.
6+
You may obtain a copy of the License at
7+
8+
http://www.apache.org/licenses/LICENSE-2.0
9+
10+
Unless required by applicable law or agreed to in writing, software
11+
distributed under the License is distributed on an "AS IS" BASIS,
12+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
See the License for the specific language governing permissions and
14+
limitations under the License.
15+
*/
16+
17+
// Package assetcopy copies cluster assets to a private file repository or container registry
18+
// (`kops get assets --copy`). It is separate from pkg/assets so that only the kops CLI links
19+
// the container registry client libraries; runtime binaries (kops-controller, nodeup) import
20+
// pkg/assets without inheriting those dependencies.
21+
package assetcopy
22+
23+
import (
24+
"fmt"
25+
"sort"
26+
27+
"k8s.io/klog/v2"
28+
"k8s.io/kops/pkg/apis/kops"
29+
"k8s.io/kops/pkg/assets"
30+
"k8s.io/kops/util/pkg/vfs"
31+
)
32+
33+
type assetTask interface {
34+
Run() error
35+
}
36+
37+
func Copy(imageAssets []*assets.ImageAsset, fileAssets []*assets.FileAsset, vfsContext *vfs.VFSContext, cluster *kops.Cluster) error {
38+
tasks := map[string]assetTask{}
39+
40+
for _, imageAsset := range imageAssets {
41+
if imageAsset.DownloadLocation != imageAsset.CanonicalLocation {
42+
copyImageTask := &CopyImage{
43+
Name: imageAsset.DownloadLocation,
44+
SourceImage: imageAsset.CanonicalLocation,
45+
TargetImage: imageAsset.DownloadLocation,
46+
}
47+
48+
if existing, ok := tasks[copyImageTask.Name]; ok {
49+
if existing.(*CopyImage).SourceImage != copyImageTask.SourceImage {
50+
return fmt.Errorf("different sources for same image target %s: %s vs %s", copyImageTask.Name, copyImageTask.SourceImage, existing.(*CopyImage).SourceImage)
51+
}
52+
}
53+
54+
tasks[copyImageTask.Name] = copyImageTask
55+
}
56+
}
57+
58+
for _, fileAsset := range fileAssets {
59+
if fileAsset.DownloadURL.String() != fileAsset.CanonicalURL.String() {
60+
copyFileTask := &CopyFile{
61+
Name: fileAsset.CanonicalURL.String(),
62+
TargetFile: fileAsset.DownloadURL.String(),
63+
SourceFile: fileAsset.CanonicalURL.String(),
64+
SHA: fileAsset.SHAValue.Hex(),
65+
VFSContext: vfsContext,
66+
Cluster: cluster,
67+
}
68+
69+
if existing, ok := tasks[copyFileTask.Name]; ok {
70+
e, ok := existing.(*CopyFile)
71+
if !ok {
72+
return fmt.Errorf("different types for copy target %s", copyFileTask.Name)
73+
}
74+
if e.TargetFile != copyFileTask.TargetFile {
75+
return fmt.Errorf("different targets for same file %s: %s vs %s", copyFileTask.Name, copyFileTask.TargetFile, e.TargetFile)
76+
}
77+
if e.SHA != copyFileTask.SHA {
78+
return fmt.Errorf("different sha for same file %s: %s vs %s", copyFileTask.Name, copyFileTask.SHA, e.SHA)
79+
}
80+
}
81+
82+
tasks[copyFileTask.Name] = copyFileTask
83+
}
84+
}
85+
86+
ch := make(chan error, 5)
87+
for i := 0; i < cap(ch); i++ {
88+
ch <- nil
89+
}
90+
91+
gotError := false
92+
names := make([]string, 0, len(tasks))
93+
for name := range tasks {
94+
names = append(names, name)
95+
}
96+
sort.Strings(names)
97+
for _, name := range names {
98+
task := tasks[name]
99+
err := <-ch
100+
if err != nil {
101+
klog.Warning(err)
102+
gotError = true
103+
}
104+
go func(n string, t assetTask) {
105+
err := t.Run()
106+
if err != nil {
107+
err = fmt.Errorf("%s: %v", n, err)
108+
}
109+
ch <- err
110+
}(name, task)
111+
}
112+
113+
for i := 0; i < cap(ch); i++ {
114+
err := <-ch
115+
if err != nil {
116+
klog.Warning(err)
117+
gotError = true
118+
}
119+
}
120+
121+
close(ch)
122+
if gotError {
123+
return fmt.Errorf("not all assets copied successfully")
124+
}
125+
return nil
126+
}

‎pkg/assets/assetcopy/copyfile.go‎

Lines changed: 220 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,220 @@
1+
/*
2+
Copyright 2017 The Kubernetes Authors.
3+
4+
Licensed under the Apache License, Version 2.0 (the "License");
5+
you may not use this file except in compliance with the License.
6+
You may obtain a copy of the License at
7+
8+
http://www.apache.org/licenses/LICENSE-2.0
9+
10+
Unless required by applicable law or agreed to in writing, software
11+
distributed under the License is distributed on an "AS IS" BASIS,
12+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
See the License for the specific language governing permissions and
14+
limitations under the License.
15+
*/
16+
17+
package assetcopy
18+
19+
import (
20+
"bytes"
21+
"context"
22+
"fmt"
23+
"net/url"
24+
"os"
25+
"strings"
26+
27+
"k8s.io/klog/v2"
28+
"k8s.io/kops/pkg/apis/kops"
29+
"k8s.io/kops/util/pkg/hashing"
30+
"k8s.io/kops/util/pkg/vfs"
31+
"k8s.io/kops/util/pkg/vfs/acls"
32+
)
33+
34+
// CopyFile copies an from a source file repository, to a target repository,
35+
// typically used for highly secure clusters.
36+
type CopyFile struct {
37+
Name string
38+
SourceFile string
39+
TargetFile string
40+
SHA string
41+
VFSContext *vfs.VFSContext
42+
Cluster *kops.Cluster
43+
}
44+
45+
// fileExtensionForSHA returns the expected extension for the given hash
46+
// If the hash length is not recognized, it returns an error.
47+
func fileExtensionForSHA(sha string) (string, error) {
48+
switch len(sha) {
49+
case 40:
50+
return ".sha1", nil
51+
case 64:
52+
return ".sha256", nil
53+
default:
54+
return "", fmt.Errorf("unhandled sha length for %q", sha)
55+
}
56+
}
57+
58+
func (e *CopyFile) Run() error {
59+
ctx := context.TODO()
60+
61+
expectedSHA := strings.TrimSpace(e.SHA)
62+
63+
shaExtension, err := fileExtensionForSHA(expectedSHA)
64+
if err != nil {
65+
return err
66+
}
67+
68+
targetSHAFile := e.TargetFile + shaExtension
69+
70+
targetSHABytes, err := e.VFSContext.ReadFile(targetSHAFile)
71+
if err != nil {
72+
if os.IsNotExist(err) {
73+
klog.V(4).Infof("unable to download: %q, assuming target file is not present, and if not present may not be an error: %v",
74+
targetSHAFile, err)
75+
} else {
76+
klog.V(4).Infof("unable to download: %q, %v", targetSHAFile, err)
77+
}
78+
} else {
79+
targetSHA := string(targetSHABytes)
80+
81+
if strings.TrimSpace(targetSHA) == expectedSHA {
82+
klog.V(8).Infof("found matching target sha for file: %q", e.TargetFile)
83+
return nil
84+
}
85+
86+
klog.V(8).Infof("did not find same file, found mismatching target sha1 for file: %q", e.TargetFile)
87+
}
88+
89+
source := e.SourceFile
90+
target := e.TargetFile
91+
sourceSha := e.SHA
92+
93+
klog.V(2).Infof("copying bits from %q to %q", source, target)
94+
95+
if err := transferFile(ctx, e.VFSContext, e.Cluster, source, target, sourceSha); err != nil {
96+
return fmt.Errorf("unable to transfer %q to %q: %v", source, target, err)
97+
}
98+
99+
return nil
100+
}
101+
102+
// transferFile downloads a file from the source location, validates the file matches the SHA,
103+
// and uploads the file to the target location.
104+
func transferFile(ctx context.Context, vfsContext *vfs.VFSContext, cluster *kops.Cluster, source string, target string, sha string) error {
105+
// TODO drop file to disk, as vfs reads file into memory. We load kubelet into memory for instance.
106+
// TODO in s3 can we do a copy file ... would need to test
107+
108+
data, err := vfsContext.ReadFile(source)
109+
if err != nil {
110+
if os.IsNotExist(err) {
111+
return fmt.Errorf("file not found %q: %v", source, err)
112+
}
113+
114+
return fmt.Errorf("error downloading file %q: %v", source, err)
115+
}
116+
117+
objectStore, err := buildVFSPath(target)
118+
if err != nil {
119+
return err
120+
}
121+
122+
uploadVFS, err := vfsContext.BuildVfsPath(objectStore)
123+
if err != nil {
124+
return fmt.Errorf("error building path %q: %v", objectStore, err)
125+
}
126+
127+
shaExtension, err := fileExtensionForSHA(sha)
128+
if err != nil {
129+
return err
130+
}
131+
132+
shaTarget := objectStore + shaExtension
133+
shaVFS, err := vfsContext.BuildVfsPath(shaTarget)
134+
if err != nil {
135+
return fmt.Errorf("error building path %q: %v", shaTarget, err)
136+
}
137+
138+
shaHash, err := hashing.FromString(strings.TrimSpace(sha))
139+
if err != nil {
140+
return fmt.Errorf("unable to parse sha: %q, %v", sha, err)
141+
}
142+
143+
in := bytes.NewReader(data)
144+
dataHash, err := shaHash.Algorithm.Hash(in)
145+
if err != nil {
146+
return fmt.Errorf("unable to hash file %q downloaded: %v", source, err)
147+
}
148+
149+
if !shaHash.Equal(dataHash) {
150+
return fmt.Errorf("the sha value in %q does not match %q calculated value %q", shaTarget, source, dataHash.String())
151+
}
152+
153+
klog.Infof("uploading %q to %q", source, objectStore)
154+
if err := writeFile(ctx, cluster, uploadVFS, data); err != nil {
155+
return err
156+
}
157+
158+
b := []byte(shaHash.Hex())
159+
if err := writeFile(ctx, cluster, shaVFS, b); err != nil {
160+
return err
161+
}
162+
163+
return nil
164+
}
165+
166+
func writeFile(ctx context.Context, cluster *kops.Cluster, p vfs.Path, data []byte) error {
167+
acl, err := acls.GetACL(ctx, p, cluster)
168+
if err != nil {
169+
return err
170+
}
171+
172+
if err = p.WriteFile(ctx, bytes.NewReader(data), acl); err != nil {
173+
return fmt.Errorf("error writing path %v: %v", p, err)
174+
}
175+
176+
return nil
177+
}
178+
179+
// buildVFSPath task a recognizable https url and transforms that URL into the equivalent url with the object
180+
// store prefix.
181+
func buildVFSPath(target string) (string, error) {
182+
if !strings.Contains(target, "://") || strings.HasPrefix(target, "memfs://") || strings.HasPrefix(target, "file://") {
183+
return target, nil
184+
}
185+
186+
var vfsPath string
187+
188+
// Matches all S3 regional naming conventions:
189+
// https://docs.aws.amazon.com/general/latest/gr/rande.html#s3_region
190+
// and converts to a s3://<bucket>/<path> vfsPath
191+
s3VfsPath, err := vfs.VFSPath(target)
192+
if err == nil {
193+
vfsPath = s3VfsPath
194+
} else {
195+
// These matches only cover a subset of the URLs that you can use, but I am uncertain how to cover more of the possible
196+
// options.
197+
// This code parses the HOST and determines gs URLs.
198+
// For instance you can have the bucket name in the gs url hostname.
199+
u, err := url.Parse(target)
200+
if err != nil {
201+
return "", fmt.Errorf("Unable to parse Google Cloud Storage URL: %q", target)
202+
}
203+
if u.Host == "storage.googleapis.com" {
204+
vfsPath = "gs:/" + u.Path
205+
}
206+
}
207+
208+
if vfsPath == "" {
209+
klog.Errorf("Unable to determine VFS path from supplied URL: %s", target)
210+
klog.Errorf("S3, Google Cloud Storage, and File Paths are supported.")
211+
klog.Errorf("For S3, please make sure that the supplied file repository URL adhere to S3 naming conventions, https://docs.aws.amazon.com/general/latest/gr/rande.html#s3_region.")
212+
klog.Errorf("For GCS, please make sure that the supplied file repository URL adheres to https://storage.googleapis.com/")
213+
if err != nil { // print the S3 error for more details
214+
return "", fmt.Errorf("Error Details: %v", err)
215+
}
216+
return "", fmt.Errorf("unable to determine vfs type for %q", target)
217+
}
218+
219+
return vfsPath, nil
220+
}

0 commit comments

Comments
 (0)