|
| 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