diff --git a/docs/en/antalya/cas/bucket-requirements.md b/docs/en/antalya/cas/bucket-requirements.md index 700d0a94d371..ea4bfe0dbafe 100644 --- a/docs/en/antalya/cas/bucket-requirements.md +++ b/docs/en/antalya/cas/bucket-requirements.md @@ -23,6 +23,7 @@ these conditions is refused rather than trusted. | Unconditional complete-object publication | `Backend::publishBlob` | An absent or condemned content-addressed body is replaced atomically; native stores may use multipart | | Native same-store copy when `cas_staging_backend = s3` | `IObjectStorage::copyObject` with `ObjectStorageCopyMode::NativeOnly` | The first absent staged publication may copy its complete object without a client-side fallback | | Exact-token delete | `Backend::deleteExact` | GC must delete only the incarnation it condemned, never a replacement | +| Ranged same-store copy, for backups only | `UploadPartCopy` through `copyS3File` | A `BACKUP` to an `S3(...)` destination on the pool's own endpoint copies a blob's payload without its header. Optional, but only when the lack is recognized: `GCS` lacks it by its current XML API, so `Client::supportsMultiPartCopy` reports `false` there and the copy goes through the ClickHouse server, which is slower but correct. A store that instead refuses `UploadPartCopy` with an error other than `AccessDenied` fails the `BACKUP`; set `s3_allow_multipart_copy = 0` there | | Ranged `GET` | `Backend::get` / `Backend::getStream` with a `Range` | Opening one column file of a part costs one bounded read, not a whole-object fetch | | `LIST` with a resumable cursor | `Backend::list` | GC discovery and the orphan-manifest sweep page through the pool without a separate index | | No versioning / no delete markers | probed by `runCapabilityProbe`; `created_delete_marker` on `DeleteOutcome` | A delete marker over a live key would break exact-token semantics — GC would archive instead of reclaim | diff --git a/docs/en/antalya/cas/index.md b/docs/en/antalya/cas/index.md index cb1d563eebf0..377f03b7d40f 100644 --- a/docs/en/antalya/cas/index.md +++ b/docs/en/antalya/cas/index.md @@ -94,4 +94,5 @@ per disk, so adopting it never requires migrating an existing deployment. | [Architecture overview](/antalya/cas/architecture/) | The object model, the Git analogy, and the safety invariants | | [Correctness](/antalya/cas/architecture/correctness) | How the design was verified: TLA+ models, counterexamples, soak methodology | | [Design history](/antalya/cas/architecture/design-history) | What earlier designs were tried and rejected, and why | +| [Backup](/antalya/cas/operations/backup) | How `BACKUP` and `RESTORE` behave on a content-addressed disk | | [Roadmap](/antalya/cas/roadmap) | What is shipped, planned, and deliberately not pursued | diff --git a/docs/en/antalya/cas/operations/backup.md b/docs/en/antalya/cas/operations/backup.md new file mode 100644 index 000000000000..0a8d64a48840 --- /dev/null +++ b/docs/en/antalya/cas/operations/backup.md @@ -0,0 +1,150 @@ +--- +description: 'How BACKUP and RESTORE work for a table on a content-addressed disk: what holds the data during a backup, when the copy runs inside S3, and what is not supported yet.' +sidebar_label: 'Backup' +sidebar_position: 5 +slug: /antalya/cas/operations/backup +title: 'CAS Operations — Backup' +doc_type: 'guide' +--- + +# Operations — backup {#backup} + +Ordinary `BACKUP` and `RESTORE` work for a table on a content-addressed (`CAS`) disk. This page +covers what happens during one, how it differs from a plain disk, and what the limits are. + +The `CAS`-native backup model — `snapshot` / `mirror` / `fetch` / `restore` — is designed but not +implemented and is not wired into the SQL surface. See the [roadmap](/antalya/cas/roadmap#backups). + +## What is supported {#supported} + +```sql +BACKUP TABLE t TO S3('https://bucket.s3.amazonaws.com/backups/b1', 'key', 'secret'); +RESTORE TABLE t AS t_restored FROM S3('https://bucket.s3.amazonaws.com/backups/b1', 'key', 'secret'); +``` + +The destination can be anything: `S3`, `Disk`, `File`, or an archive. A backup can be restored onto a +disk of any type, because it holds the table's files rather than the pool's objects. + +**An `Atomic` database is required.** That has been the default since 20.x. On the deprecated +`Ordinary` engine the backup fails with `SUPPORT_IS_DISABLED`: that path pins files with temporary +hard links, which object storage does not have. + +## What holds the data during a backup {#holding} + +On a plain disk a backup pins files against deletion with a hard link. `CAS` uses a different +mechanism — pointer holding: + +- the backup holds a `shared_ptr` to the table and to each part; +- the outdated-part cleanup skips those parts; +- while a part is alive so is its [ref](/antalya/cas/architecture/manifests-and-refs#ref-table) — the + name under which the part is registered in its namespace's ref table, and through which it points + at its manifest; +- while the ref is alive, garbage collection sees the manifest and its blobs as reachable. + +This mechanism lives in the process's memory and **does not survive a server restart**. An +interrupted backup leaves nothing behind in the pool, but it also stops protecting the data once the +process is gone. For durable pinning there is `FREEZE`, which publishes a real ref. + +## How the bytes move {#copy-path} + +Which path runs depends on the destination. + +**A destination outside the pool** — the common case: another bucket, a local disk, an archive. Files +are read through the `CAS` read path and written to the destination. Pool deduplication is lost: +what was one blob shared by several replicas becomes ordinary files in the backup. + +**An `S3(...)` destination on the same `S3` endpoint as the pool** — the copy of a blob then runs +inside the `S3` store itself, without sending the bytes through ClickHouse. The server issues one `UploadPartCopy` per upload part, so a payload larger than one part takes several. + +```sql +BACKUP TABLE t TO S3('http://s3.example.com/bucket/backups/b1', 'key', 'secret'); +``` + +Here the pool is under `http://s3.example.com/bucket/pool/`. + +The files of a part fall into two categories: + +| Category | Example | How it is copied | +|---|---|---| +| Blob | `data.bin`, marks, `primary.idx` | an `UploadPartCopy` naming a byte range — only the payload moves, without the blob's internal header | +| Inside the manifest | `checksums.txt`, `count.txt`, `columns.txt` | read through the `CAS` read path and written through ClickHouse's buffers | + +A blob object is `[header][payload]`, so a file never starts at the beginning of its object. Every +copy of a blob therefore names the payload range. A copy of the whole object would put the header +into the backup, and a later restore would read wrong data. + +If a byte range cannot be copied inside `S3`, `CAS` does not fall back to copying the whole object — +the file is read and written through ClickHouse instead. That is slower, but correct. ClickHouse takes +that path on its own in three cases: the query sets `s3_allow_multipart_copy = 0`; the store is +recognized as `GCS`, whose current XML API has no `UploadPartCopy`; or `UploadPartCopy` is refused with +`AccessDenied`. + +A store that refuses `UploadPartCopy` with any other error fails the `BACKUP` instead, because +ClickHouse cannot tell a missing capability from a transient failure. On such a store, set +`s3_allow_multipart_copy = 0` on the `BACKUP` query. + +| Turn off the copy inside `S3` | Turn off the range copy | +|---|---| +| `SETTINGS allow_s3_native_copy = 0` in the `BACKUP` query | the query setting `s3_allow_multipart_copy = 0` | + +```sql +BACKUP TABLE t TO S3(...) SETTINGS allow_s3_native_copy = 0; +``` + +**A `Disk(...)` destination** never copies inside `S3`, even when the disk is an `s3` or `s3_plain` +disk on the same endpoint as the pool: + +```sql +BACKUP TABLE t TO Disk('backups_s3', 'b1'); +``` + +Every file is read through the `CAS` read path and written through ClickHouse's buffers. The +`allow_s3_native_copy` and `s3_allow_multipart_copy` settings have no effect on this path. + +## Restore {#restore} + +Each part is materialized in **one disk transaction** and published as one manifest and one ref. A +partially restored part can never appear in the pool: either the whole part is published or nothing +is. + +Restore onto a `CAS` disk never copies objects inside `S3`, even when the backup is on the same +endpoint as the pool. Each file of the part is read from the backup and written through the `CAS` +write path, because only that path can build the manifest and the blobs. A restore onto a plain disk +works as usual and can copy inside `S3`. + +Restored data is packed afresh — on a `CAS` disk it gets new blobs and new refs. Deduplication +against data already in the pool works as usual: identical content hashes to the same blob and is +not written twice. + +## `FREEZE` is not a backup {#freeze} + +`ALTER TABLE ... FREEZE` works on `CAS` and publishes parts into a separate shadow namespace, which +is a garbage-collection root in its own right. `DROP PARTITION` removes the live refs and leaves the +snapshot alone; `SYSTEM UNFREEZE` removes only the shadow refs. + +It is still not a snapshot of a table: there is no SQL metadata, no single commit marker, no +portable object with a listing and a restore API, and its lifetime is tied to a manual `UNFREEZE`. +It is a useful building block, not a replacement for `BACKUP`. + +## Limitations {#limitations} + +- The `CAS`-native backup model (`snapshot` / `mirror` / `fetch`) is not implemented. +- The `Ordinary` database engine is not supported. +- Pool deduplication is lost in the backup: its size follows the logical files, not the unique blobs. +- Pointer holding does not survive a server restart. +- The copy inside the store works only for an `S3(...)` destination on an `S3` or `S3`-compatible + store, and only when the destination is on the same endpoint as the pool. Other destinations, + including `Disk(...)`, get the copy through ClickHouse's buffers. +- A blob is copied inside `S3` only with multipart copy (`UploadPartCopy`), because only it can name a + byte range. With `s3_allow_multipart_copy = 0`, on a store recognized as `GCS`, and when + `UploadPartCopy` is refused with `AccessDenied`, every blob goes through ClickHouse's buffers + instead. A store that refuses `UploadPartCopy` with another error fails the `BACKUP`; set + `s3_allow_multipart_copy = 0` there. +- Restore onto a `CAS` disk always writes through ClickHouse, see [restore](#restore). +- A disk-level write of a part file onto a `CAS` disk outside of a part transaction is rejected with + `NOT_IMPLEMENTED` and the message `Autocommit writes are not supported for content part files`: a + `CAS` disk publishes a part's manifest and ref at commit, so it has nowhere to put a single + autocommitted file. +- A `CAS` disk cannot be a backup destination: `BACKUP TABLE t TO Disk('', 'b1')` is rejected + the same way. A backup's own layout mirrors the table's data directory, so its files sit under a part + directory as well and count as part files. diff --git a/docs/en/antalya/cas/roadmap.md b/docs/en/antalya/cas/roadmap.md index 4df03af9ab3c..14578bd42058 100644 --- a/docs/en/antalya/cas/roadmap.md +++ b/docs/en/antalya/cas/roadmap.md @@ -91,6 +91,10 @@ positioning. ## Backups {#backups} +Ordinary `BACKUP` and `RESTORE` already work for a table on a `CAS` disk — see +[backup](/antalya/cas/operations/backup) for how they behave and what the limits are. What follows is +about the `CAS`-native model, which is a different thing. + A `snapshot` / `mirror` / `fetch` / `restore` design is **approved but not implemented**. The model is deliberately git-shaped: `snapshot` is instant and free (like `git tag` — it references existing manifests, copies nothing); `mirror` is a continuous pull from a production pool into a diff --git a/src/Backups/BackupIO_AzureBlobStorage.cpp b/src/Backups/BackupIO_AzureBlobStorage.cpp index cdd7b3ca326a..c3ceb2682c84 100644 --- a/src/Backups/BackupIO_AzureBlobStorage.cpp +++ b/src/Backups/BackupIO_AzureBlobStorage.cpp @@ -37,7 +37,15 @@ BackupReaderAzureBlobStorage::BackupReaderAzureBlobStorage( const WriteSettings & write_settings_, const ContextPtr & context_) : BackupReaderDefault(read_settings_, write_settings_, getLogger("BackupReaderAzureBlobStorage")) - , data_source_description{DataSourceType::ObjectStorage, ObjectStorageType::Azure, MetadataStorageType::None, connection_params_.getConnectionURL(), false, false, ""} + , data_source_description{ + .type = DataSourceType::ObjectStorage, + .object_storage_type = ObjectStorageType::Azure, + .metadata_type = MetadataStorageType::None, + .description = connection_params_.getConnectionURL(), + .is_encrypted = false, + .is_cached = false, + .zookeeper_name = "", + .files_are_whole_objects = true} , connection_params(connection_params_) , blob_path(blob_path_) { @@ -87,7 +95,9 @@ void BackupReaderAzureBlobStorage::copyFileToDisk(const String & path_in_backup, auto destination_data_source_description = destination_disk->getDataSourceDescription(); LOG_TRACE(log, "Source description {}, destination description {}", data_source_description.description, destination_data_source_description.description); if (destination_data_source_description.object_storage_type == ObjectStorageType::Azure - && destination_data_source_description.is_encrypted == encrypted_in_backup) + && destination_data_source_description.is_encrypted == encrypted_in_backup + && destination_data_source_description.files_are_whole_objects + && data_source_description.files_are_whole_objects) { LOG_TRACE(log, "Copying {} from AzureBlobStorage to disk {}", path_in_backup, destination_disk->getName()); auto write_blob_function = [&](const Strings & dst_blob_path, WriteMode mode, const std::optional &) -> size_t @@ -133,7 +143,15 @@ BackupWriterAzureBlobStorage::BackupWriterAzureBlobStorage( const ContextPtr & context_, bool attempt_to_create_container) : BackupWriterDefault(read_settings_, write_settings_, getLogger("BackupWriterAzureBlobStorage")) - , data_source_description{DataSourceType::ObjectStorage, ObjectStorageType::Azure, MetadataStorageType::None, connection_params_.getConnectionURL(), false, false, ""} + , data_source_description{ + .type = DataSourceType::ObjectStorage, + .object_storage_type = ObjectStorageType::Azure, + .metadata_type = MetadataStorageType::None, + .description = connection_params_.getConnectionURL(), + .is_encrypted = false, + .is_cached = false, + .zookeeper_name = "", + .files_are_whole_objects = true} , connection_params(connection_params_) , blob_path(blob_path_) { @@ -165,7 +183,9 @@ void BackupWriterAzureBlobStorage::copyFileFromDisk( auto source_data_source_description = src_disk->getDataSourceDescription(); LOG_TRACE(log, "Source description {}, destination description {}", source_data_source_description.description, data_source_description.description); if (source_data_source_description.object_storage_type == ObjectStorageType::Azure - && source_data_source_description.is_encrypted == copy_encrypted) + && source_data_source_description.is_encrypted == copy_encrypted + && source_data_source_description.files_are_whole_objects + && data_source_description.files_are_whole_objects) { /// getBlobPath() can return more than 2 elements if the file is stored as multiple objects in AzureBlobStorage container. /// In this case we can't use the native copy. diff --git a/src/Backups/BackupIO_S3.cpp b/src/Backups/BackupIO_S3.cpp index df43076d8957..e758d6a31467 100644 --- a/src/Backups/BackupIO_S3.cpp +++ b/src/Backups/BackupIO_S3.cpp @@ -15,6 +15,8 @@ #include #include #include +#include +#include #include @@ -256,7 +258,15 @@ BackupReaderS3::BackupReaderS3( bool is_internal_backup) : BackupReaderDefault(read_settings_, write_settings_, getLogger("BackupReaderS3")) , s3_uri(s3_uri_) - , data_source_description{DataSourceType::ObjectStorage, ObjectStorageType::S3, MetadataStorageType::None, s3_uri.endpoint, false, false, ""} + , data_source_description{ + .type = DataSourceType::ObjectStorage, + .object_storage_type = ObjectStorageType::S3, + .metadata_type = MetadataStorageType::None, + .description = s3_uri.endpoint, + .is_encrypted = false, + .is_cached = false, + .zookeeper_name = "", + .files_are_whole_objects = true} { s3_settings.loadFromConfig(context_->getConfigRef(), "s3", context_->getSettingsRef()); @@ -299,7 +309,7 @@ void BackupReaderS3::copyFileToDisk(const String & path_in_backup, size_t file_s /// Use the native copy as a more optimal way to copy a file from S3 to S3 if it's possible. /// We don't check for `has_throttling` here because the native copy almost doesn't use network. auto destination_data_source_description = destination_disk->getDataSourceDescription(); - if (destination_data_source_description.sameKind(data_source_description) + if (destination_data_source_description.canUseNativeCopyWith(data_source_description) && (destination_data_source_description.is_encrypted == encrypted_in_backup)) { LOG_TRACE(log, "Copying {} from S3 to disk {}", path_in_backup, destination_disk->getName()); @@ -353,7 +363,15 @@ BackupWriterS3::BackupWriterS3( bool is_internal_backup) : BackupWriterDefault(read_settings_, write_settings_, getLogger("BackupWriterS3")) , s3_uri(s3_uri_) - , data_source_description{DataSourceType::ObjectStorage, ObjectStorageType::S3, MetadataStorageType::None, s3_uri.endpoint, false, false, ""} + , data_source_description{ + .type = DataSourceType::ObjectStorage, + .object_storage_type = ObjectStorageType::S3, + .metadata_type = MetadataStorageType::None, + .description = s3_uri.endpoint, + .is_encrypted = false, + .is_cached = false, + .zookeeper_name = "", + .files_are_whole_objects = true} , s3_capabilities(getCapabilitiesFromConfig(context_->getConfigRef(), "s3")) , disk_client_factory(S3BackupClientCreator(context_)) { @@ -385,7 +403,20 @@ void BackupWriterS3::copyFileFromDisk( /// Use the native copy as a more optimal way to copy a file from S3 to S3 if it's possible. /// We don't check for `has_throttling` here because the native copy almost doesn't use network. auto source_data_source_description = src_disk->getDataSourceDescription(); - if (source_data_source_description.sameKind(data_source_description) && (source_data_source_description.is_encrypted == copy_encrypted)) + + if (!copy_encrypted && !source_data_source_description.is_encrypted) + { + if (auto * ca = tryGetContentAddressedExchange(src_disk)) + { + if (tryNativeCopyFromContentAddressedDisk(*ca, path_in_backup, src_disk, src_path, start_pos, length)) + return; + + BackupWriterDefault::copyFileFromDisk(path_in_backup, src_disk, src_path, copy_encrypted, start_pos, length); + return; + } + } + + if (source_data_source_description.canUseNativeCopyWith(data_source_description) && (source_data_source_description.is_encrypted == copy_encrypted)) { /// getBlobPath() can return more than 2 elements if the file is stored as multiple objects in S3 bucket. /// In this case we can't use the native copy. @@ -422,6 +453,97 @@ void BackupWriterS3::copyFileFromDisk( BackupWriterDefault::copyFileFromDisk(path_in_backup, src_disk, src_path, copy_encrypted, start_pos, length); } +bool BackupWriterS3::tryNativeCopyFromContentAddressedDisk( + IContentAddressedExchange & ca, + const String & path_in_backup, + DiskPtr src_disk, + const String & src_path, + UInt64 start_pos, + UInt64 length) +{ + if (length == 0) + return false; + + auto source_data_source_description = src_disk->getDataSourceDescription(); + if (!source_data_source_description.sameKind(data_source_description)) + return false; + + auto src_client = disk_client_factory.getOrCreate(src_disk); + if (!src_client->supportsMultiPartCopy()) + return false; + + const auto plan = ca.getBlobViewPlan(src_path); + if (!plan) + return false; + + const String src_bucket = src_disk->getObjectStorage()->getObjectsNamespace(); + if (src_bucket.empty()) + throw Exception( + ErrorCodes::LOGICAL_ERROR, + "Disk {} has the same S3 endpoint as the backup but no bucket for content-addressed file {}", + src_disk->getName(), + src_path); + + if (plan->object.remote_path.empty()) + throw Exception( + ErrorCodes::LOGICAL_ERROR, + "Content-addressed file {} on disk {} resolved to a blob with an empty key", + src_path, + src_disk->getName()); + + if (plan->payload_offset == 0) + throw Exception( + ErrorCodes::LOGICAL_ERROR, + "Content-addressed file {} on disk {} resolved to a blob whose payload starts at byte 0, but a blob always " + "begins with an envelope header", + src_path, + src_disk->getName()); + + const UInt64 payload_size = plan->payload_end - plan->payload_offset; + if (start_pos > payload_size || length > payload_size - start_pos) + throw Exception( + ErrorCodes::LOGICAL_ERROR, + "Requested range [{}, {}) of content-addressed file {} lies outside its payload of {} bytes", + start_pos, + start_pos + length, + src_path, + payload_size); + + const UInt64 src_offset = plan->payload_offset + start_pos; + + LOG_TRACE( + log, + "Copying the payload of content-addressed file {} from disk {} to S3 as a ranged server-side copy", + src_path, + src_disk->getName()); + + copyS3File( + std::move(src_client), + src_bucket, + /* src_key */ plan->object.remote_path, + src_offset, + length, + /* dest_s3_client */ client, + /* dest_bucket */ s3_uri.bucket, + /* dest_key */ fs::path(s3_uri.key) / path_in_backup, + s3_settings.request_settings, + read_settings, + blob_storage_log, + threadPoolCallbackRunnerUnsafe(getBackupsIOThreadPool().get(), ThreadName::S3_BACKUP_WRITER), + [&, this] + { + LOG_TRACE( + log, + "Falling back to copy the raw object of content-addressed file {} from disk {} to S3 through buffers", + src_path, + src_disk->getName()); + + return src_disk->getObjectStorage()->readObject(plan->object, read_settings); + }); + + return true; +} + void BackupWriterS3::copyFile(const String & destination, const String & source, size_t size) { LOG_TRACE(log, "Copying file inside backup from {} to {}", source, destination); diff --git a/src/Backups/BackupIO_S3.h b/src/Backups/BackupIO_S3.h index 78bd9031450c..bdf0fd51584b 100644 --- a/src/Backups/BackupIO_S3.h +++ b/src/Backups/BackupIO_S3.h @@ -18,6 +18,8 @@ namespace DB { +class IContentAddressedExchange; + class S3BackupDiskClientFactory { public: @@ -108,6 +110,14 @@ class BackupWriterS3 : public BackupWriterDefault private: std::unique_ptr readFile(const String & file_name, size_t expected_file_size) override; + bool tryNativeCopyFromContentAddressedDisk( + IContentAddressedExchange & ca, + const String & path_in_backup, + DiskPtr src_disk, + const String & src_path, + UInt64 start_pos, + UInt64 length); + const S3::URI s3_uri; const DataSourceDescription data_source_description; S3Settings s3_settings; diff --git a/src/Disks/DiskLocal.cpp b/src/Disks/DiskLocal.cpp index 472944b64432..fee1adc17a93 100644 --- a/src/Disks/DiskLocal.cpp +++ b/src/Disks/DiskLocal.cpp @@ -658,6 +658,7 @@ DataSourceDescription DiskLocal::getLocalDataSourceDescription(const String & pa res.description = path; res.is_encrypted = false; res.is_cached = false; + res.files_are_whole_objects = true; return res; } diff --git a/src/Disks/DiskObjectStorage/DiskObjectStorage.cpp b/src/Disks/DiskObjectStorage/DiskObjectStorage.cpp index 2c2d071a7fbb..bebb76139b64 100644 --- a/src/Disks/DiskObjectStorage/DiskObjectStorage.cpp +++ b/src/Disks/DiskObjectStorage/DiskObjectStorage.cpp @@ -54,6 +54,7 @@ namespace ErrorCodes { extern const int INCORRECT_DISK_INDEX; extern const int CANNOT_RMDIR; + extern const int LOGICAL_ERROR; } namespace @@ -128,6 +129,7 @@ DiskObjectStorage::DiskObjectStorage( .is_encrypted = false, .is_cached = object_storages->takePointingTo(cluster->getLocalLocation())->supportsCache(), .zookeeper_name = metadata_storage->getZooKeeperName(), + .files_are_whole_objects = !metadata_storage->isContentAddressed(), }; resource_changes_subscription = Context::getGlobalContextInstance()->getWorkloadEntityStoragePtr()->getAllEntitiesAndSubscribe( [this] (const std::vector & events) @@ -297,11 +299,14 @@ void DiskObjectStorage::copyFile( /// NOLINT const std::function & cancellation_hook) { auto component_guard = Coordination::setCurrentComponent("DiskObjectStorage::copyFile"); - if (getDataSourceDescription() == to_disk.getDataSourceDescription()) + const auto source_description = getDataSourceDescription(); + const auto destination_description = to_disk.getDataSourceDescription(); + auto * to_disk_object_storage = dynamic_cast(&to_disk); + if (to_disk_object_storage && source_description == destination_description + && source_description.canUseNativeCopyWith(destination_description)) { /// It may use s3-server-side copy - auto & to_disk_object_storage = dynamic_cast(to_disk); - auto transaction = createObjectStorageTransactionToAnotherDisk(to_disk_object_storage); + auto transaction = createObjectStorageTransactionToAnotherDisk(*to_disk_object_storage); try { transaction->copyFile(from_file_path, to_file_path, read_settings, write_settings); @@ -824,12 +829,16 @@ void DiskObjectStorage::prepareRead( if (metadata_storage->isContentAddressed()) { const auto * ca = dynamic_cast(metadata_storage.get()); - if (ca) - { - if (ca->prepareInManifestRead(path, settings, pipeline)) - return; - ca_blob_view = ca->getBlobViewPlan(path); - } + if (!ca) + throw Exception( + ErrorCodes::LOGICAL_ERROR, + "Metadata storage of disk {} reports itself content-addressed but does not implement " + "IContentAddressedExchange, so the payload window of {} cannot be resolved", + getName(), path); + + if (ca->prepareInManifestRead(path, settings, pipeline)) + return; + ca_blob_view = ca->getBlobViewPlan(path); } const auto storage_objects = ca_blob_view diff --git a/src/Disks/DiskObjectStorage/DiskObjectStorage.h b/src/Disks/DiskObjectStorage/DiskObjectStorage.h index 09ec478b02d0..4d4aa6d3ad02 100644 --- a/src/Disks/DiskObjectStorage/DiskObjectStorage.h +++ b/src/Disks/DiskObjectStorage/DiskObjectStorage.h @@ -179,6 +179,7 @@ friend class DiskObjectStorageReservation; const WriteSettings & settings) override; Strings getBlobPath(const String & path) const override; + bool areBlobPathsRandom() const override; void writeFileUsingBlobWritingFunction(const String & path, WriteMode mode, WriteBlobFunction && write_blob_function) override; diff --git a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedExchange.cpp b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedExchange.cpp index 513b0ae0a753..b6fa4e7fd5f7 100644 --- a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedExchange.cpp +++ b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedExchange.cpp @@ -1,5 +1,8 @@ #include +#include +#include + #include namespace DB @@ -167,4 +170,11 @@ std::optional decodeCasRelinkSourceToken(std::string_view return token; } +IContentAddressedExchange * tryGetContentAddressedExchange(const DiskPtr & disk) +{ + if (!disk || !disk->isContentAddressed()) + return nullptr; + return dynamic_cast(disk->getMetadataStorage().get()); +} + } diff --git a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedExchange.h b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedExchange.h index ac05a13bc78d..de714d3ede31 100644 --- a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedExchange.h +++ b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedExchange.h @@ -9,6 +9,9 @@ namespace DB { +class IDisk; +using DiskPtr = std::shared_ptr; + class ReadPipeline; struct ReadSettings; @@ -257,4 +260,6 @@ class IContentAddressedExchange virtual std::optional getBlobViewPlan(const std::string & path) const = 0; }; +IContentAddressedExchange * tryGetContentAddressedExchange(const DiskPtr & disk); + } diff --git a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedMetadataStorage.cpp b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedMetadataStorage.cpp index b20d299df60c..c409cd318f59 100644 --- a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedMetadataStorage.cpp +++ b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedMetadataStorage.cpp @@ -2146,17 +2146,15 @@ std::optional ContentAddressedMet /// `partAccess()` then `store()` pair (each an independent `pointer_mutex` acquisition) -- see /// `poolAccess()`. const auto snap = poolAccess(); - auto view = snap.part_access->getView(r->refKey(), Cas::Freshness::CachedForLoad); - if (!view) + const auto manifest_view = snap.part_access->getView(r->refKey(), Cas::Freshness::CachedForLoad); + if (!manifest_view) return std::nullopt; - if (const auto * entry = view->findFile(r->file)) + if (const auto * entry = manifest_view->findFile(r->file)) { + if (entry->placement != Cas::EntryPlacement::Blob) + return std::nullopt; const auto location = snap.pool->locate(*entry); BlobViewPlan plan; - /// bytes_size is the readable extent of THIS file's window, NOT the whole blob: a - /// right-bounded read stops at payload_end, and a shared blob's bytes beyond it belong - /// to other files. The caches key on the physical blob key, so payload ranges are - /// shared between every part that references the same blob. plan.object = StoredObject(physicalKey(location.key), path, location.offset + location.length); plan.payload_offset = location.offset; plan.payload_end = location.offset + location.length; diff --git a/src/Disks/DiskType.cpp b/src/Disks/DiskType.cpp index 6af38a7fa6b8..177e704b275c 100644 --- a/src/Disks/DiskType.cpp +++ b/src/Disks/DiskType.cpp @@ -51,6 +51,11 @@ bool DataSourceDescription::sameKind(const DataSourceDescription & other) const == std::tie(other.type, other.object_storage_type, other_description); } +bool DataSourceDescription::canUseNativeCopyWith(const DataSourceDescription & other) const +{ + return files_are_whole_objects && other.files_are_whole_objects && sameKind(other); +} + String DataSourceDescription::name() const { switch (type) @@ -66,8 +71,9 @@ String DataSourceDescription::name() const String DataSourceDescription::toString() const { - return fmt::format("{} (description = '{}', is_encrypted = {}, is_cached = {}, zookeeper_name = '{}')", - name(), description, is_encrypted, is_cached, zookeeper_name); + return fmt::format( + "{} (description = '{}', is_encrypted = {}, is_cached = {}, zookeeper_name = '{}', files_are_whole_objects = {})", + name(), description, is_encrypted, is_cached, zookeeper_name, files_are_whole_objects); } ObjectStorageType objectStorageTypeFromString(const std::string & type) diff --git a/src/Disks/DiskType.h b/src/Disks/DiskType.h index 003395b3a9d4..ee39b23c7029 100644 --- a/src/Disks/DiskType.h +++ b/src/Disks/DiskType.h @@ -55,8 +55,11 @@ struct DataSourceDescription String zookeeper_name; + bool files_are_whole_objects = false; + bool operator==(const DataSourceDescription & other) const; bool sameKind(const DataSourceDescription & other) const; + bool canUseNativeCopyWith(const DataSourceDescription & other) const; String name() const; diff --git a/src/Disks/ReadOnlyDiskWrapper.h b/src/Disks/ReadOnlyDiskWrapper.h index 9a38e85cde77..9d812cb7da1c 100644 --- a/src/Disks/ReadOnlyDiskWrapper.h +++ b/src/Disks/ReadOnlyDiskWrapper.h @@ -1,5 +1,7 @@ #pragma once +#include "config.h" + #include #include @@ -95,6 +97,11 @@ class ReadOnlyDiskWrapper : public IDisk /// drops out of the CAS introspection paths. bool isContentAddressed() const override { return delegate->isContentAddressed(); } +#if USE_AWS_S3 + std::shared_ptr getS3StorageClient() const override { return delegate->getS3StorageClient(); } + std::shared_ptr tryGetS3StorageClient() const override { return delegate->tryGetS3StorageClient(); } +#endif + std::unordered_map getSerializedMetadata(const std::vector & file_paths) const override { return delegate->getSerializedMetadata(file_paths); } UInt32 getRefCount(const String & path) const override { return delegate->getRefCount(path); } diff --git a/src/Disks/tests/gtest_ca_wiring.cpp b/src/Disks/tests/gtest_ca_wiring.cpp index db32c32b43b5..098be83ff5ce 100644 --- a/src/Disks/tests/gtest_ca_wiring.cpp +++ b/src/Disks/tests/gtest_ca_wiring.cpp @@ -534,13 +534,7 @@ TEST(CASWiringRead, BlobViewPlanRidesTheStandardPipeline) DB::readStringUntilEOF(manifest_bytes, *buf); } EXPECT_EQ(manifest_bytes, "u-123"); - /// Not a `getBlobViewPlan` call on the in-manifest path here (all-tree Task 6/9: uuid.txt is now - /// a real Inline manifest entry): `getBlobViewPlan`'s only production caller - /// (`DiskObjectStorage::prepareRead`) never reaches it once `prepareInManifestRead` returns true - /// above — `getBlobViewPlan`'s precondition is "confirmed not in-manifest-servable," which calling - /// it directly on an Inline path violates. Pre-Task-9 this assertion passed only by coincidence - /// (uuid.txt was not a manifest entry at all, so `findFile` returned not-found, not because - /// `getBlobViewPlan` gracefully handles an Inline entry it does find). + EXPECT_FALSE(storage->getBlobViewPlan("a11/a11a11a1-1111-4111-8111-111111111111/all_1_1_0/uuid.txt").has_value()); /// Blob-backed file: a real physical key and a payload-sized window whose extent equals /// the object's readable size (a right-bounded read never overshoots the window). diff --git a/src/Disks/tests/gtest_data_source_description.cpp b/src/Disks/tests/gtest_data_source_description.cpp new file mode 100644 index 000000000000..8e14cdd4e50c --- /dev/null +++ b/src/Disks/tests/gtest_data_source_description.cpp @@ -0,0 +1,48 @@ +#include + +#include + +using namespace DB; + +namespace +{ + +DataSourceDescription makeS3Description(bool files_are_whole_objects) +{ + DataSourceDescription description; + description.type = DataSourceType::ObjectStorage; + description.object_storage_type = ObjectStorageType::S3; + description.description = "http://storage.example.com/bucket/"; + description.files_are_whole_objects = files_are_whole_objects; + return description; +} + +} + +TEST(DataSourceDescription, NativeCopyNeedsWholeObjectsOnBothSides) +{ + const auto whole = makeS3Description(true); + const auto windowed = makeS3Description(false); + + EXPECT_TRUE(whole.canUseNativeCopyWith(whole)); + EXPECT_FALSE(whole.canUseNativeCopyWith(windowed)); + EXPECT_FALSE(windowed.canUseNativeCopyWith(whole)); + EXPECT_FALSE(windowed.canUseNativeCopyWith(windowed)); +} + +TEST(DataSourceDescription, EqualityStaysReflexiveForWindowedFiles) +{ + const auto windowed = makeS3Description(false); + + EXPECT_TRUE(windowed == windowed); + EXPECT_TRUE(windowed.sameKind(windowed)); +} + +TEST(DataSourceDescription, NativeCopyStillNeedsTheSameKind) +{ + auto one = makeS3Description(true); + auto other = makeS3Description(true); + other.description = "http://other.example.com/bucket/"; + + EXPECT_FALSE(one.canUseNativeCopyWith(other)); +} diff --git a/src/IO/S3/copyS3File.cpp b/src/IO/S3/copyS3File.cpp index 4b1f5e14ece5..4f9e8b0a9743 100644 --- a/src/IO/S3/copyS3File.cpp +++ b/src/IO/S3/copyS3File.cpp @@ -640,8 +640,32 @@ namespace void performCopy() { LOG_TEST(log, "Copy object {} to {} using native copy", src_key, dest_key); - bool use_single_operation_copy = !supports_multipart_copy || !request_settings[S3RequestSetting::allow_multipart_copy] - || (size <= request_settings[S3RequestSetting::max_single_operation_copy_size]); + + const bool multipart_copy_available + = supports_multipart_copy && request_settings[S3RequestSetting::allow_multipart_copy]; + + if (offset != 0 && !multipart_copy_available) + { + if (!allow_fallback) + throw Exception( + ErrorCodes::NOT_IMPLEMENTED, + "Cannot copy a byte range of object {} server-side: only UploadPartCopy can express a " + "range, and multipart copy is unavailable", + src_key); + + LOG_DEBUG( + log, + "Multipart copy is unavailable, so the byte range [{}, {}) of {} cannot be copied " + "server-side, will copy through the server instead", + offset, + offset + size, + src_key); + fallback_method(); + return; + } + + const bool use_single_operation_copy = offset == 0 + && (!multipart_copy_available || (size <= request_settings[S3RequestSetting::max_single_operation_copy_size])); if (use_single_operation_copy) performSingleOperationCopy(); diff --git a/src/Storages/MergeTree/DataPartsExchange.cpp b/src/Storages/MergeTree/DataPartsExchange.cpp index 5d7e3ba6ac93..be021038579c 100644 --- a/src/Storages/MergeTree/DataPartsExchange.cpp +++ b/src/Storages/MergeTree/DataPartsExchange.cpp @@ -158,17 +158,6 @@ constexpr auto CA_CONFIRM_ANSWER_PROVEN = "yes"; /// same safe outcome as a refusal. constexpr auto CA_CONFIRM_ANSWER_UNPROVEN = "unproven"; -/// Resolve a disk to the content-addressed exchange facade, or nullptr if the disk is not CA. The -/// cast targets the purpose-built INTERFACE (IContentAddressedExchange), never the concrete -/// metadata-storage class. Used by both the relink sender (the part's -/// disk) and the relink receiver (the target disk). -IContentAddressedExchange * tryGetContentAddressedExchange(const DiskPtr & disk) -{ - if (!disk || !disk->isContentAddressed()) - return nullptr; - return dynamic_cast(disk->getMetadataStorage().get()); -} - /// Simple functor for tracking fetch progress in system.replicated_fetches table. struct ReplicatedFetchReadCallback { diff --git a/tests/integration/test_cas_backup_s3_native_copy/__init__.py b/tests/integration/test_cas_backup_s3_native_copy/__init__.py new file mode 100644 index 000000000000..e69de29bb2d1 diff --git a/tests/integration/test_cas_backup_s3_native_copy/configs/storage_conf.xml b/tests/integration/test_cas_backup_s3_native_copy/configs/storage_conf.xml new file mode 100644 index 000000000000..44d97e47d5fe --- /dev/null +++ b/tests/integration/test_cas_backup_s3_native_copy/configs/storage_conf.xml @@ -0,0 +1,110 @@ + + + + + object_storage + s3 + cas + 30 + 10000 + itest-cas-backup-s3-native-copy + + http://rustfs1:11121/test/cas_backup_data/ + clickhouse + clickhouse + + + s3_plain + http://rustfs1:11121/test/backup_disk_s3_plain/ + clickhouse + clickhouse + + + object_storage + s3 + local + http://rustfs1:11121/test/backup_disk_s3/ + clickhouse + clickhouse + + + cache + disk_cas_backup_s3 + cas_cache/ + 1073741824 + + + object_storage + s3 + local + http://rustfs1:11121/test/plain_second/ + clickhouse + clickhouse + + + encrypted + backup_disk_s3 + encrypted/ + 1234567812345678 + + + + + +
+ disk_cas_backup_s3 +
+
+
+ + +
+ disk_cas_backup_s3 +
+ + backup_disk_s3 + +
+
+ + +
+ disk_cas_cached +
+ + backup_disk_s3 + +
+
+ + +
+ backup_disk_s3 +
+ + disk_plain_s3_second + +
+
+ + +
+ backup_disk_s3 +
+ + disk_plain_s3_encrypted + +
+
+
+
+ + backup_disk_s3_plain + backup_disk_s3 + disk_plain_s3_encrypted + disk_cas_backup_s3 + +
diff --git a/tests/integration/test_cas_backup_s3_native_copy/test.py b/tests/integration/test_cas_backup_s3_native_copy/test.py new file mode 100644 index 000000000000..7c9a02d179b9 --- /dev/null +++ b/tests/integration/test_cas_backup_s3_native_copy/test.py @@ -0,0 +1,709 @@ +"""BACKUP/RESTORE of a CAS table to an S3 destination sharing the pool's authority. + +`sameKind` then matches and `BackupWriterS3` copies objects server-side instead of reading +through the CAS read path. That path must handle two shapes: a blob, whose object is +`[envelope][payload]` so the file starts at a non-zero offset, and an inline manifest entry, +which has no object at all. `getStorageObjects` reports neither -- it drops the offset and +returns an empty key. + +Columns are chosen so one table yields both shapes: placement keys on the file name, so every +column makes blobs while per-part metadata stays inline. `n` and `arr` add the `.null.bin` and +`.size0.bin` substreams. +""" + +import uuid + +import pytest + +from helpers.client import QueryRuntimeException +from helpers.cluster import ClickHouseCluster + +cluster = ClickHouseCluster(__file__) + +STORAGE_POLICY = "cas_backup_s3" + +S3_AUTHORITY = "http://rustfs1:11121" +S3_CREDENTIALS = "'clickhouse', 'clickhouse'" + +NUM_ROWS = 100000 + +COLUMNS = ["k", "s", "n", "arr"] + +RUN_TOKEN = uuid.uuid4().hex + + +@pytest.fixture(scope="module", autouse=True) +def start_cluster(): + cluster.add_instance( + "node", + main_configs=["configs/storage_conf.xml"], + with_rustfs=True, + with_remote_database_disk=False, + stay_alive=True, + ) + try: + cluster.start() + yield + finally: + cluster.shutdown() + + +def backup_s3_destination(name): + return f"S3('{S3_AUTHORITY}/test/backups/{RUN_TOKEN}/{name}', {S3_CREDENTIALS})" + + +def backup_disk_destination(disk, name): + return f"Disk('{disk}', '{RUN_TOKEN}/{name}')" + + +def create_and_fill(node, table, storage_policy=STORAGE_POLICY): + node.query(f"DROP TABLE IF EXISTS {table} SYNC") + node.query( + f""" + CREATE TABLE {table} (k UInt64, s String, n Nullable(Int64), arr Array(UInt32)) + ENGINE = MergeTree ORDER BY k + SETTINGS storage_policy = '{storage_policy}', min_bytes_for_wide_part = 0 + """ + ) + node.query( + f""" + INSERT INTO {table} + SELECT + number, + randomPrintableASCII(64), + if(number % 7 = 0, NULL, toInt64(number)), + [toUInt32(number), toUInt32(number + 1)] + FROM numbers({NUM_ROWS}) + """ + ) + + +def column_fingerprints(node, table): + """Order-independent per-column hash; reading every column fetches every blob.""" + exprs = ", ".join( + f"sum(cityHash64(ifNull(toString({column}), '')))" for column in COLUMNS + ) + row = node.query(f"SELECT count(), {exprs} FROM {table}").strip().split("\t") + return dict(zip(["count"] + COLUMNS, row)) + + +def transfer_events(node, query_id): + """How many S3 operations of each mechanism the query issued. + + `server_side_part_copy` is `UploadPartCopy`, the only copy that can name a byte range. + `server_side_object_copy` is `CopyObject`, which always takes the whole object. + `buffered_write` and `buffered_read` are the puts and gets a copy through the server pays. + """ + node.query("SYSTEM FLUSH LOGS query_log") + row = node.query( + f""" + SELECT + ProfileEvents['S3UploadPartCopy'], + ProfileEvents['S3CopyObject'], + ProfileEvents['S3PutObject'] + ProfileEvents['S3UploadPart'], + ProfileEvents['S3GetObject'] + FROM system.query_log + WHERE type = 'QueryFinish' AND query_id = '{query_id}' + ORDER BY event_time DESC LIMIT 1 + """ + ).strip() + assert row, f"no query_log row for {query_id}" + names = ["server_side_part_copy", "server_side_object_copy", "buffered_write", "buffered_read"] + return dict(zip(names, (int(value) for value in row.split("\t")))) + + +@pytest.mark.parametrize("allow_native_copy", [True, False]) +def test_backup_to_s3_round_trip(allow_native_copy): + """A blob is `[envelope][payload]`, so its copy must be ranged: `UploadPartCopy`, never + `CopyObject`. Inline entries have no object and go through buffers either way, and a restore onto + a CAS disk writes every file through the CAS write path. + """ + node = cluster.instances["node"] + suffix = "native" if allow_native_copy else "buffered" + table = f"cas_backup_{suffix}" + restored = f"{table}_restored" + s3_destination = backup_s3_destination(suffix) + backup_query_id = f"{table}_backup_{RUN_TOKEN}" + restore_query_id = f"{table}_restore_{RUN_TOKEN}" + + create_and_fill(node, table) + expected = column_fingerprints(node, table) + + node.query( + f"BACKUP TABLE {table} TO {s3_destination} " + f"SETTINGS allow_s3_native_copy = {int(allow_native_copy)}", + query_id=backup_query_id, + ) + + backup = transfer_events(node, backup_query_id) + assert backup["server_side_object_copy"] == 0, ( + "CopyObject cannot express a range, so the envelope would land in the backup" + ) + assert backup["buffered_write"] > 0, "inline entries have no object and go through buffers" + if allow_native_copy: + assert backup["server_side_part_copy"] > 0, "blobs must be copied server-side with a range" + else: + assert backup["server_side_part_copy"] == 0, ( + "allow_s3_native_copy = 0 leaves no copy inside S3, so every file goes through buffers" + ) + + node.query(f"DROP TABLE IF EXISTS {restored} SYNC") + node.query( + f"RESTORE TABLE {table} AS {restored} FROM {s3_destination} " + f"SETTINGS allow_s3_native_copy = {int(allow_native_copy)}", + query_id=restore_query_id, + ) + + restore = transfer_events(node, restore_query_id) + assert (restore["server_side_part_copy"], restore["server_side_object_copy"]) == (0, 0), ( + "a restore onto a CAS disk writes every file through the CAS write path" + ) + assert restore["buffered_read"] > 0, "the backup must be read through buffers" + assert restore["buffered_write"] > 0, "the restored part must be written through buffers" + + actual = column_fingerprints(node, restored) + + assert actual["count"] == expected["count"] + differing = [c for c in COLUMNS if actual[c] != expected[c]] + assert not differing, f"columns differ after restore: {differing}" + + assert ( + node.query( + f"CHECK TABLE {restored} SETTINGS check_query_single_value_result = 1" + ).strip() + == "1" + ) + + node.query(f"DROP TABLE {table} SYNC") + node.query(f"DROP TABLE {restored} SYNC") + + +@pytest.mark.parametrize("backup_disk", ["backup_disk_s3_plain", "backup_disk_s3"]) +def test_backup_to_disk_only_buffers(backup_disk): + node = cluster.instances["node"] + table = f"cas_backup_to_{backup_disk}" + restored = f"{table}_restored" + disk_destination = backup_disk_destination(backup_disk, table) + backup_query_id = f"{table}_backup_{RUN_TOKEN}" + restore_query_id = f"{table}_restore_{RUN_TOKEN}" + + create_and_fill(node, table) + expected = column_fingerprints(node, table) + + node.query(f"BACKUP TABLE {table} TO {disk_destination}", query_id=backup_query_id) + + backup = transfer_events(node, backup_query_id) + assert (backup["server_side_part_copy"], backup["server_side_object_copy"]) == (0, 0), ( + "a Disk(...) destination goes through IDisk::copyFile, which no longer copies CAS objects" + ) + assert backup["buffered_write"] > 0, "every file must reach the destination through buffers" + + node.query(f"DROP TABLE IF EXISTS {restored} SYNC") + node.query( + f"RESTORE TABLE {table} AS {restored} FROM {disk_destination}", + query_id=restore_query_id, + ) + + restore = transfer_events(node, restore_query_id) + assert (restore["server_side_part_copy"], restore["server_side_object_copy"]) == (0, 0), ( + "RESTORE onto a CAS disk must write through the CAS path, not copy objects into the pool" + ) + assert restore["buffered_write"] > 0, "the restored part must be written through buffers" + + assert ( + node.query( + f"SELECT storage_policy FROM system.tables WHERE name = '{restored}'" + ).strip() + == STORAGE_POLICY + ) + + actual = column_fingerprints(node, restored) + + assert actual["count"] == expected["count"] + differing = [c for c in COLUMNS if actual[c] != expected[c]] + assert not differing, f"columns differ after restore: {differing}" + + assert ( + node.query( + f"CHECK TABLE {restored} SETTINGS check_query_single_value_result = 1" + ).strip() + == "1" + ) + + node.query(f"DROP TABLE {table} SYNC") + node.query(f"DROP TABLE {restored} SYNC") + + +def test_backup_to_s3_falls_back_without_multipart(): + """Only `UploadPartCopy` can express a range. With multipart copy off there is no server-side + operation left, so the copy must go through buffers instead of taking the whole object. + """ + node = cluster.instances["node"] + table = "cas_backup_no_multipart" + restored = f"{table}_restored" + s3_destination = backup_s3_destination("no_multipart") + query_id = f"cas_backup_no_multipart_{RUN_TOKEN}" + no_multipart = {"s3_allow_multipart_copy": 0} + + create_and_fill(node, table) + expected = column_fingerprints(node, table) + + node.query( + f"BACKUP TABLE {table} TO {s3_destination} SETTINGS allow_s3_native_copy = 1", + query_id=query_id, + settings=no_multipart, + ) + events = transfer_events(node, query_id) + + assert events["server_side_part_copy"] == 0, ( + "multipart copy was disabled but UploadPartCopy still ran" + ) + assert events["server_side_object_copy"] == 0, ( + "a ranged copy fell back to CopyObject, which would take the envelope" + ) + assert events["buffered_write"] > 0, ( + "nothing was uploaded through the server, so nothing was copied at all" + ) + + node.query(f"DROP TABLE IF EXISTS {restored} SYNC") + node.query( + f"RESTORE TABLE {table} AS {restored} FROM {s3_destination}", settings=no_multipart + ) + + actual = column_fingerprints(node, restored) + assert actual["count"] == expected["count"] + assert not [c for c in COLUMNS if actual[c] != expected[c]] + + node.query(f"DROP TABLE {table} SYNC") + node.query(f"DROP TABLE {restored} SYNC") + + +def test_backup_to_s3_ranged_copy_spans_several_parts(): + """A payload larger than one upload part makes the ranged copy issue several `UploadPartCopy` + requests. Every part but the first carries an offset of its own, so the payload window must be + applied to each one and not only to the first. + """ + node = cluster.instances["node"] + table = "cas_backup_multipart_range" + restored = f"{table}_restored" + s3_destination = backup_s3_destination("multipart_range") + query_id = f"{table}_backup_{RUN_TOKEN}" + small_parts = {"s3_min_upload_part_size": 5 * 1024 * 1024} + fingerprint = "SELECT count(), sum(cityHash64(s)) FROM {}" + + node.query(f"DROP TABLE IF EXISTS {table} SYNC") + node.query( + f""" + CREATE TABLE {table} (k UInt64, s String) + ENGINE = MergeTree ORDER BY k + SETTINGS storage_policy = '{STORAGE_POLICY}', min_bytes_for_wide_part = 0 + """ + ) + node.query( + f""" + INSERT INTO {table} + SELECT number, randomPrintableASCII(128) FROM numbers({NUM_ROWS}) + """ + ) + expected = node.query(fingerprint.format(table)).strip() + + node.query( + f"BACKUP TABLE {table} TO {s3_destination}", + query_id=query_id, + settings=small_parts, + ) + + events = transfer_events(node, query_id) + assert events["server_side_part_copy"] > 1, ( + "no blob was copied in more than one part, so the per-part offset stayed untested" + ) + assert events["server_side_object_copy"] == 0, ( + "CopyObject has no range: the envelope would land in the backup" + ) + + node.query(f"DROP TABLE IF EXISTS {restored} SYNC") + node.query(f"RESTORE TABLE {table} AS {restored} FROM {s3_destination}") + + assert node.query(fingerprint.format(restored)).strip() == expected + assert ( + node.query( + f"CHECK TABLE {restored} SETTINGS check_query_single_value_result = 1" + ).strip() + == "1" + ) + + node.query(f"DROP TABLE {table} SYNC") + node.query(f"DROP TABLE {restored} SYNC") + + +def test_move_partition_from_cas_to_s3_disk_with_empty_arrays(): + """A zero-size `.bin` still becomes a blob: `partFileMustStayBlob` keys on the file name, not + the size. Moving it off a CAS disk must copy it through buffers. + """ + node = cluster.instances["node"] + table = "cas_move_empty_arrays" + query_id = f"{table}_move_{RUN_TOKEN}" + + node.query(f"DROP TABLE IF EXISTS {table} SYNC") + node.query( + f""" + CREATE TABLE {table} (k UInt64, s String, empty Array(UInt32)) + ENGINE = MergeTree ORDER BY k + SETTINGS storage_policy = 'cas_then_plain', min_bytes_for_wide_part = 0 + """ + ) + node.query( + f""" + INSERT INTO {table} + SELECT number, randomPrintableASCII(64), [] + FROM numbers({NUM_ROWS}) + """ + ) + + expected = node.query( + f"SELECT count(), sum(cityHash64(s)), sum(length(empty)) FROM {table}" + ).strip() + + node.query( + f"ALTER TABLE {table} MOVE PARTITION tuple() TO DISK 'backup_disk_s3'", + query_id=query_id, + ) + + events = transfer_events(node, query_id) + assert (events["server_side_part_copy"], events["server_side_object_copy"]) == (0, 0), ( + "a CAS source is windowed, so no copy may run inside S3" + ) + assert events["buffered_write"] > 0, "every file must be moved through buffers" + + assert ( + node.query( + f"SELECT count(), sum(cityHash64(s)), sum(length(empty)) FROM {table}" + ).strip() + == expected + ) + assert ( + node.query( + f"CHECK TABLE {table} SETTINGS check_query_single_value_result = 1" + ).strip() + == "1" + ) + + node.query(f"DROP TABLE {table} SYNC") + + +def test_backup_to_s3_with_empty_arrays(): + """A zero-size `.bin` is a blob with an empty payload. A ranged copy of zero bytes reaches + `calculatePartSize(0)`, which throws, so it must be copied through buffers. Only a backup with + `deduplicate_files = 0` passes empty files to the writer; a deduplicated one drops them earlier. + """ + node = cluster.instances["node"] + table = "cas_backup_empty_arrays" + restored = f"{table}_restored" + s3_destination = backup_s3_destination("empty_arrays") + query_id = f"{table}_backup_{RUN_TOKEN}" + fingerprint = "SELECT count(), sum(cityHash64(s)), sum(length(empty)) FROM {}" + + node.query(f"DROP TABLE IF EXISTS {table} SYNC") + node.query( + f""" + CREATE TABLE {table} (k UInt64, s String, empty Array(UInt32)) + ENGINE = MergeTree ORDER BY k + SETTINGS storage_policy = '{STORAGE_POLICY}', min_bytes_for_wide_part = 0 + """ + ) + node.query( + f""" + INSERT INTO {table} + SELECT number, randomPrintableASCII(64), [] + FROM numbers({NUM_ROWS}) + """ + ) + expected = node.query(fingerprint.format(table)).strip() + + node.query( + f"BACKUP TABLE {table} TO {s3_destination} SETTINGS deduplicate_files = 0", + query_id=query_id, + ) + + events = transfer_events(node, query_id) + assert events["server_side_part_copy"] > 0, ( + "a zero-size blob must not cost the other blobs their ranged copy" + ) + assert events["server_side_object_copy"] == 0, ( + "CopyObject has no range: the envelope would land in the backup" + ) + + empty_files_in_backup = int( + node.query( + f""" + SELECT count() + FROM s3('{S3_AUTHORITY}/test/backups/{RUN_TOKEN}/empty_arrays/**', {S3_CREDENTIALS}, 'One') + WHERE _size = 0 + SETTINGS s3_skip_empty_files = 0 + """ + ).strip() + ) + assert empty_files_in_backup > 0, "no zero-size file reached the backup writer" + + node.query(f"DROP TABLE IF EXISTS {restored} SYNC") + node.query(f"RESTORE TABLE {table} AS {restored} FROM {s3_destination}") + + assert node.query(fingerprint.format(restored)).strip() == expected + assert ( + node.query( + f"CHECK TABLE {restored} SETTINGS check_query_single_value_result = 1" + ).strip() + == "1" + ) + + node.query(f"DROP TABLE {table} SYNC") + node.query(f"DROP TABLE {restored} SYNC") + + +def test_incremental_backup_to_s3(): + """An incremental backup copies only the files that changed since the base backup, so the + ranged copy runs against a subset of the part files and the restore reads from both backups. + """ + node = cluster.instances["node"] + table = "cas_backup_incremental" + restored = f"{table}_restored" + base = backup_s3_destination("incremental_base") + incremental = backup_s3_destination("incremental") + query_id = f"{table}_backup_{RUN_TOKEN}" + + create_and_fill(node, table) + node.query(f"BACKUP TABLE {table} TO {base}") + + node.query( + f""" + INSERT INTO {table} + SELECT + number, + randomPrintableASCII(64), + if(number % 5 = 0, NULL, toInt64(number)), + [toUInt32(number)] + FROM numbers({NUM_ROWS}, {NUM_ROWS}) + """ + ) + expected = column_fingerprints(node, table) + + node.query( + f"BACKUP TABLE {table} TO {incremental} SETTINGS base_backup = {base}", + query_id=query_id, + ) + + events = transfer_events(node, query_id) + assert events["server_side_part_copy"] > 0, ( + "the files of the new part must still be copied inside S3 with a range" + ) + assert events["server_side_object_copy"] == 0, ( + "CopyObject has no range: the envelope would land in the backup" + ) + assert events["buffered_write"] > 0, "inline entries have no object and go through buffers" + + node.query(f"DROP TABLE IF EXISTS {restored} SYNC") + node.query(f"RESTORE TABLE {table} AS {restored} FROM {incremental}") + + actual = column_fingerprints(node, restored) + assert actual["count"] == expected["count"] + differing = [c for c in COLUMNS if actual[c] != expected[c]] + assert not differing, f"columns differ after restore: {differing}" + + node.query(f"DROP TABLE {table} SYNC") + node.query(f"DROP TABLE {restored} SYNC") + + +def test_move_partition_between_plain_and_encrypted_s3_disks(): + """No CAS here. The capability predicate must only narrow: `sameKind` ignores `is_encrypted`, and + `DiskEncrypted` reports its delegate's description, so a predicate that replaced `operator==` + would let this pair through and the cast to `DiskObjectStorage` would throw. + """ + node = cluster.instances["node"] + table = "plain_to_encrypted" + to_encrypted_query_id = f"{table}_to_encrypted_{RUN_TOKEN}" + from_encrypted_query_id = f"{table}_from_encrypted_{RUN_TOKEN}" + + node.query(f"DROP TABLE IF EXISTS {table} SYNC") + node.query( + f""" + CREATE TABLE {table} (k UInt64, s String) + ENGINE = MergeTree ORDER BY k + SETTINGS storage_policy = 'plain_then_encrypted', min_bytes_for_wide_part = 0 + """ + ) + node.query( + f"INSERT INTO {table} SELECT number, randomPrintableASCII(64) FROM numbers({NUM_ROWS})" + ) + + expected = node.query(f"SELECT count(), sum(cityHash64(s)) FROM {table}").strip() + + node.query( + f"ALTER TABLE {table} MOVE PARTITION tuple() TO DISK 'disk_plain_s3_encrypted'", + query_id=to_encrypted_query_id, + ) + to_encrypted = transfer_events(node, to_encrypted_query_id) + assert to_encrypted["buffered_write"] > 0, ( + "one side encrypts and the other does not, so the bytes must pass through the server" + ) + assert ( + node.query(f"SELECT count(), sum(cityHash64(s)) FROM {table}").strip() == expected + ) + + node.query( + f"ALTER TABLE {table} MOVE PARTITION tuple() TO DISK 'backup_disk_s3'", + query_id=from_encrypted_query_id, + ) + from_encrypted = transfer_events(node, from_encrypted_query_id) + assert from_encrypted["buffered_write"] > 0, ( + "the way back decrypts, so the bytes must pass through the server again" + ) + assert ( + node.query(f"SELECT count(), sum(cityHash64(s)) FROM {table}").strip() == expected + ) + + node.query(f"DROP TABLE {table} SYNC") + + +def test_move_partition_between_plain_s3_disks_copies_server_side(): + """The capability defaults to false, so a description that forgets to claim it silently loses the + server-side copy. Two ordinary s3 disks must keep it. + """ + node = cluster.instances["node"] + table = "plain_to_plain" + query_id = f"plain_to_plain_move_{RUN_TOKEN}" + + node.query(f"DROP TABLE IF EXISTS {table} SYNC") + node.query( + f""" + CREATE TABLE {table} (k UInt64, s String) + ENGINE = MergeTree ORDER BY k + SETTINGS storage_policy = 'plain_then_plain', min_bytes_for_wide_part = 0 + """ + ) + node.query( + f"INSERT INTO {table} SELECT number, randomPrintableASCII(64) FROM numbers({NUM_ROWS})" + ) + expected = node.query(f"SELECT count(), sum(cityHash64(s)) FROM {table}").strip() + + node.query( + f"ALTER TABLE {table} MOVE PARTITION tuple() TO DISK 'disk_plain_s3_second'", + query_id=query_id, + ) + + events = transfer_events(node, query_id) + assert events["server_side_object_copy"] > 0, ( + "a non-CAS disk pair lost its server-side copy: a whole file is a whole object here, so " + "CopyObject is the expected operation" + ) + assert ( + node.query(f"SELECT count(), sum(cityHash64(s)) FROM {table}").strip() == expected + ) + + node.query(f"DROP TABLE {table} SYNC") + + +def test_backup_to_file_keeps_fs_copy(): + """Both sides take their description from `DiskLocal::getLocalDataSourceDescription`. A forgotten + claim there gives `false && false`, and `BackupWriterFile` silently stops using `fs::copy`. + """ + node = cluster.instances["node"] + table = "plain_local_to_file" + query_id = f"plain_local_backup_{RUN_TOKEN}" + + node.query(f"DROP TABLE IF EXISTS {table} SYNC") + node.query( + f""" + CREATE TABLE {table} (k UInt64, s String) + ENGINE = MergeTree ORDER BY k + SETTINGS min_bytes_for_wide_part = 0 + """ + ) + node.query( + f"INSERT INTO {table} SELECT number, randomPrintableASCII(64) FROM numbers({NUM_ROWS})" + ) + + node.query( + f"BACKUP TABLE {table} TO File('{RUN_TOKEN}/file_backup')", query_id=query_id + ) + node.query("SYSTEM FLUSH LOGS query_log") + + read_bytes = int( + node.query( + f""" + SELECT ProfileEvents['ReadBufferFromFileDescriptorReadBytes'] + FROM system.query_log + WHERE type = 'QueryFinish' AND query_id = '{query_id}' + ORDER BY event_time DESC LIMIT 1 + """ + ).strip() + ) + bytes_on_disk = int( + node.query( + f""" + SELECT sum(bytes_on_disk) FROM system.parts + WHERE table = '{table}' AND active + """ + ).strip() + ) + + assert bytes_on_disk > 0, "the table has no active parts, so the assertion below proves nothing" + assert read_bytes < bytes_on_disk * 3 // 2, ( + f"BACKUP TO File(...) read {read_bytes} bytes with {bytes_on_disk} on disk. One pass is the " + "checksum pass every entry pays; a second pass means fs::copy was replaced by a buffered copy" + ) + + node.query(f"DROP TABLE {table} SYNC") + + +def test_move_partition_from_cached_cas_to_s3_disk_only_buffers(): + """`wrapWithCache` reuses the CAS metadata storage for a CAS disk, so the cache disk must inherit + the same answer and stay out of the server-side copy. + """ + node = cluster.instances["node"] + table = "cas_cached_move" + query_id = f"cas_cached_move_{RUN_TOKEN}" + + create_and_fill(node, table, "cas_cached_then_plain") + expected = column_fingerprints(node, table) + + node.query( + f"ALTER TABLE {table} MOVE PARTITION tuple() TO DISK 'backup_disk_s3'", + query_id=query_id, + ) + + events = transfer_events(node, query_id) + assert (events["server_side_part_copy"], events["server_side_object_copy"]) == (0, 0), ( + "a cache disk over CAS reported itself as whole-object" + ) + assert events["buffered_write"] > 0, "the moved part must be written through buffers" + assert column_fingerprints(node, table) == expected + + node.query(f"DROP TABLE {table} SYNC") + + +def test_backup_to_disk_cas_is_rejected(): + """A `CAS` disk takes a part only as a whole part in one transaction. A backup's own layout mirrors + the table's data directory, so its files sit under a part directory too and `isPartFilePath` matches + them - which is why writing a backup into a `CAS` disk is refused rather than silently accepted. + """ + node = cluster.instances["node"] + table = "backup_into_cas" + disk_destination = backup_disk_destination("disk_cas_backup_s3", table) + + node.query(f"DROP TABLE IF EXISTS {table} SYNC") + node.query( + f""" + CREATE TABLE {table} (k UInt64, s String) + ENGINE = MergeTree ORDER BY k + SETTINGS min_bytes_for_wide_part = 0 + """ + ) + node.query( + f"INSERT INTO {table} SELECT number, randomPrintableASCII(64) FROM numbers({NUM_ROWS})" + ) + + with pytest.raises(QueryRuntimeException) as raised: + node.query(f"BACKUP TABLE {table} TO {disk_destination}") + assert "Autocommit writes are not supported for content part files" in str(raised.value) + + node.query(f"DROP TABLE {table} SYNC")