mirror of
https://github.com/pgsty/minio.git
synced 2026-09-10 04:24:03 +03:00
Merge pull request #163 from pgsty/fix/issue-158-federated-copy-sse
fix: forward plaintext on federated CopyObject of SSE objects (#158)
This commit is contained in:
@@ -433,7 +433,16 @@ func putOptsFromHeaders(ctx context.Context, hdr http.Header, metadata map[strin
|
||||
if err != nil {
|
||||
return ObjectOptions{}, err
|
||||
}
|
||||
sseKms, err := encrypt.NewSSEKMS(keyID, context)
|
||||
// kms.Context implements encoding.TextMarshaler, so handing it to the
|
||||
// SDK's interface{} parameter would serialize the context as a JSON
|
||||
// string, which the receiving ParseHTTP rejects; a nil Context is a
|
||||
// typed nil there and would be sent as "{}". Pass a plain map, or
|
||||
// nothing when no context was requested.
|
||||
var sdkContext any
|
||||
if context != nil {
|
||||
sdkContext = map[string]string(context)
|
||||
}
|
||||
sseKms, err := encrypt.NewSSEKMS(keyID, sdkContext)
|
||||
if err != nil {
|
||||
return ObjectOptions{}, err
|
||||
}
|
||||
|
||||
@@ -20,20 +20,28 @@ package cmd
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"crypto/md5"
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"encoding/xml"
|
||||
"io"
|
||||
"maps"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
|
||||
miniogo "github.com/minio/minio-go/v7"
|
||||
miniocredentials "github.com/minio/minio-go/v7/pkg/credentials"
|
||||
"github.com/minio/minio-go/v7/pkg/set"
|
||||
"github.com/minio/minio/internal/auth"
|
||||
"github.com/minio/minio/internal/config/dns"
|
||||
"github.com/minio/minio/internal/crypto"
|
||||
"github.com/minio/minio/internal/hash"
|
||||
xhttp "github.com/minio/minio/internal/http"
|
||||
"github.com/minio/minio/internal/kms"
|
||||
)
|
||||
|
||||
// federationRemoteCapture records the exact request headers the remote
|
||||
@@ -91,6 +99,26 @@ func (c *federationRemoteCapture) reservedKeys() []string {
|
||||
// to the default region, exactly as TestAPIFederatedCopyObjectPartChecksum does.
|
||||
func setupCopyObjectFederation(t *testing.T, objectAPI ObjectLayer, apiRouter http.Handler,
|
||||
instanceType, srcBucket string, responseFilters ...func(http.Header),
|
||||
) (remoteBucket string, capture *federationRemoteCapture, cleanup func()) {
|
||||
t.Helper()
|
||||
return setupCopyObjectFederationRemote(t, objectAPI, apiRouter, instanceType, srcBucket, false, responseFilters...)
|
||||
}
|
||||
|
||||
// setupCopyObjectFederationTLS is setupCopyObjectFederation with a TLS remote
|
||||
// endpoint and globalIsTLS set, so an SSE-C copy is accepted on both hops: the
|
||||
// proxy's own SSE-C transport gate and the minio-go client's SSE-C policy. The
|
||||
// proxy client is built exactly as in production, trusting the test
|
||||
// certificate for the duration.
|
||||
func setupCopyObjectFederationTLS(t *testing.T, objectAPI ObjectLayer, apiRouter http.Handler,
|
||||
instanceType, srcBucket string,
|
||||
) (remoteBucket string, cleanup func()) {
|
||||
t.Helper()
|
||||
remoteBucket, _, cleanup = setupCopyObjectFederationRemote(t, objectAPI, apiRouter, instanceType, srcBucket, true)
|
||||
return remoteBucket, cleanup
|
||||
}
|
||||
|
||||
func setupCopyObjectFederationRemote(t *testing.T, objectAPI ObjectLayer, apiRouter http.Handler,
|
||||
instanceType, srcBucket string, secure bool, responseFilters ...func(http.Header),
|
||||
) (remoteBucket string, capture *federationRemoteCapture, cleanup func()) {
|
||||
t.Helper()
|
||||
remoteBucket = getRandomBucketName()
|
||||
@@ -99,13 +127,19 @@ func setupCopyObjectFederation(t *testing.T, objectAPI ObjectLayer, apiRouter ht
|
||||
}
|
||||
|
||||
capture = &federationRemoteCapture{}
|
||||
remote := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
handler := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
capture.record(r.Header)
|
||||
if len(responseFilters) != 0 && r.Method == http.MethodPut {
|
||||
w = federationResponseFilter{ResponseWriter: w, filter: responseFilters[0]}
|
||||
}
|
||||
apiRouter.ServeHTTP(w, r)
|
||||
}))
|
||||
})
|
||||
var remote *httptest.Server
|
||||
if secure {
|
||||
remote = httptest.NewTLSServer(handler)
|
||||
} else {
|
||||
remote = httptest.NewServer(handler)
|
||||
}
|
||||
host, port, _ := strings.Cut(remote.Listener.Addr().String(), ":")
|
||||
|
||||
globalObjLayerMutex.Lock()
|
||||
@@ -113,6 +147,7 @@ func setupCopyObjectFederation(t *testing.T, objectAPI ObjectLayer, apiRouter ht
|
||||
globalObjectAPI = remoteBucketObjectLayer{ObjectLayer: previousLayer, remoteBucket: remoteBucket}
|
||||
globalObjLayerMutex.Unlock()
|
||||
previousDNS, previousFederation, previousIPs := globalDNSConfig, globalBucketFederation, globalDomainIPs
|
||||
previousTLS, previousClient := globalIsTLS, getRemoteInstanceClient
|
||||
globalDNSConfig = federationTestDNS{records: map[string][]dns.SrvRecord{
|
||||
srcBucket: {{Host: host, Port: json.Number(port)}},
|
||||
remoteBucket: {{Host: host, Port: json.Number(port)}},
|
||||
@@ -121,6 +156,23 @@ func setupCopyObjectFederation(t *testing.T, objectAPI ObjectLayer, apiRouter ht
|
||||
// middleware always serves locally and only the handler proxies.
|
||||
globalDomainIPs = set.CreateStringSet(remote.Listener.Addr().String())
|
||||
globalBucketFederation = true
|
||||
if secure {
|
||||
globalIsTLS = true
|
||||
transport := remote.Client().Transport
|
||||
getRemoteInstanceClient = func(r *http.Request, host string) (*miniogo.Core, error) {
|
||||
cred := getReqAccessCred(r, globalSite.Region())
|
||||
core, err := miniogo.NewCore(host, &miniogo.Options{
|
||||
Creds: miniocredentials.NewStaticV4(cred.AccessKey, cred.SecretKey, ""),
|
||||
Secure: true,
|
||||
Transport: transport,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
core.SetAppInfo(federatedInternalAppName, ReleaseTag)
|
||||
return core, nil
|
||||
}
|
||||
}
|
||||
|
||||
cleanup = func() {
|
||||
remote.Close()
|
||||
@@ -128,6 +180,7 @@ func setupCopyObjectFederation(t *testing.T, objectAPI ObjectLayer, apiRouter ht
|
||||
globalObjectAPI = previousLayer
|
||||
globalObjLayerMutex.Unlock()
|
||||
globalDNSConfig, globalBucketFederation, globalDomainIPs = previousDNS, previousFederation, previousIPs
|
||||
globalIsTLS, getRemoteInstanceClient = previousTLS, previousClient
|
||||
}
|
||||
return remoteBucket, capture, cleanup
|
||||
}
|
||||
@@ -437,3 +490,303 @@ func testAPIFederatedCopyObjectInheritedChecksum(objectAPI ObjectLayer, instance
|
||||
assertCopyChecksumResponse(t, rec, hash.ChecksumCRC32, data)
|
||||
assertCopyChecksum(t, objectAPI, remoteBucket, dstObject, hash.ChecksumCRC32, data, false, nil)
|
||||
}
|
||||
|
||||
// federationTestKMSKeyID names the single key of the builtin KMS the SSE test
|
||||
// installs; SSE-S3 seals with it implicitly, SSE-KMS names it by ID.
|
||||
const federationTestKMSKeyID = "federation-test-key"
|
||||
|
||||
// federationSSEHeaders returns the request headers that select one server-side
|
||||
// encryption kind: "plain", "s3", "kms", "kms-context" (SSE-KMS with an
|
||||
// explicit encryption context) or "c". SSE-C keys derive from keyByte so a
|
||||
// case can name two distinct customer keys. With copySource the SSE-C key is
|
||||
// returned in its x-amz-copy-source-* form, the only kind a copy has to name
|
||||
// for its source; the other kinds then return nothing.
|
||||
func federationSSEHeaders(kind string, keyByte byte, copySource bool) map[string]string {
|
||||
h := map[string]string{}
|
||||
if copySource && kind != "c" {
|
||||
return h
|
||||
}
|
||||
switch kind {
|
||||
case "plain":
|
||||
case "s3":
|
||||
h[xhttp.AmzServerSideEncryption] = xhttp.AmzEncryptionAES
|
||||
case "kms", "kms-context":
|
||||
h[xhttp.AmzServerSideEncryption] = xhttp.AmzEncryptionKMS
|
||||
h[xhttp.AmzServerSideEncryptionKmsID] = federationTestKMSKeyID
|
||||
if kind == "kms-context" {
|
||||
h[xhttp.AmzServerSideEncryptionKmsContext] = base64.StdEncoding.EncodeToString([]byte(`{"tenant":"federation"}`))
|
||||
}
|
||||
case "c":
|
||||
key := bytes.Repeat([]byte{keyByte}, 32)
|
||||
sum := md5.Sum(key)
|
||||
algorithm, customerKey, keyMD5 := xhttp.AmzServerSideEncryptionCustomerAlgorithm,
|
||||
xhttp.AmzServerSideEncryptionCustomerKey, xhttp.AmzServerSideEncryptionCustomerKeyMD5
|
||||
if copySource {
|
||||
algorithm, customerKey, keyMD5 = xhttp.AmzServerSideEncryptionCopyCustomerAlgorithm,
|
||||
xhttp.AmzServerSideEncryptionCopyCustomerKey, xhttp.AmzServerSideEncryptionCopyCustomerKeyMD5
|
||||
}
|
||||
h[algorithm] = xhttp.AmzEncryptionAES
|
||||
h[customerKey] = base64.StdEncoding.EncodeToString(key)
|
||||
h[keyMD5] = base64.StdEncoding.EncodeToString(sum[:])
|
||||
default:
|
||||
panic("unknown SSE kind " + kind)
|
||||
}
|
||||
return h
|
||||
}
|
||||
|
||||
// federationStoredSSE reports which SSE kind stored object metadata declares.
|
||||
func federationStoredSSE(metadata map[string]string) string {
|
||||
switch {
|
||||
case crypto.SSEC.IsEncrypted(metadata):
|
||||
return "c"
|
||||
case crypto.S3KMS.IsEncrypted(metadata):
|
||||
return "kms"
|
||||
case crypto.S3.IsEncrypted(metadata):
|
||||
return "s3"
|
||||
}
|
||||
return "plain"
|
||||
}
|
||||
|
||||
// federationGetObject reads an object through GetObjectHandler.
|
||||
func federationGetObject(t *testing.T, apiRouter http.Handler, credentials auth.Credentials,
|
||||
bucket, object string, headers map[string]string,
|
||||
) *httptest.ResponseRecorder {
|
||||
t.Helper()
|
||||
req, err := newTestSignedRequestV4(http.MethodGet, getGetObjectURL("", bucket, object),
|
||||
0, nil, credentials.AccessKey, credentials.SecretKey, headers)
|
||||
if err != nil {
|
||||
t.Fatalf("failed to build GetObject request: %v", err)
|
||||
}
|
||||
rec := httptest.NewRecorder()
|
||||
apiRouter.ServeHTTP(rec, req)
|
||||
return rec
|
||||
}
|
||||
|
||||
// TestAPIFederatedCopyObjectSSE guards the legacy etcd federation branch of
|
||||
// CopyObjectHandler for encrypted sources and destinations (#158). The proxy
|
||||
// reads its source through getObjectNInfo, which yields the decrypted and
|
||||
// decompressed bytes, but it used to run the destination encryption locally
|
||||
// as well and then forward that stream with the source's stored size under
|
||||
// the destination SSE option. SSE to plain and plain to SSE failed on the
|
||||
// length mismatch; SSE to SSE matched by coincidence, so the remote encrypted
|
||||
// the ciphertext a second time and a destination GET returned the inner
|
||||
// ciphertext with HTTP 200. The proxy must forward the logical bytes with
|
||||
// their logical size and let the remote encrypt exactly once.
|
||||
func TestAPIFederatedCopyObjectSSE(t *testing.T) {
|
||||
defer DetectTestLeak(t)()
|
||||
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
|
||||
t: t,
|
||||
objAPITest: testAPIFederatedCopyObjectSSE,
|
||||
// The test router registers routes in this order, and the plain
|
||||
// PutObject route has no query matcher, so the multipart routes must
|
||||
// precede it or they would never be reached.
|
||||
endpoints: []string{
|
||||
"NewMultipart", "PutObjectPart", "CompleteMultipart",
|
||||
"CopyObject", "PutObject", "HeadObject", "GetObject",
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
// putFederationMultipartSource uploads parts as one multipart object through
|
||||
// the API router; headers apply to the NewMultipartUpload request.
|
||||
func putFederationMultipartSource(t *testing.T, apiRouter http.Handler, credentials auth.Credentials,
|
||||
bucket, object string, parts [][]byte, headers map[string]string,
|
||||
) {
|
||||
t.Helper()
|
||||
do := func(method, url string, body []byte, headers map[string]string) *httptest.ResponseRecorder {
|
||||
req, err := newTestSignedRequestV4(method, url, int64(len(body)), bytes.NewReader(body),
|
||||
credentials.AccessKey, credentials.SecretKey, headers)
|
||||
if err != nil {
|
||||
t.Fatalf("failed to build %s %s request: %v", method, url, err)
|
||||
}
|
||||
rec := httptest.NewRecorder()
|
||||
apiRouter.ServeHTTP(rec, req)
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Fatalf("%s %s failed: %d %s", method, url, rec.Code, rec.Body.String())
|
||||
}
|
||||
return rec
|
||||
}
|
||||
var upload InitiateMultipartUploadResponse
|
||||
rec := do(http.MethodPost, getNewMultipartURL("", bucket, object), nil, headers)
|
||||
if err := xml.Unmarshal(rec.Body.Bytes(), &upload); err != nil {
|
||||
t.Fatalf("failed to decode NewMultipartUpload response: %v", err)
|
||||
}
|
||||
completion := CompleteMultipartUpload{}
|
||||
for i, part := range parts {
|
||||
rec := do(http.MethodPut, getPutObjectPartURL("", bucket, object, upload.UploadID, strconv.Itoa(i+1)), part, nil)
|
||||
completion.Parts = append(completion.Parts, CompletePart{PartNumber: i + 1, ETag: canonicalizeETag(rec.Header()[xhttp.ETag][0])})
|
||||
}
|
||||
body, err := xml.Marshal(completion)
|
||||
if err != nil {
|
||||
t.Fatalf("failed to encode CompleteMultipartUpload: %v", err)
|
||||
}
|
||||
do(http.MethodPost, getCompleteMultipartUploadURL("", bucket, object, upload.UploadID), body, nil)
|
||||
}
|
||||
|
||||
func testAPIFederatedCopyObjectSSE(objectAPI ObjectLayer, instanceType, bucketName string,
|
||||
apiRouter http.Handler, credentials auth.Credentials, t *testing.T,
|
||||
) {
|
||||
testKMS, err := kms.NewBuiltin(federationTestKMSKeyID, bytes.Repeat([]byte{0x58}, 32))
|
||||
if err != nil {
|
||||
t.Fatalf("unable to create the test KMS: %v", err)
|
||||
}
|
||||
previousKMS := GlobalKMS
|
||||
GlobalKMS = testKMS
|
||||
defer func() { GlobalKMS = previousKMS }()
|
||||
// .txt sources are stored compressed, SSE-S3 and SSE-KMS ones included;
|
||||
// SSE-C data is never compressed.
|
||||
restoreCompression := setCopyChecksumCompression(true)
|
||||
defer restoreCompression()
|
||||
|
||||
remoteBucket, cleanup := setupCopyObjectFederationTLS(t, objectAPI, apiRouter, instanceType, bucketName)
|
||||
defer cleanup()
|
||||
|
||||
// A DARE package holds 64 KiB of plaintext, so the larger bodies span two
|
||||
// packages and catch a size that accounts for only one.
|
||||
large := bytes.Repeat([]byte("federated sse copy body "), 64*1024/24+1)[:64*1024+1]
|
||||
bodies := []struct {
|
||||
name, ext string
|
||||
data []byte
|
||||
}{
|
||||
{name: "small", ext: ".bin", data: []byte("abc")},
|
||||
{name: "two-packages", ext: ".bin", data: large},
|
||||
{name: "compressed", ext: ".txt", data: large},
|
||||
}
|
||||
pairs := []struct{ src, dst string }{
|
||||
{"plain", "s3"},
|
||||
{"s3", "plain"},
|
||||
{"s3", "s3"},
|
||||
{"plain", "c"},
|
||||
{"c", "plain"},
|
||||
{"c", "c"},
|
||||
{"s3", "c"},
|
||||
{"plain", "kms"},
|
||||
{"kms", "plain"},
|
||||
{"kms", "kms"},
|
||||
{"plain", "kms-context"},
|
||||
{"kms-context", "kms-context"},
|
||||
}
|
||||
const srcKeyByte, dstKeyByte = 0x11, 0x22
|
||||
|
||||
for _, body := range bodies {
|
||||
for _, pair := range pairs {
|
||||
t.Run(body.name+"/"+pair.src+"-to-"+pair.dst, func(t *testing.T) {
|
||||
prefix := "federation/sse-" + body.name + "-" + pair.src + "-to-" + pair.dst
|
||||
srcObject, dstObject := prefix+"-source"+body.ext, prefix+"-destination"+body.ext
|
||||
|
||||
putCopyChecksumSource(t, apiRouter, credentials, bucketName, srcObject, body.data,
|
||||
federationSSEHeaders(pair.src, srcKeyByte, false))
|
||||
before, err := objectAPI.GetObjectInfo(t.Context(), bucketName, srcObject, ObjectOptions{})
|
||||
if err != nil {
|
||||
t.Fatalf("%s: GetObjectInfo(source) failed: %v", instanceType, err)
|
||||
}
|
||||
if got, want := federationStoredSSE(before.UserDefined), strings.TrimSuffix(pair.src, "-context"); got != want {
|
||||
t.Fatalf("%s: source stored as %s, want %s", instanceType, got, want)
|
||||
}
|
||||
if compressed := body.ext == ".txt" && pair.src != "c"; before.IsCompressed() != compressed {
|
||||
t.Fatalf("%s: source compressed=%v, want %v", instanceType, before.IsCompressed(), compressed)
|
||||
}
|
||||
|
||||
headers := federationSSEHeaders(pair.dst, dstKeyByte, false)
|
||||
maps.Copy(headers, federationSSEHeaders(pair.src, srcKeyByte, true))
|
||||
rec := federatedCopyRequest(t, apiRouter, credentials, bucketName, srcObject, remoteBucket, dstObject, headers)
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Fatalf("%s: federated CopyObject failed: %d %s", instanceType, rec.Code, rec.Body.String())
|
||||
}
|
||||
// The remote computes the S3 default CRC-64NVME over the bytes it
|
||||
// received, so the returned checksum must be the plaintext's.
|
||||
assertCopyChecksumResponse(t, rec, hash.ChecksumCRC64NVME, body.data)
|
||||
|
||||
// A destination GET must return the original plaintext at its length.
|
||||
var getHeaders map[string]string
|
||||
if pair.dst == "c" {
|
||||
getHeaders = federationSSEHeaders("c", dstKeyByte, false)
|
||||
}
|
||||
got := federationGetObject(t, apiRouter, credentials, remoteBucket, dstObject, getHeaders)
|
||||
if got.Code != http.StatusOK {
|
||||
t.Fatalf("%s: destination GET failed: %d %s", instanceType, got.Code, got.Body.String())
|
||||
}
|
||||
if !bytes.Equal(got.Body.Bytes(), body.data) {
|
||||
t.Fatalf("%s: destination GET returned %d bytes that differ from the %d-byte plaintext",
|
||||
instanceType, got.Body.Len(), len(body.data))
|
||||
}
|
||||
if want := strconv.Itoa(len(body.data)); got.Header().Get(xhttp.ContentLength) != want {
|
||||
t.Fatalf("%s: destination Content-Length = %q, want %q",
|
||||
instanceType, got.Header().Get(xhttp.ContentLength), want)
|
||||
}
|
||||
|
||||
// The destination is stored under the requested SSE kind and was
|
||||
// encrypted exactly once: a second layer would add its own DARE
|
||||
// package overhead to the stored size.
|
||||
after, err := objectAPI.GetObjectInfo(t.Context(), remoteBucket, dstObject, ObjectOptions{})
|
||||
if err != nil {
|
||||
t.Fatalf("%s: GetObjectInfo(destination) failed: %v", instanceType, err)
|
||||
}
|
||||
if got, want := federationStoredSSE(after.UserDefined), strings.TrimSuffix(pair.dst, "-context"); got != want {
|
||||
t.Fatalf("%s: destination stored as %s, want %s", instanceType, got, want)
|
||||
}
|
||||
if pair.dst != "plain" && !after.IsCompressed() {
|
||||
once := ObjectInfo{Size: int64(len(body.data))}
|
||||
if want := once.EncryptedSize(); after.Size != want {
|
||||
t.Fatalf("%s: destination stored %d bytes, want %d for %d bytes encrypted once",
|
||||
instanceType, after.Size, want, len(body.data))
|
||||
}
|
||||
}
|
||||
|
||||
// The source must be untouched.
|
||||
source, err := objectAPI.GetObjectInfo(t.Context(), bucketName, srcObject, ObjectOptions{})
|
||||
if err != nil {
|
||||
t.Fatalf("%s: GetObjectInfo(source) after the copy failed: %v", instanceType, err)
|
||||
}
|
||||
if source.Size != before.Size || source.ETag != before.ETag || !source.ModTime.Equal(before.ModTime) {
|
||||
t.Fatalf("%s: federated copy modified its source: size %d->%d etag %s->%s",
|
||||
instanceType, before.Size, source.Size, before.ETag, source.ETag)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// A multipart SSE-S3 source is encrypted per part, so its logical size is
|
||||
// the sum of the parts' decrypted sizes and the decrypting reader crosses
|
||||
// a part boundary; the copy must still forward exactly the plaintext.
|
||||
parts := [][]byte{bytes.Repeat([]byte("p"), globalMinPartSize+1), []byte("q!")}
|
||||
data := bytes.Join(parts, nil)
|
||||
for _, dst := range []string{"plain", "s3"} {
|
||||
t.Run("multipart/s3-to-"+dst, func(t *testing.T) {
|
||||
srcObject, dstObject := "federation/sse-multipart-source-"+dst+".bin", "federation/sse-multipart-destination-"+dst+".bin"
|
||||
putFederationMultipartSource(t, apiRouter, credentials, bucketName, srcObject, parts,
|
||||
federationSSEHeaders("s3", srcKeyByte, false))
|
||||
before, err := objectAPI.GetObjectInfo(t.Context(), bucketName, srcObject, ObjectOptions{})
|
||||
if err != nil {
|
||||
t.Fatalf("%s: GetObjectInfo(source) failed: %v", instanceType, err)
|
||||
}
|
||||
if len(before.Parts) != len(parts) || federationStoredSSE(before.UserDefined) != "s3" {
|
||||
t.Fatalf("%s: source has %d parts stored as %s, want %d parts as s3",
|
||||
instanceType, len(before.Parts), federationStoredSSE(before.UserDefined), len(parts))
|
||||
}
|
||||
|
||||
rec := federatedCopyRequest(t, apiRouter, credentials, bucketName, srcObject, remoteBucket, dstObject,
|
||||
federationSSEHeaders(dst, dstKeyByte, false))
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Fatalf("%s: federated CopyObject failed: %d %s", instanceType, rec.Code, rec.Body.String())
|
||||
}
|
||||
assertCopyChecksumResponse(t, rec, hash.ChecksumCRC64NVME, data)
|
||||
got := federationGetObject(t, apiRouter, credentials, remoteBucket, dstObject, nil)
|
||||
if got.Code != http.StatusOK || !bytes.Equal(got.Body.Bytes(), data) {
|
||||
t.Fatalf("%s: destination GET = %d with %d bytes, want 200 with the %d-byte plaintext",
|
||||
instanceType, got.Code, got.Body.Len(), len(data))
|
||||
}
|
||||
after, err := objectAPI.GetObjectInfo(t.Context(), remoteBucket, dstObject, ObjectOptions{})
|
||||
if err != nil {
|
||||
t.Fatalf("%s: GetObjectInfo(destination) failed: %v", instanceType, err)
|
||||
}
|
||||
if got := federationStoredSSE(after.UserDefined); got != dst {
|
||||
t.Fatalf("%s: destination stored as %s, want %s", instanceType, got, dst)
|
||||
}
|
||||
if once := (ObjectInfo{Size: int64(len(data))}); dst == "s3" && after.Size != once.EncryptedSize() {
|
||||
t.Fatalf("%s: destination stored %d bytes, want %d for %d bytes encrypted once",
|
||||
instanceType, after.Size, once.EncryptedSize(), len(data))
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
+21
-8
@@ -1518,12 +1518,22 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re
|
||||
}
|
||||
}
|
||||
|
||||
// Federation only: the destination bucket lives on another deployment and
|
||||
// the copy is forwarded to it as a PutObject. That remote write owns the
|
||||
// destination's storage transformations, so this handler hands it the
|
||||
// logical (decompressed, decrypted) bytes at their logical size and lets
|
||||
// the remote compress and encrypt once. Encrypting here as well would
|
||||
// forward ciphertext under the destination's own SSE option, which either
|
||||
// fails the length check or, for SSE to SSE, has the remote encrypt the
|
||||
// ciphertext a second time and store an unreadable object (#158).
|
||||
remoteCallRequired := isRemoteCopyRequired(ctx, srcBucket, dstBucket, objectAPI)
|
||||
|
||||
var compressMetadata map[string]string
|
||||
// No need to compress for remote etcd calls
|
||||
// Pass the decompressed stream to such calls.
|
||||
isDstCompressed := isCompressible(r.Header, dstObject) &&
|
||||
length > minCompressibleSize &&
|
||||
!isRemoteCopyRequired(ctx, srcBucket, dstBucket, objectAPI)
|
||||
!remoteCallRequired
|
||||
if isDstCompressed {
|
||||
compressMetadata = make(map[string]string, 2)
|
||||
// Preserving the compression metadata.
|
||||
@@ -1652,6 +1662,9 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re
|
||||
var targetSize int64
|
||||
|
||||
switch {
|
||||
case remoteCallRequired:
|
||||
// The remote receives the logical bytes, see above.
|
||||
targetSize = actualSize
|
||||
case isDstCompressed:
|
||||
targetSize = -1
|
||||
case !isSourceEncrypted && !isTargetEncrypted:
|
||||
@@ -1722,7 +1735,7 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re
|
||||
pReader.setChecksumReader(checksumReader)
|
||||
}
|
||||
|
||||
if isTargetEncrypted {
|
||||
if isTargetEncrypted && !remoteCallRequired {
|
||||
var encReader io.Reader
|
||||
kind, _ := crypto.IsRequested(r.Header)
|
||||
encReader, objEncKey, err = newEncryptReader(ctx, srcInfo.Reader, kind, newKeyID, newKey, dstBucket, dstObject, encMetadata, kmsCtx)
|
||||
@@ -1746,7 +1759,7 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re
|
||||
return
|
||||
}
|
||||
|
||||
if isTargetEncrypted {
|
||||
if isTargetEncrypted && !remoteCallRequired {
|
||||
pReader, err = pReader.WithEncryption(srcInfo.Reader, &objEncKey)
|
||||
if err != nil {
|
||||
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
|
||||
@@ -1899,9 +1912,6 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re
|
||||
srcInfo.metadataOnly = false
|
||||
}
|
||||
|
||||
// Federation only.
|
||||
remoteCallRequired := isRemoteCopyRequired(ctx, srcBucket, dstBucket, objectAPI)
|
||||
|
||||
var objInfo ObjectInfo
|
||||
var os *objSweeper
|
||||
if remoteCallRequired {
|
||||
@@ -1948,7 +1958,7 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re
|
||||
var checksumHeaderValue string
|
||||
switch {
|
||||
case wantChecksumType.IsSet():
|
||||
if srcInfo.Size == 0 {
|
||||
if actualSize == 0 {
|
||||
// minio-go streams no trailing checksum for an empty body, so the
|
||||
// remote would compute none, the bind below would fail, and every
|
||||
// empty-object federated copy would 500. Forward the empty-content
|
||||
@@ -1977,8 +1987,11 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re
|
||||
if checksumHeaderValue != "" {
|
||||
opts.UserMetadata[wantChecksumType.Key()] = checksumHeaderValue
|
||||
}
|
||||
// srcInfo.Reader yields the logical bytes, so declare the logical size:
|
||||
// srcInfo.Size is the stored size, which differs for an encrypted or
|
||||
// compressed source.
|
||||
remoteObjInfo, rerr := core.PutObject(ctx, dstBucket, dstObject, srcInfo.Reader,
|
||||
srcInfo.Size, "", "", opts)
|
||||
actualSize, "", "", opts)
|
||||
if rerr != nil {
|
||||
writeErrorResponse(ctx, w, toAPIError(ctx, rerr), r.URL)
|
||||
return
|
||||
|
||||
Reference in New Issue
Block a user