Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions .github/workflows/CI.yml
Original file line number Diff line number Diff line change
Expand Up @@ -47,9 +47,9 @@ jobs:
- uses: actions/checkout@v4
- name: Install Azure CLI
run: curl -sL https://aka.ms/InstallAzureCLIDeb | sudo bash
- name: Start LocalStack (AWS S3)
- name: Start Floci (AWS S3)
run: |
container=$(docker create -p 4566:4566 localstack/localstack:s3-latest)
container=$(docker create -p 4566:4566 floci/floci:latest)
docker start $container
sleep 5
docker logs $container
Expand All @@ -61,7 +61,7 @@ jobs:
docker logs $container
- name: Create test AWS S3 bucket
run: |
AWS_ACCESS_KEY_ID=test AWS_SECRET_ACCESS_KEY=test AWS_DEFAULT_REGION=${DEFAULT_REGION:-$AWS_DEFAULT_REGION} aws --endpoint-url=http://localhost:4566 s3api create-bucket --bucket cloud-copy-test
AWS_ACCESS_KEY_ID=test AWS_SECRET_ACCESS_KEY=test aws --endpoint-url=http://localhost:4566 s3api create-bucket --bucket cloud-copy-test
- name: Create test Azure Storage container
run: |
az storage container create --name cloud-copy-test --connection-string "DefaultEndpointsProtocol=http;AccountName=devstoreaccount1;AccountKey=Eby8vdM02xNOcqFlqUwJPLlmEtlCDXJ1OUzFT50uSRZ6IFsuFq2UVErCz4I6tq/K1SZFPTOtr/KBHBeksoGMGw==;BlobEndpoint=http://localhost:10000/devstoreaccount1"
Expand Down
8 changes: 8 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,14 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## Unreleased

#### Fixed

* Fixed `walk` implementation to filter "directory" blobs from Azure Blob ([#36](https://github.com/stjude-rust-labs/cloud-copy/pull/36)).

#### Changed

* Switched to `floci` for testing S3 in CI ([#36](https://github.com/stjude-rust-labs/cloud-copy/pull/36)).

## 0.8.0 - 03-20-2026

#### Added
Expand Down
12 changes: 6 additions & 6 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -135,18 +135,18 @@ Automated tests rely on having the following cloud service emulators installed:

* [Azurite](https://learn.microsoft.com/en-us/azure/storage/common/storage-use-azurite?tabs=visual-studio%2Cblob-storage)
for Azure Blob Storage.
* [Localstack](https://github.com/localstack/localstack) for AWS S3.
* [Floci](https://github.com/floci-io/floci) for AWS S3.

Use the following command to run Azurite in a Docker container:

```bash
docker run -p 10000:10000 mcr.microsoft.com/azure-storage/azurite azurite -l /data --blobHost 0.0.0.0 --loose
docker run -p 10000:10000 mcr.microsoft.com/azure-storage/azurite azurite -l /data --blobHost 0.0.0.0 --loose --skipApiVersionCheck
```

Use the following command to run Localstack in a Docker container:
Use the following command to run Floci in a Docker container:

```bash
docker run -p 4566:4566 localstack/localstack:s3-latest
docker run -p 4566:4566 floci/floci:latest
```

The tests expect a container/bucket with the name `cloud-copy-test` to be
Expand All @@ -158,10 +158,10 @@ To create the container with Azurite, use the Azure CLI:
az storage container create --name cloud-copy-test --connection-string "DefaultEndpointsProtocol=http;AccountName=devstoreaccount1;AccountKey=Eby8vdM02xNOcqFlqUwJPLlmEtlCDXJ1OUzFT50uSRZ6IFsuFq2UVErCz4I6tq/K1SZFPTOtr/KBHBeksoGMGw==;BlobEndpoint=http://127.0.0.1:10000/devstoreaccount1"
```

To create the bucket with Localstack, use the `awslocal` tool:
To create the bucket with Floci, use the `aws` tool:

```bash
awslocal s3api create-bucket --bucket cloud-copy-test
AWS_ACCESS_KEY_ID=test AWS_SECRET_ACCESS_KEY=test aws --endpoint-url=http://localhost:4566 s3api create-bucket --bucket cloud-copy-test
```

Finally, run the tests:
Expand Down
2 changes: 2 additions & 0 deletions src/backend/azure.rs
Original file line number Diff line number Diff line change
Expand Up @@ -746,6 +746,8 @@ impl StorageBackend for AzureBlobStorageBackend {
pairs.append_pair("comp", "list");
// The prefix to use for listing blobs in the container.
pairs.append_pair("prefix", &prefix);
// Only include files in the output.
pairs.append_pair("showonly", "files");

// Only return at most one result if we're returning the first only
if first_only {
Expand Down
15 changes: 6 additions & 9 deletions src/backend/s3.rs
Original file line number Diff line number Diff line change
Expand Up @@ -42,9 +42,6 @@ use crate::streams::TransferStream;
/// The root domain for AWS.
const AWS_ROOT_DOMAIN: &str = "amazonaws.com";

/// The root domain for localstack.
const LOCALSTACK_ROOT_DOMAIN: &str = "localhost.localstack.cloud";

/// The maximum number of parts in an upload.
const MAX_PARTS: u64 = 10000;

Expand Down Expand Up @@ -555,8 +552,8 @@ impl StorageBackend for S3StorageBackend {
// There must be at least two path segments
!region.is_empty()
&& (domain.eq_ignore_ascii_case(AWS_ROOT_DOMAIN)
|| (config.s3().use_localstack()
&& domain.eq_ignore_ascii_case(LOCALSTACK_ROOT_DOMAIN)))
|| (config.s3().use_floci()
&& domain.eq_ignore_ascii_case("localhost")))
&& url
.path_segments()
.map(|mut s| s.nth(1).is_some())
Expand All @@ -571,8 +568,8 @@ impl StorageBackend for S3StorageBackend {
&& !region.is_empty()
&& service.eq_ignore_ascii_case("s3")
&& (domain.eq_ignore_ascii_case(AWS_ROOT_DOMAIN)
|| (config.s3().use_localstack()
&& domain.eq_ignore_ascii_case(LOCALSTACK_ROOT_DOMAIN)))
|| (config.s3().use_floci()
&& domain.eq_ignore_ascii_case("localhost")))
&& url
.path_segments()
.map(|mut s| s.next().is_some())
Expand All @@ -597,8 +594,8 @@ impl StorageBackend for S3StorageBackend {
return Err(S3Error::InvalidScheme.into());
}

let (scheme, root, port) = if config.s3().use_localstack() {
("http", LOCALSTACK_ROOT_DOMAIN, ":4566")
let (scheme, root, port) = if config.s3().use_floci() {
("http", "localhost", ":4566")
} else {
("https", AWS_ROOT_DOMAIN, "")
};
Expand Down
18 changes: 9 additions & 9 deletions src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -228,9 +228,9 @@ pub struct S3Config {
/// If `None`, no authentication header will be put on requests.
#[serde(default)]
auth: Option<S3AuthConfig>,
/// Stores whether or not localstack is being used.
/// Stores whether or not Floci is being used.
#[serde(default)]
use_localstack: bool,
use_floci: bool,
}

impl S3Config {
Expand Down Expand Up @@ -265,15 +265,15 @@ impl S3Config {
self
}

/// Sets whether or not [localstack](https://github.com/localstack/localstack) is being used.
/// Sets whether or not [Floci](https://github.com/floci-io/floci) is being used.
///
/// The domain suffix is expected to be `localhost.localstack.cloud`.
/// The domain suffix is expected to be `localhost` for Floci requests.
///
/// Any URLs that use the `s3` scheme will be rewritten to use that suffix.
///
/// This setting is primarily intended for local testing.
pub fn with_use_localstack(mut self, use_localstack: bool) -> Self {
self.use_localstack = use_localstack;
pub fn with_use_floci(mut self, use_floci: bool) -> Self {
self.use_floci = use_floci;
self
}

Expand All @@ -291,9 +291,9 @@ impl S3Config {
self.auth.as_ref()
}

/// Gets whether or not [localstack](https://github.com/localstack/localstack) is being used.
pub fn use_localstack(&self) -> bool {
self.use_localstack
/// Gets whether or not [Floci](https://github.com/floci-io/floci) is being used.
pub fn use_floci(&self) -> bool {
self.use_floci
}
}

Expand Down
88 changes: 82 additions & 6 deletions src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ use std::borrow::Cow;
use std::fmt;
use std::io::ErrorKind;
use std::ops::Deref;
use std::path::Component;
use std::path::Path;
use std::path::PathBuf;
use std::sync::Arc;
Expand Down Expand Up @@ -114,6 +115,58 @@ fn sha256_hex_string(bytes: impl AsRef<[u8]>) -> String {
hex::encode(hash.finalize())
}

/// Sorts a set of walk entries.
///
/// Entries are expected to be resource subpaths, e.g. `foo` and `bar/baz`.
///
/// Each entry is sorted by path segment.
///
/// If an entry prefixes another, it means that a file would conflict with a
/// directory and an error is returned.
///
/// Each entry may not contain a path segment that references the root or a
/// parent directory, otherwise an error is returned.
fn sort_walk_entries(url: &Url, entries: &mut [String]) -> Result<()> {
// Sort each entry by path components
entries.sort_by(|a, b| {
let a = Path::new(a);
let b = Path::new(b);
a.components().cmp(b.components())
});

// Check for invalid or conflicting entries
let mut iter = entries.iter().peekable();
while let Some(entry) = iter.next() {
// Ensure entry segments don't reference the root or parent directory
if Path::new(entry)
.components()
.any(|c| !matches!(c, Component::Normal(_) | Component::CurDir))
{
return Err(Error::InvalidWalkEntry {
url: url.display().to_string(),
entry: entry.clone(),
});
}

// Ensure the next entry in the sorted list is not prefixed with this entry,
// otherwise a file conflicts with a directory
if let Some(next) = iter.peek()
Comment thread
peterhuene marked this conversation as resolved.
&& next
.strip_prefix(entry)
.and_then(|p| p.strip_prefix('/'))
.is_some()
{
return Err(Error::WalkEntryConflict {
url: url.display().to_string(),
first: entry.clone(),
second: (*next).clone(),
});
}
}

Ok(())
}

trait DateTimeExt {
/// Converts a [`DateTime`] to a HTTP date string.
///
Expand Down Expand Up @@ -421,6 +474,24 @@ pub enum Error {
/// The remote content was modified during a download.
#[error("the remote content was modified during the download")]
RemoteContentModified,
/// Failed to walk a URL due to an invalid entry.
#[error("failed to walk URL `{url}` due to invalid entry `{entry}`")]
InvalidWalkEntry {
/// The URL that failed to walk.
url: String,
/// The invalid entry
entry: String,
},
/// Failed to walk a URL due to a conflict in its entries.
#[error("failed to walk URL `{url}` due to conflicting entries `{first}` and `{second}`")]
WalkEntryConflict {
/// The URL that failed to walk.
url: String,
/// The first conflicting entry.
first: String,
/// The second conflicting entry.
second: String,
},
/// Failed to create a directory.
#[error("failed to create directory `{path}`: {error}", path = .path.display())]
DirectoryCreationFailed {
Expand Down Expand Up @@ -819,7 +890,8 @@ pub fn rewrite_url<'a>(config: &Config, url: &'a Url) -> Result<Cow<'a, Url>> {

/// Walks a given storage URL as if it were a directory.
///
/// Returns a list of relative paths from the given URL.
/// Returns a lexicographically-sorted list of relative paths from the given
/// URL.
///
/// If the given storage URL is not a directory, an empty list is returned.
pub async fn walk(config: Config, client: HttpClient, mut url: Url) -> Result<Vec<String>> {
Expand All @@ -829,7 +901,7 @@ pub async fn walk(config: Config, client: HttpClient, mut url: Url) -> Result<Ve
segments.pop_if_empty().push("");
}

if AzureBlobStorageBackend::is_supported_url(&config, &url) {
let mut entries = if AzureBlobStorageBackend::is_supported_url(&config, &url) {
let url = AzureBlobStorageBackend::rewrite_url(&config, &url)?;
AzureBlobStorageBackend::new(config, client, None)
.walk(url.into_owned(), false)
Expand All @@ -845,8 +917,12 @@ pub async fn walk(config: Config, client: HttpClient, mut url: Url) -> Result<Ve
.walk(url.into_owned(), false)
.await
} else {
Err(Error::UnsupportedUrl(url))
}
Err(Error::UnsupportedUrl(url.clone()))
}?;

// Sort the entries
sort_walk_entries(&url, &mut entries)?;
Ok(entries)
}

/// Represents the content digest of a resource.
Expand Down Expand Up @@ -1097,7 +1173,7 @@ mod test {

let config = Config::builder()
.with_azure(AzureConfig::default().with_use_azurite(true))
.with_s3(S3Config::default().with_use_localstack(true))
.with_s3(S3Config::default().with_use_floci(true))
.build();

assert_eq!(
Expand All @@ -1111,7 +1187,7 @@ mod test {
rewrite_url(&config, &"s3://foo/bar/baz".parse().unwrap())
.unwrap()
.as_str(),
"http://foo.s3.us-east-1.localhost.localstack.cloud:4566/bar/baz"
"http://foo.s3.us-east-1.localhost:4566/bar/baz",
);
}
}
12 changes: 8 additions & 4 deletions src/transfer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ use crate::backend::Upload;
use crate::notify_retry;
use crate::pool::BufferGuard;
use crate::pool::BufferPool;
use crate::sort_walk_entries;
use crate::streams::TransferStream;

/// The supported `accept-range` unit.
Expand Down Expand Up @@ -795,7 +796,7 @@ where
let destination = destination.as_ref();

// Start by walking the given URL for files to download
let paths = Retry::spawn_notify(
let mut entries = Retry::spawn_notify(
self.inner.backend.config().retry_durations(),
|| async {
select! {
Expand All @@ -809,6 +810,9 @@ where
)
.await?;

// Sort the entries
sort_walk_entries(&source, &mut entries)?;

// Delete the destination if it exists
if let Ok(metadata) = destination.metadata() {
if metadata.is_file() {
Expand All @@ -822,15 +826,15 @@ where
let mut set = JoinSet::new();

// If there are no files relative to the given URL, download just the given URL
if paths.is_empty() {
if entries.is_empty() {
let inner = self.inner.clone();
let destination = destination.to_path_buf();
let permits = permits.clone();
let cancel = self.cancel.clone();
set.spawn(async move { inner.download(source, &destination, permits, cancel).await });
} else {
// Otherwise, download each file in turn
for path in paths {
for entry in entries {
let inner = self.inner.clone();
let mut source = source.clone();
let mut destination = destination.to_path_buf();
Expand All @@ -841,7 +845,7 @@ where
{
let mut segments = source.path_segments_mut().expect("URL should have a path");
segments.pop_if_empty();
for segment in path.split('/') {
for segment in entry.split('/') {
segments.push(segment);
destination.push(segment);
}
Expand Down
6 changes: 3 additions & 3 deletions tests/copy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -161,10 +161,10 @@ fn urls(test: &str) -> Vec<Url> {
vec![
// S3 URLs
format!("s3://{TEST_BUCKET_NAME}/1/{test}").parse().unwrap(),
format!("http://s3.us-east-1.localhost.localstack.cloud:4566/{TEST_BUCKET_NAME}/2/{test}")
format!("http://s3.us-east-1.localhost:4566/{TEST_BUCKET_NAME}/2/{test}")
.parse()
.unwrap(),
format!("http://{TEST_BUCKET_NAME}.s3.us-east-1.localhost.localstack.cloud:4566/3/{test}")
format!("http://{TEST_BUCKET_NAME}.s3.us-east-1.localhost:4566/3/{test}")
.parse()
.unwrap(),

Expand All @@ -184,7 +184,7 @@ fn azure_config() -> AzureConfig {

fn s3_config() -> S3Config {
S3Config::default()
.with_use_localstack(true)
.with_use_floci(true)
.with_auth("test", "test")
}

Expand Down
Loading