diff --git a/inventory/file.go b/inventory/file.go index 3b344bd9..070b2805 100644 --- a/inventory/file.go +++ b/inventory/file.go @@ -134,8 +134,12 @@ type ( } RelocateEntityParameter struct { - Entity *ent.Entity - NewSource string + Entity *ent.Entity + NewSource string + // NewEncryptMetadata describes the encryption of the relocated blob. It is set when the + // blob has been encrypted or decrypted while being transferred, so that the entity is + // committed along with the key its new blob can be read with. + NewEncryptMetadata *types.EncryptMetadata ParentFiles []int PrimaryEntityParentFiles []int } diff --git a/pkg/filemanager/manager/entity.go b/pkg/filemanager/manager/entity.go index 27ce4108..3627d087 100644 --- a/pkg/filemanager/manager/entity.go +++ b/pkg/filemanager/manager/entity.go @@ -45,6 +45,11 @@ type ( ListPhysical(ctx context.Context, path string, policyID int, recursive bool, progress driver.ListProgressFunc) ([]fs.PhysicalObject, error) // ImportPhysical imports a physical file to a Cloudreve file ImportPhysical(ctx context.Context, dst *fs.URI, policyId int, src fs.PhysicalObject, completeHook bool) error + // RelocateBlob transfers the blob of an entity into another storage policy, encrypting it + // on the fly if the destination policy has file encryption enabled. It returns the + // encryption metadata to be persisted onto the entity, or nil if the encryption state of + // the entity is unchanged. + RelocateBlob(ctx context.Context, args *RelocateBlobArgs) (*types.EncryptMetadata, error) } DirectLink struct { File fs.File diff --git a/pkg/filemanager/manager/relocate.go b/pkg/filemanager/manager/relocate.go new file mode 100644 index 00000000..48f3f8d5 --- /dev/null +++ b/pkg/filemanager/manager/relocate.go @@ -0,0 +1,115 @@ +package manager + +import ( + "context" + "fmt" + + "github.com/cloudreve/Cloudreve/v4/ent" + "github.com/cloudreve/Cloudreve/v4/inventory/types" + "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/fs" + "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/manager/entitysource" + "github.com/cloudreve/Cloudreve/v4/pkg/serializer" +) + +// RelocateBlobArgs args to relocate the blob of an entity into another storage policy. +type RelocateBlobArgs struct { + // Entity whose blob is being relocated. + Entity fs.Entity + // Uri of the file owning the entity. Storage drivers fall back to it for MIME type + // detection, so it should be given whenever the owning file is known. + Uri *fs.URI + // DstPolicy is the storage policy the blob is relocated into. + DstPolicy *ent.StoragePolicy + // SavePath is the path of the blob in the destination storage policy. + SavePath string + // ProgressFunc, if given, is called while the blob is being transferred. + ProgressFunc fs.ProgressFunc +} + +// RelocateBlob copies the blob of args.Entity into args.DstPolicy at args.SavePath, applying the +// file encryption setting of the destination policy along the way. +// +// A blob that is already encrypted is transferred in its ciphertext form, keeping the data key it +// was encrypted with, so that relocation costs no crypto and never invalidates the key stored on +// the entity. A plaintext blob relocated into a policy with encryption enabled is encrypted during +// the transfer under a freshly generated data key. +// +// When not nil, the returned metadata describes the encryption of the relocated blob and has to be +// persisted onto the entity once the caller commits the relocation; a nil return means the +// encryption state of the entity is left as it is. The metadata is returned rather than persisted +// here so that it never needs to be carried in a resumable task state, where the plaintext data +// key would end up written into the database. +func (m *manager) RelocateBlob(ctx context.Context, args *RelocateBlobArgs) (*types.EncryptMetadata, error) { + if args == nil || args.Entity == nil || args.DstPolicy == nil { + return nil, fmt.Errorf("invalid relocate blob arguments") + } + + es, err := m.GetEntitySource(ctx, 0, fs.WithEntity(args.Entity)) + if err != nil { + return nil, fmt.Errorf("failed to open source entity: %w", err) + } + defer es.Close() + + var metadata *types.EncryptMetadata + if args.Entity.Encrypted() { + // Present the blob as stored, the data key travels with the entity row. + es.Apply(entitysource.WithDisableCryptor()) + } else if args.DstPolicy.Settings != nil && args.DstPolicy.Settings.Encryption { + cryptor, err := m.dep.EncryptorFactory(ctx)(types.CipherAES256CTR) + if err != nil { + return nil, fmt.Errorf("failed to create cryptor: %w", err) + } + + metadata, err = cryptor.GenerateMetadata(ctx) + if err != nil { + return nil, fmt.Errorf("failed to generate encrypt metadata: %w", err) + } + } + + req := &fs.UploadRequest{ + Props: &fs.UploadProps{ + Uri: args.Uri, + Size: args.Entity.Size(), + SavePath: args.SavePath, + }, + Mode: fs.ModeOverwrite, + File: es, + Seeker: es, + ProgressFunc: args.ProgressFunc, + } + + if args.Uri != nil { + req.Props.MimeType = m.dep.MimeDetector(ctx).TypeByName(args.Uri.Name()) + } + + if metadata != nil { + // The stream is encrypted from its very beginning, hence a zero counter offset. AES-256-CTR + // preserves the length of the stream, so Props.Size stays valid for the ciphertext and the + // size recorded on the entity needs no adjustment. + if err := encryptUploadRequest(ctx, m.dep.EncryptorFactory(ctx), req, metadata, 0); err != nil { + return nil, err + } + } + + d, err := m.GetStorageDriver(ctx, m.CastStoragePolicyOnSlave(ctx, args.DstPolicy)) + if err != nil { + return nil, err + } + + if err := d.Put(ctx, req); err != nil { + return nil, serializer.NewError(serializer.CodeIOFailed, "Failed to relocate blob", err) + } + + if metadata == nil { + return nil, nil + } + + // The plaintext data key is dropped from the returned metadata: it is only needed while the + // stream is being encrypted, whereas what is returned here is meant to be persisted, on the + // entity and possibly in a resumable task state along the way. + return &types.EncryptMetadata{ + Algorithm: metadata.Algorithm, + Key: metadata.Key, + IV: metadata.IV, + }, nil +} diff --git a/pkg/filemanager/manager/upload.go b/pkg/filemanager/manager/upload.go index f42d2665..6962aae9 100644 --- a/pkg/filemanager/manager/upload.go +++ b/pkg/filemanager/manager/upload.go @@ -15,6 +15,7 @@ import ( "github.com/cloudreve/Cloudreve/v4/inventory/types" "github.com/cloudreve/Cloudreve/v4/pkg/cluster" "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/driver" + "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/encrypt" "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/fs" "github.com/cloudreve/Cloudreve/v4/pkg/logging" "github.com/cloudreve/Cloudreve/v4/pkg/queue" @@ -229,29 +230,46 @@ func (m *manager) Upload(ctx context.Context, req *fs.UploadRequest, policy *ent } if session != nil && session.EncryptMetadata != nil && !req.Props.ClientSideEncrypted { - cryptor, err := m.dep.EncryptorFactory(ctx)(session.EncryptMetadata.Algorithm) - if err != nil { - return fmt.Errorf("failed to create cryptor: %w", err) + if err := encryptUploadRequest(ctx, m.dep.EncryptorFactory(ctx), req, session.EncryptMetadata, req.Offset); err != nil { + return err } + } - err = cryptor.LoadMetadata(ctx, session.EncryptMetadata) - if err != nil { - return fmt.Errorf("failed to load encrypt metadata: %w", err) - } + if err := d.Put(ctx, req); err != nil { + return serializer.NewError(serializer.CodeIOFailed, "Failed to upload file", err) + } - if err := cryptor.SetSource(req.File, req.Seeker, req.Props.Size, 0); err != nil { - return fmt.Errorf("failed to set source: %w", err) - } + return nil +} - req.File = cryptor +// encryptUploadRequest wraps the file stream of req with a cryptor described by metadata, so +// that the stream is encrypted on the fly while being written into a storage policy. +// +// counterOffset is the absolute offset of the stream within the blob it belongs to. It has to +// be given so that the keystream stays aligned when only a part of the blob is being written, +// and so that a retried chunk resumes from the right position. +func encryptUploadRequest(ctx context.Context, factory encrypt.CryptorFactory, req *fs.UploadRequest, + metadata *types.EncryptMetadata, counterOffset int64) error { + cryptor, err := factory(metadata.Algorithm) + if err != nil { + return fmt.Errorf("failed to create cryptor: %w", err) + } - if req.Seeker != nil { - req.Seeker = cryptor - } + if err := cryptor.LoadMetadata(ctx, metadata); err != nil { + return fmt.Errorf("failed to load encrypt metadata: %w", err) } - if err := d.Put(ctx, req); err != nil { - return serializer.NewError(serializer.CodeIOFailed, "Failed to upload file", err) + if err := cryptor.SetSource(req.File, req.Seeker, req.Props.Size, counterOffset); err != nil { + return fmt.Errorf("failed to set source: %w", err) + } + + req.File = cryptor + + // The seeker has to be replaced along with the file stream: chunked drivers rewind the + // request on retry, and rewinding the plaintext without rewinding the keystream would + // silently produce corrupted ciphertext. + if req.Seeker != nil { + req.Seeker = cryptor } return nil diff --git a/pkg/filemanager/manager/upload_test.go b/pkg/filemanager/manager/upload_test.go new file mode 100644 index 00000000..3eea9a55 --- /dev/null +++ b/pkg/filemanager/manager/upload_test.go @@ -0,0 +1,153 @@ +package manager + +import ( + "bytes" + "context" + "crypto/rand" + "io" + "strings" + "testing" + + "github.com/cloudreve/Cloudreve/v4/inventory/types" + "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/encrypt" + "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/fs" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +type staticMasterKeyVault struct { + key []byte +} + +func (v staticMasterKeyVault) GetMasterKey(context.Context) ([]byte, error) { + return v.key, nil +} + +type readSeekCloser struct { + *bytes.Reader +} + +func (readSeekCloser) Close() error { + return nil +} + +func testCryptorFactory(t *testing.T) encrypt.CryptorFactory { + t.Helper() + + masterKey := make([]byte, 32) + _, err := rand.Read(masterKey) + require.NoError(t, err) + + return encrypt.NewCryptorFactory(staticMasterKeyVault{key: masterKey}) +} + +// persistedMetadata returns encryption metadata in the shape it is stored on an entity, that is +// with the plaintext data key dropped, so that tests exercise the unwrapping path. +func persistedMetadata(t *testing.T, factory encrypt.CryptorFactory) *types.EncryptMetadata { + t.Helper() + + cryptor, err := factory(types.CipherAES256CTR) + require.NoError(t, err) + + generated, err := cryptor.GenerateMetadata(context.Background()) + require.NoError(t, err) + + return &types.EncryptMetadata{ + Algorithm: generated.Algorithm, + Key: generated.Key, + IV: generated.IV, + } +} + +func decryptBlob(t *testing.T, factory encrypt.CryptorFactory, metadata *types.EncryptMetadata, + ciphertext []byte, counterOffset int64) []byte { + t.Helper() + + cryptor, err := factory(metadata.Algorithm) + require.NoError(t, err) + require.NoError(t, cryptor.LoadMetadata(context.Background(), metadata)) + require.NoError(t, cryptor.SetSource(readSeekCloser{bytes.NewReader(ciphertext)}, nil, + int64(len(ciphertext)), counterOffset)) + + plaintext, err := io.ReadAll(cryptor) + require.NoError(t, err) + + return plaintext +} + +func newTestUploadRequest(plaintext []byte) *fs.UploadRequest { + src := readSeekCloser{bytes.NewReader(plaintext)} + return &fs.UploadRequest{ + Props: &fs.UploadProps{Size: int64(len(plaintext))}, + File: src, + Seeker: src, + } +} + +// A blob encrypted on its way into a storage policy has to decrypt back to the original bytes, +// and has to keep its length: the size recorded on the entity is the plaintext size, and storage +// drivers write exactly that many bytes. +func TestEncryptUploadRequestRoundTripPreservesLength(t *testing.T) { + factory := testCryptorFactory(t) + metadata := persistedMetadata(t, factory) + plaintext := []byte(strings.Repeat("cloudreve relocation payload ", 200)) + + req := newTestUploadRequest(plaintext) + require.NoError(t, encryptUploadRequest(context.Background(), factory, req, metadata, 0)) + + ciphertext, err := io.ReadAll(req) + require.NoError(t, err) + + assert.Len(t, ciphertext, len(plaintext)) + assert.NotEqual(t, plaintext, ciphertext) + assert.Equal(t, plaintext, decryptBlob(t, factory, metadata, ciphertext, 0)) +} + +// Chunked storage drivers rewind the upload request to replay a chunk on retry. Rewinding the +// plaintext without rewinding the keystream would silently produce corrupted ciphertext, and +// nothing downstream computes a checksum that would catch it. +func TestEncryptUploadRequestRewindOnRetryIsStable(t *testing.T) { + factory := testCryptorFactory(t) + metadata := persistedMetadata(t, factory) + plaintext := []byte(strings.Repeat("retried chunk ", 300)) + + req := newTestUploadRequest(plaintext) + require.NoError(t, encryptUploadRequest(context.Background(), factory, req, metadata, 0)) + + // Abandon a partially read chunk, the way a failed request does. + discarded := make([]byte, 64) + _, err := io.ReadFull(req, discarded) + require.NoError(t, err) + + _, err = req.Seek(0, io.SeekStart) + require.NoError(t, err) + + replayed, err := io.ReadAll(req) + require.NoError(t, err) + + assert.Len(t, replayed, len(plaintext)) + assert.Equal(t, plaintext, decryptBlob(t, factory, metadata, replayed, 0)) +} + +// A non-zero counter offset keeps the keystream aligned when only a part of a blob is written, so +// that a chunk written at an offset stays decryptable as part of the whole blob. +func TestEncryptUploadRequestCounterOffsetAlignsKeystream(t *testing.T) { + const offset = 512 + + factory := testCryptorFactory(t) + metadata := persistedMetadata(t, factory) + plaintext := []byte(strings.Repeat("offset aligned payload ", 100)) + require.Greater(t, len(plaintext), offset) + + whole := newTestUploadRequest(plaintext) + require.NoError(t, encryptUploadRequest(context.Background(), factory, whole, metadata, 0)) + full, err := io.ReadAll(whole) + require.NoError(t, err) + + tail := newTestUploadRequest(plaintext[offset:]) + require.NoError(t, encryptUploadRequest(context.Background(), factory, tail, metadata, offset)) + partial, err := io.ReadAll(tail) + require.NoError(t, err) + + assert.Equal(t, full[offset:], partial) +}