feat: add corex-storage crate for unified storage abstraction

- Introduced corex-storage crate with support for local disk, S3-compatible, and Google Cloud Storage backends.
- Implemented StorageDriver trait for various storage backends.
- Added LocalFileStorage for local disk operations.
- Added S3Storage for S3-compatible object storage with multipart upload support.
- Added GcsStorage for Google Cloud Storage operations.
- Included error handling for storage operations.
- Added tests for each storage backend to ensure functionality.
- Created README.md for documentation and usage examples.
- Added Apache and MIT licenses for open-source compliance.
This commit is contained in:
asepharyana
2026-08-28 22:24:18 +07:00
parent 6a4b2b8f54
commit db4f277336
47 changed files with 0 additions and 0 deletions
+44
View File
@@ -0,0 +1,44 @@
[package]
name = "corex-storage"
version = "1.2.0"
edition = "2021"
rust-version = "1.75"
license = "MIT OR Apache-2.0"
repository = "https://github.com/asepharyana/mytheclipse"
homepage = "https://github.com/asepharyana/mytheclipse"
documentation = "https://docs.rs/corex-storage"
authors = ["asepharyana <superaseph@gmail.com>"]
description = "Unified storage & file system abstraction: one driver interface over Local Disk, S3/MinIO, and Google Cloud Storage, with stream-based upload/download."
readme = "README.md"
keywords = ["storage", "s3", "gcs", "minio", "filesystem"]
categories = ["filesystem", "asynchronous"]
[features]
default = ["local"]
# The core `StorageDriver` trait (always compiled) needs `tokio`'s io-util for
# `AsyncRead`/`ReadBuf`; `local` additionally needs `fs`/`rt`.
local = ["tokio/fs", "tokio/rt"]
# S3-compatible object storage (also covers MinIO via a custom endpoint).
s3 = ["dep:aws-sdk-s3", "dep:aws-config", "dep:aws-credential-types"]
# Google Cloud Storage.
gcs = ["dep:google-cloud-storage", "dep:google-cloud-auth"]
[dependencies]
tracing = "0.1"
# Required unconditionally: the core `StorageDriver` trait uses both.
async-trait = "0.1"
tokio = { version = "1.53", features = ["io-util"] }
bytes = "1"
# S3 / MinIO
aws-sdk-s3 = { version = "1", optional = true }
aws-config = { version = "1", optional = true, features = ["behavior-version-latest"] }
aws-credential-types = { version = "1", optional = true }
# Google Cloud Storage
google-cloud-storage = { version = "0.24", optional = true }
google-cloud-auth = { version = "0.17", optional = true }
[dev-dependencies]
tokio = { version = "1.53", features = ["full"] }
tempfile = "3"
+201
View File
@@ -0,0 +1,201 @@
Apache License
Version 2.0, January 2004
http://www.apache.org/licenses/
TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION
1. Definitions.
"License" shall mean the terms and conditions for use, reproduction,
and distribution as defined by Sections 1 through 9 of this document.
"Licensor" shall mean the copyright owner or entity authorized by
the copyright owner that is granting the License.
"Legal Entity" shall mean the union of the acting entity and all
other entities that control, are controlled by, or are under common
control with that entity. For the purposes of this definition,
"control" means (i) the power, direct or indirect, to cause the
direction or management of such entity, whether by contract or
otherwise, or (ii) ownership of fifty percent (50%) or more of the
outstanding shares, or (iii) beneficial ownership of such entity.
"You" (or "Your") shall mean an individual or Legal Entity
exercising permissions granted by this License.
"Source" form shall mean the preferred form for making modifications,
including but not limited to software source code, documentation
source, and configuration files.
"Object" form shall mean any form resulting from mechanical
transformation or translation of a Source form, including but
not limited to compiled object code, generated documentation,
and conversions to other media types.
"Work" shall mean the work of authorship, whether in Source or
Object form, made available under the License, as indicated by a
copyright notice that is included in or attached to the work
(an example is provided in the Appendix below).
"Derivative Works" shall mean any work, whether in Source or Object
form, that is based on (or derived from) the Work and for which the
editorial revisions, annotations, elaborations, or other modifications
represent, as a whole, an original work of authorship. For the purposes
of this License, Derivative Works shall not include works that remain
separable from, or merely link (or bind by name) to the interfaces of,
the Work and Derivative Works thereof.
"Contribution" shall mean any work of authorship, including
the original version of the Work and any modifications or additions
to that Work or Derivative Works thereof, that is intentionally
submitted to Licensor for inclusion in the Work by the copyright owner
or by an individual or Legal Entity authorized to submit on behalf of
the copyright owner. For the purposes of this definition, "submitted"
means any form of electronic, verbal, or written communication sent
to the Licensor or its representatives, including but not limited to
communication on electronic mailing lists, source code control systems,
and issue tracking systems that are managed by, or on behalf of, the
Licensor for the purpose of discussing and improving the Work, but
excluding communication that is conspicuously marked or otherwise
designated in writing by the copyright owner as "Not a Contribution."
"Contributor" shall mean Licensor and any individual or Legal Entity
on behalf of whom a Contribution has been received by Licensor and
subsequently incorporated within the Work.
2. Grant of Copyright License. Subject to the terms and conditions of
this License, each Contributor hereby grants to You a perpetual,
worldwide, non-exclusive, no-charge, royalty-free, irrevocable
copyright license to reproduce, prepare Derivative Works of,
publicly display, publicly perform, sublicense, and distribute the
Work and such Derivative Works in Source or Object form.
3. Grant of Patent License. Subject to the terms and conditions of
this License, each Contributor hereby grants to You a perpetual,
worldwide, non-exclusive, no-charge, royalty-free, irrevocable
(except as stated in this section) patent license to make, have made,
use, offer to sell, sell, import, and otherwise transfer the Work,
where such license applies only to those patent claims licensable
by such Contributor that are necessarily infringed by their
Contribution(s) alone or by combination of their Contribution(s)
with the Work to which such Contribution(s) was submitted. If You
institute patent litigation against any entity (including a
cross-claim or counterclaim in a lawsuit) alleging that the Work
or a Contribution incorporated within the Work constitutes direct
or contributory patent infringement, then any patent licenses
granted to You under this License for that Work shall terminate
as of the date such litigation is filed.
4. Redistribution. You may reproduce and distribute copies of the
Work or Derivative Works thereof in any medium, with or without
modifications, and in Source or Object form, provided that You
meet the following conditions:
(a) You must give any other recipients of the Work or
Derivative Works a copy of this License; and
(b) You must cause any modified files to carry prominent notices
stating that You changed the files; and
(c) You must retain, in the Source form of any Derivative Works
that You distribute, all copyright, patent, trademark, and
attribution notices from the Source form of the Work,
excluding those notices that do not pertain to any part of
the Derivative Works; and
(d) If the Work includes a "NOTICE" text file as part of its
distribution, then any Derivative Works that You distribute must
include a readable copy of the attribution notices contained
within such NOTICE file, excluding those notices that do not
pertain to any part of the Derivative Works, in at least one
of the following places: within a NOTICE text file distributed
as part of the Derivative Works; within the Source form or
documentation, if provided along with the Derivative Works; or,
within a display generated by the Derivative Works, if and
wherever such third-party notices normally appear. The contents
of the NOTICE file are for informational purposes only and
do not modify the License. You may add Your own attribution
notices within Derivative Works that You distribute, alongside
or as an addendum to the NOTICE text from the Work, provided
that such additional attribution notices cannot be construed
as modifying the License.
You may add Your own copyright statement to Your modifications and
may provide additional or different license terms and conditions
for use, reproduction, or distribution of Your modifications, or
for any such Derivative Works as a whole, provided Your use,
reproduction, and distribution of the Work otherwise complies with
the conditions stated in this License.
5. Submission of Contributions. Unless You explicitly state otherwise,
any Contribution intentionally submitted for inclusion in the Work
by You to the Licensor shall be under the terms and conditions of
this License, without any additional terms or conditions.
Notwithstanding the above, nothing herein shall supersede or modify
the terms of any separate license agreement you may have executed
with Licensor regarding such Contributions.
6. Trademarks. This License does not grant permission to use the trade
names, trademarks, service marks, or product names of the Licensor,
except as required for reasonable and customary use in describing the
origin of the Work and reproducing the content of the NOTICE file.
7. Disclaimer of Warranty. Unless required by applicable law or
agreed to in writing, Licensor provides the Work (and each
Contributor provides its Contributions) on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
implied, including, without limitation, any warranties or conditions
of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A
PARTICULAR PURPOSE. You are solely responsible for determining the
appropriateness of using or redistributing the Work and assume any
risks associated with Your exercise of permissions under this License.
8. Limitation of Liability. In no event and under no legal theory,
whether in tort (including negligence), contract, or otherwise,
unless required by applicable law (such as deliberate and grossly
negligent acts) or agreed to in writing, shall any Contributor be
liable to You for damages, including any direct, indirect, special,
incidental, or consequential damages of any character arising as a
result of this License or out of the use or inability to use the
Work (including but not limited to damages for loss of goodwill,
work stoppage, computer failure or malfunction, or any and all
other commercial damages or losses), even if such Contributor
has been advised of the possibility of such damages.
9. Accepting Warranty or Additional Liability. While redistributing
the Work or Derivative Works thereof, You may choose to offer,
and charge a fee for, acceptance of support, warranty, indemnity,
or other liability obligations and/or rights consistent with this
License. However, in accepting such obligations, You may act only
on Your own behalf and on Your sole responsibility, not on behalf
of any other Contributor, and only if You agree to indemnify,
defend, and hold each Contributor harmless for any liability
incurred by, or claims asserted against, such Contributor by reason
of your accepting any such warranty or additional liability.
END OF TERMS AND CONDITIONS
APPENDIX: How to apply the Apache License to your work.
To apply the Apache License to your work, attach the following
boilerplate notice, with the fields enclosed by brackets "[]"
replaced with your own identifying information. (Don't include
the brackets!) The text should be enclosed in the appropriate
comment syntax for the file format. We also recommend that a
file or class name and description of purpose be included on the
same "printed page" as the copyright notice for easier
identification within third-party archives.
Copyright 2026 The corex Authors
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
+21
View File
@@ -0,0 +1,21 @@
MIT License
Copyright (c) 2026 The corex Authors
Permission is hereby granted, free of charge, to any person obtaining a copy
of this software and associated documentation files (the "Software"), to deal
in the Software without restriction, including without limitation the rights
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
copies of the Software, and to permit persons to whom the Software is
furnished to do so, subject to the following conditions:
The above copyright notice and this permission notice shall be included in all
copies or substantial portions of the Software.
THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
SOFTWARE.
+26
View File
@@ -0,0 +1,26 @@
# corex-storage
A unified storage & file system abstraction so file handling isn't locked to
one physical location.
- **Local disk** (default) — plain files under a root directory.
- **S3-compatible object storage** (`s3`) — Amazon S3, and any S3-compatible
endpoint (MinIO included) via a custom endpoint URL.
- **Google Cloud Storage** (`gcs`).
All operations are stream-based, so uploading/downloading a huge object never
requires holding it entirely in memory.
## Usage
```rust
use corex_storage::{LocalFileStorage, StorageDriver, bytes_stream, read_to_vec};
let storage = LocalFileStorage::new("/var/data");
storage.put("uploads/report.csv", bytes_stream(data)).await?;
let bytes = read_to_vec(storage.get("uploads/report.csv").await?).await?;
```
Swap `LocalFileStorage::new(path)` for `S3Storage::connect("bucket").await?` or
`GcsStorage::connect("bucket").await?` to move to a distributed backend
without touching call sites.
+165
View File
@@ -0,0 +1,165 @@
//! A Google Cloud Storage [`StorageDriver`] (feature `gcs`), via
//! `google-cloud-storage`.
//!
//! [`get`](GcsStorage::get) streams the downloaded object; [`put`](GcsStorage::put)
//! currently buffers the input before uploading in a single request (the
//! `google-cloud-storage` crate's simple upload API takes an owned body). For
//! very large uploads prefer chunked application-level batching until this
//! crate grows resumable-upload support.
use async_trait::async_trait;
use google_cloud_storage::client::{Client, ClientConfig};
use google_cloud_storage::http::objects::delete::DeleteObjectRequest;
use google_cloud_storage::http::objects::download::Range;
use google_cloud_storage::http::objects::get::GetObjectRequest;
use google_cloud_storage::http::objects::upload::{Media, UploadObjectRequest, UploadType};
use tokio::io::AsyncReadExt;
use crate::traits::{ObjectMeta, ObjectStream, StorageDriver, StorageError};
fn map_err<E: std::fmt::Display>(e: E) -> StorageError {
StorageError::Io(e.to_string())
}
/// A [`StorageDriver`] backed by a Google Cloud Storage bucket.
#[derive(Clone)]
pub struct GcsStorage {
client: Client,
bucket: String,
}
impl GcsStorage {
/// Connects using Application Default Credentials (`GOOGLE_APPLICATION_CREDENTIALS`,
/// workload identity, etc.).
pub async fn connect(bucket: impl Into<String>) -> Result<Self, StorageError> {
let config = ClientConfig::default().with_auth().await.map_err(map_err)?;
Ok(Self {
client: Client::new(config),
bucket: bucket.into(),
})
}
/// Wraps an already-configured client.
pub fn from_client(client: Client, bucket: impl Into<String>) -> Self {
Self {
client,
bucket: bucket.into(),
}
}
}
#[async_trait]
impl StorageDriver for GcsStorage {
async fn get(&self, path: &str) -> Result<ObjectStream, StorageError> {
let bytes = self
.client
.download_object(
&GetObjectRequest {
bucket: self.bucket.clone(),
object: path.to_string(),
..Default::default()
},
&Range::default(),
)
.await
.map_err(|e| {
if e.to_string().contains("404") {
StorageError::NotFound(path.to_string())
} else {
map_err(e)
}
})?;
Ok(crate::traits::bytes_stream(bytes))
}
async fn put(&self, path: &str, mut data: ObjectStream) -> Result<u64, StorageError> {
let mut buf = Vec::new();
data.read_to_end(&mut buf)
.await
.map_err(|e| StorageError::Io(e.to_string()))?;
let len = buf.len() as u64;
let upload_type = UploadType::Simple(Media::new(path.to_string()));
self.client
.upload_object(
&UploadObjectRequest {
bucket: self.bucket.clone(),
..Default::default()
},
buf,
&upload_type,
)
.await
.map_err(map_err)?;
Ok(len)
}
async fn delete(&self, path: &str) -> Result<(), StorageError> {
self.client
.delete_object(&DeleteObjectRequest {
bucket: self.bucket.clone(),
object: path.to_string(),
..Default::default()
})
.await
.map_err(map_err)
}
async fn exists(&self, path: &str) -> Result<bool, StorageError> {
match self.stat(path).await {
Ok(_) => Ok(true),
Err(StorageError::NotFound(_)) => Ok(false),
Err(e) => Err(e),
}
}
async fn stat(&self, path: &str) -> Result<ObjectMeta, StorageError> {
let obj = self
.client
.get_object(&GetObjectRequest {
bucket: self.bucket.clone(),
object: path.to_string(),
..Default::default()
})
.await
.map_err(|e| {
if e.to_string().contains("404") {
StorageError::NotFound(path.to_string())
} else {
map_err(e)
}
})?;
Ok(ObjectMeta {
size: obj.size.max(0) as u64,
etag: Some(obj.etag),
last_modified: obj.updated.map(|t| t.into()),
})
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::traits::{bytes_stream, read_to_vec};
/// Requires a live GCS bucket with Application Default Credentials
/// configured, and `GCS_BUCKET` set. Run with:
/// `GCS_BUCKET=my-bucket cargo test -p corex-storage --features gcs -- --ignored`.
#[tokio::test]
#[ignore = "requires a live GCS bucket + Application Default Credentials"]
async fn put_get_roundtrip_live() {
let bucket = std::env::var("GCS_BUCKET").expect("set GCS_BUCKET");
let storage = GcsStorage::connect(bucket).await.unwrap();
storage
.put(
"corex-storage-test.txt",
bytes_stream(b"hello gcs".to_vec()),
)
.await
.unwrap();
let data = read_to_vec(storage.get("corex-storage-test.txt").await.unwrap())
.await
.unwrap();
assert_eq!(data, b"hello gcs");
storage.delete("corex-storage-test.txt").await.unwrap();
}
}
+51
View File
@@ -0,0 +1,51 @@
//! # corex-storage
//!
//! A unified storage & file system abstraction so file handling doesn't get
//! locked to one physical storage location.
//!
//! - **Local disk** ([`local::LocalFileStorage`], feature `local`, default).
//! - **S3-compatible object storage** ([`s3::S3Storage`], feature `s3`) —
//! works with Amazon S3 and any S3-compatible endpoint, including MinIO
//! (pass a custom endpoint via [`s3::S3Storage::connect_with_endpoint`]).
//! - **Google Cloud Storage** ([`gcs::GcsStorage`], feature `gcs`).
//!
//! All operations are stream-based (an [`traits::ObjectStream`] is a boxed
//! [`tokio::io::AsyncRead`]), so uploading/downloading a huge object never
//! requires holding it entirely in memory.
//!
//! ## Example
//!
//! ```
//! use corex_storage::{LocalFileStorage, StorageDriver, bytes_stream, read_to_vec};
//! # #[tokio::main] async fn main() {
//! # let dir = tempfile::tempdir().unwrap();
//! let storage = LocalFileStorage::new(dir.path());
//! storage.put("hello.txt", bytes_stream(b"hi".to_vec())).await.unwrap();
//! let data = read_to_vec(storage.get("hello.txt").await.unwrap()).await.unwrap();
//! assert_eq!(data, b"hi");
//! # }
//! ```
pub mod traits;
#[cfg(feature = "local")]
pub mod local;
#[cfg(feature = "s3")]
pub mod s3;
#[cfg(feature = "gcs")]
pub mod gcs;
pub use traits::{
bytes_stream, read_to_vec, ObjectMeta, ObjectStream, StorageDriver, StorageError,
};
#[cfg(feature = "local")]
pub use local::LocalFileStorage;
#[cfg(feature = "s3")]
pub use s3::S3Storage;
#[cfg(feature = "gcs")]
pub use gcs::GcsStorage;
+165
View File
@@ -0,0 +1,165 @@
//! A local-disk [`StorageDriver`] (feature `local`, default).
use std::path::{Component, Path, PathBuf};
use async_trait::async_trait;
use crate::traits::{ObjectMeta, ObjectStream, StorageDriver, StorageError};
/// Stores objects as files under a root directory.
///
/// Paths are always resolved relative to the configured root; `..` path
/// components are rejected to prevent escaping the root.
#[derive(Clone)]
pub struct LocalFileStorage {
root: PathBuf,
}
impl LocalFileStorage {
/// Builds a driver rooted at `root`. The directory is not required to
/// exist yet; it's created lazily on first write.
pub fn new(root: impl Into<PathBuf>) -> Self {
Self { root: root.into() }
}
fn resolve(&self, path: &str) -> Result<PathBuf, StorageError> {
let rel = Path::new(path.trim_start_matches('/'));
for component in rel.components() {
if matches!(component, Component::ParentDir | Component::Prefix(_)) {
return Err(StorageError::InvalidPath(path.to_string()));
}
}
Ok(self.root.join(rel))
}
}
#[async_trait]
impl StorageDriver for LocalFileStorage {
async fn get(&self, path: &str) -> Result<ObjectStream, StorageError> {
let full = self.resolve(path)?;
let file = tokio::fs::File::open(&full).await.map_err(|e| {
if e.kind() == std::io::ErrorKind::NotFound {
StorageError::NotFound(path.to_string())
} else {
StorageError::Io(e.to_string())
}
})?;
Ok(Box::pin(file))
}
async fn put(&self, path: &str, mut data: ObjectStream) -> Result<u64, StorageError> {
let full = self.resolve(path)?;
if let Some(parent) = full.parent() {
tokio::fs::create_dir_all(parent)
.await
.map_err(|e| StorageError::Io(e.to_string()))?;
}
let mut file = tokio::fs::File::create(&full)
.await
.map_err(|e| StorageError::Io(e.to_string()))?;
let written = tokio::io::copy(&mut data, &mut file)
.await
.map_err(|e| StorageError::Io(e.to_string()))?;
Ok(written)
}
async fn delete(&self, path: &str) -> Result<(), StorageError> {
let full = self.resolve(path)?;
match tokio::fs::remove_file(&full).await {
Ok(()) => Ok(()),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
Err(StorageError::NotFound(path.to_string()))
}
Err(e) => Err(StorageError::Io(e.to_string())),
}
}
async fn exists(&self, path: &str) -> Result<bool, StorageError> {
let full = self.resolve(path)?;
Ok(tokio::fs::metadata(&full).await.is_ok())
}
async fn stat(&self, path: &str) -> Result<ObjectMeta, StorageError> {
let full = self.resolve(path)?;
let meta = tokio::fs::metadata(&full).await.map_err(|e| {
if e.kind() == std::io::ErrorKind::NotFound {
StorageError::NotFound(path.to_string())
} else {
StorageError::Io(e.to_string())
}
})?;
Ok(ObjectMeta {
size: meta.len(),
etag: None,
last_modified: meta.modified().ok(),
})
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::traits::{bytes_stream, read_to_vec};
fn driver() -> (LocalFileStorage, tempfile::TempDir) {
let dir = tempfile::tempdir().unwrap();
(LocalFileStorage::new(dir.path()), dir)
}
#[tokio::test]
async fn put_get_roundtrip() {
let (storage, _dir) = driver();
let written = storage
.put("a/b/file.txt", bytes_stream(b"hello world".to_vec()))
.await
.unwrap();
assert_eq!(written, 11);
let data = read_to_vec(storage.get("a/b/file.txt").await.unwrap())
.await
.unwrap();
assert_eq!(data, b"hello world");
}
#[tokio::test]
async fn get_missing_returns_not_found() {
let (storage, _dir) = driver();
// `ObjectStream` isn't `Debug`, so match rather than `unwrap_err()`.
match storage.get("missing.txt").await {
Err(StorageError::NotFound(_)) => {}
other => panic!("expected NotFound, got {:?}", other.is_ok()),
}
}
#[tokio::test]
async fn delete_removes_object() {
let (storage, _dir) = driver();
storage
.put("x.txt", bytes_stream(b"x".to_vec()))
.await
.unwrap();
assert!(storage.exists("x.txt").await.unwrap());
storage.delete("x.txt").await.unwrap();
assert!(!storage.exists("x.txt").await.unwrap());
}
#[tokio::test]
async fn stat_reports_size() {
let (storage, _dir) = driver();
storage
.put("s.bin", bytes_stream(vec![0u8; 1024]))
.await
.unwrap();
let meta = storage.stat("s.bin").await.unwrap();
assert_eq!(meta.size, 1024);
}
#[tokio::test]
async fn path_traversal_is_rejected() {
let (storage, _dir) = driver();
let err = storage
.put("../escape.txt", bytes_stream(b"x".to_vec()))
.await
.unwrap_err();
assert!(matches!(err, StorageError::InvalidPath(_)));
}
}
+334
View File
@@ -0,0 +1,334 @@
//! An S3-compatible [`StorageDriver`] (feature `s3`), via `aws-sdk-s3`.
//!
//! Works against Amazon S3 or any S3-compatible endpoint — pass a custom
//! endpoint via [`S3Storage::connect_with_endpoint`] to target MinIO or
//! another compatible service.
//!
//! [`put`](S3Storage::put) streams the input in fixed-size chunks through
//! S3's multipart upload API rather than buffering the whole object in
//! memory, so large uploads stay within a bounded memory footprint.
use async_trait::async_trait;
use aws_sdk_s3::config::{Credentials, Region};
use aws_sdk_s3::primitives::ByteStream as AwsByteStream;
use aws_sdk_s3::types::{CompletedMultipartUpload, CompletedPart};
use aws_sdk_s3::Client;
use tokio::io::AsyncReadExt;
use crate::traits::{ObjectMeta, ObjectStream, StorageDriver, StorageError};
/// S3's minimum multipart part size (5 MiB), except for the final part.
const MULTIPART_CHUNK_SIZE: usize = 8 * 1024 * 1024;
fn map_sdk_err<E: std::fmt::Debug>(e: E) -> StorageError {
StorageError::Io(format!("{e:?}"))
}
/// A [`StorageDriver`] backed by an S3 (or S3-compatible) bucket.
#[derive(Clone)]
pub struct S3Storage {
client: Client,
bucket: String,
}
impl S3Storage {
/// Connects using the ambient AWS environment/credential chain
/// (`AWS_ACCESS_KEY_ID`, IAM role, profile, etc.) against real S3.
pub async fn connect(bucket: impl Into<String>) -> Self {
let config = aws_config::load_from_env().await;
let client = Client::new(&config);
Self {
client,
bucket: bucket.into(),
}
}
/// Connects to a custom S3-compatible endpoint (e.g. MinIO) with static
/// credentials and path-style addressing.
pub async fn connect_with_endpoint(
bucket: impl Into<String>,
endpoint: &str,
region: &str,
access_key: &str,
secret_key: &str,
) -> Self {
let credentials = Credentials::new(access_key, secret_key, None, None, "corex-storage");
let config = aws_sdk_s3::Config::builder()
.region(Region::new(region.to_string()))
.endpoint_url(endpoint)
.credentials_provider(credentials)
.force_path_style(true)
.behavior_version(aws_sdk_s3::config::BehaviorVersion::latest())
.build();
Self {
client: Client::from_conf(config),
bucket: bucket.into(),
}
}
/// Wraps an already-configured client (advanced use: sharing a client
/// across multiple `S3Storage` instances pointed at different buckets).
pub fn from_client(client: Client, bucket: impl Into<String>) -> Self {
Self {
client,
bucket: bucket.into(),
}
}
}
#[async_trait]
impl StorageDriver for S3Storage {
async fn get(&self, path: &str) -> Result<ObjectStream, StorageError> {
let resp = self
.client
.get_object()
.bucket(&self.bucket)
.key(path)
.send()
.await
.map_err(|e| {
if is_not_found(&e) {
StorageError::NotFound(path.to_string())
} else {
map_sdk_err(e)
}
})?;
Ok(Box::pin(resp.body.into_async_read()))
}
async fn put(&self, path: &str, mut data: ObjectStream) -> Result<u64, StorageError> {
// Read the first chunk to decide between a simple `PutObject` (small
// objects) and a streamed multipart upload (anything larger).
let mut first_chunk = vec![0u8; MULTIPART_CHUNK_SIZE];
let mut filled = 0usize;
while filled < first_chunk.len() {
let n = data
.read(&mut first_chunk[filled..])
.await
.map_err(|e| StorageError::Io(e.to_string()))?;
if n == 0 {
break;
}
filled += n;
}
first_chunk.truncate(filled);
if filled < MULTIPART_CHUNK_SIZE {
// The whole object fit in one chunk: a single PutObject suffices.
let len = first_chunk.len() as u64;
self.client
.put_object()
.bucket(&self.bucket)
.key(path)
.body(AwsByteStream::from(first_chunk))
.send()
.await
.map_err(map_sdk_err)?;
return Ok(len);
}
// Larger objects: stream the rest through multipart upload.
let create = self
.client
.create_multipart_upload()
.bucket(&self.bucket)
.key(path)
.send()
.await
.map_err(map_sdk_err)?;
let upload_id = create
.upload_id()
.ok_or_else(|| StorageError::Io("S3 did not return an upload id".to_string()))?;
let result = upload_parts(
&self.client,
&self.bucket,
path,
upload_id,
first_chunk,
&mut data,
)
.await;
match result {
Ok((total, parts)) => {
self.client
.complete_multipart_upload()
.bucket(&self.bucket)
.key(path)
.upload_id(upload_id)
.multipart_upload(
CompletedMultipartUpload::builder()
.set_parts(Some(parts))
.build(),
)
.send()
.await
.map_err(map_sdk_err)?;
Ok(total)
}
Err(e) => {
let _ = self
.client
.abort_multipart_upload()
.bucket(&self.bucket)
.key(path)
.upload_id(upload_id)
.send()
.await;
Err(e)
}
}
}
async fn delete(&self, path: &str) -> Result<(), StorageError> {
self.client
.delete_object()
.bucket(&self.bucket)
.key(path)
.send()
.await
.map_err(map_sdk_err)?;
Ok(())
}
async fn exists(&self, path: &str) -> Result<bool, StorageError> {
match self
.client
.head_object()
.bucket(&self.bucket)
.key(path)
.send()
.await
{
Ok(_) => Ok(true),
Err(e) if is_not_found(&e) => Ok(false),
Err(e) => Err(map_sdk_err(e)),
}
}
async fn stat(&self, path: &str) -> Result<ObjectMeta, StorageError> {
let resp = self
.client
.head_object()
.bucket(&self.bucket)
.key(path)
.send()
.await
.map_err(|e| {
if is_not_found(&e) {
StorageError::NotFound(path.to_string())
} else {
map_sdk_err(e)
}
})?;
Ok(ObjectMeta {
size: resp.content_length().unwrap_or(0).max(0) as u64,
etag: resp.e_tag().map(|s| s.to_string()),
last_modified: resp
.last_modified()
.and_then(|d| d.to_owned().try_into().ok()),
})
}
}
/// Uploads `first_chunk` as part 1, then continues reading `data` in
/// [`MULTIPART_CHUNK_SIZE`] chunks until exhausted, uploading each as a part.
async fn upload_parts(
client: &Client,
bucket: &str,
key: &str,
upload_id: &str,
first_chunk: Vec<u8>,
data: &mut ObjectStream,
) -> Result<(u64, Vec<CompletedPart>), StorageError> {
let mut total = 0u64;
let mut parts = Vec::new();
let mut part_number = 1i32;
let mut chunk = first_chunk;
loop {
total += chunk.len() as u64;
let resp = client
.upload_part()
.bucket(bucket)
.key(key)
.upload_id(upload_id)
.part_number(part_number)
.body(AwsByteStream::from(chunk))
.send()
.await
.map_err(map_sdk_err)?;
parts.push(
CompletedPart::builder()
.e_tag(resp.e_tag().unwrap_or_default())
.part_number(part_number)
.build(),
);
part_number += 1;
let mut next = vec![0u8; MULTIPART_CHUNK_SIZE];
let mut filled = 0usize;
while filled < next.len() {
let n = data
.read(&mut next[filled..])
.await
.map_err(|e| StorageError::Io(e.to_string()))?;
if n == 0 {
break;
}
filled += n;
}
next.truncate(filled);
if next.is_empty() {
break;
}
chunk = next;
}
Ok((total, parts))
}
/// Best-effort detection of a "not found" S3 error across the SDK's error
/// variants (`GetObject`/`HeadObject` surface this differently).
fn is_not_found<E: std::fmt::Debug>(e: &E) -> bool {
format!("{e:?}").contains("NotFound") || format!("{e:?}").contains("404")
}
#[cfg(test)]
mod tests {
use super::*;
use crate::traits::{bytes_stream, read_to_vec};
/// Requires a live S3-compatible endpoint (e.g. MinIO) configured via
/// `S3_ENDPOINT`, `S3_BUCKET`, `S3_ACCESS_KEY`, `S3_SECRET_KEY`. Run with:
/// `S3_ENDPOINT=http://127.0.0.1:9000 S3_BUCKET=test S3_ACCESS_KEY=... S3_SECRET_KEY=... \
/// cargo test -p corex-storage --features s3 -- --ignored`.
#[tokio::test]
#[ignore = "requires a live S3-compatible endpoint (S3_ENDPOINT, S3_BUCKET, ...)"]
async fn put_get_roundtrip_live() {
let endpoint = std::env::var("S3_ENDPOINT").expect("set S3_ENDPOINT");
let bucket = std::env::var("S3_BUCKET").expect("set S3_BUCKET");
let access_key = std::env::var("S3_ACCESS_KEY").expect("set S3_ACCESS_KEY");
let secret_key = std::env::var("S3_SECRET_KEY").expect("set S3_SECRET_KEY");
let storage = S3Storage::connect_with_endpoint(
bucket,
&endpoint,
"us-east-1",
&access_key,
&secret_key,
)
.await;
storage
.put("corex-storage-test.txt", bytes_stream(b"hello s3".to_vec()))
.await
.unwrap();
let data = read_to_vec(storage.get("corex-storage-test.txt").await.unwrap())
.await
.unwrap();
assert_eq!(data, b"hello s3");
storage.delete("corex-storage-test.txt").await.unwrap();
}
}
+109
View File
@@ -0,0 +1,109 @@
//! The core [`StorageDriver`] trait and supporting types.
use std::pin::Pin;
use std::task::{Context, Poll};
use std::time::SystemTime;
use async_trait::async_trait;
use tokio::io::{AsyncRead, ReadBuf};
/// Errors returned by storage operations.
#[derive(Debug)]
pub enum StorageError {
/// The object does not exist.
NotFound(String),
/// The path was rejected (e.g. attempted directory traversal).
InvalidPath(String),
/// An I/O or backend transport error occurred.
Io(String),
}
impl std::fmt::Display for StorageError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::NotFound(p) => write!(f, "object not found: {p}"),
Self::InvalidPath(p) => write!(f, "invalid path: {p}"),
Self::Io(s) => write!(f, "storage io error: {s}"),
}
}
}
impl std::error::Error for StorageError {}
/// A boxed, unpin, streaming byte source used for both `get` results and
/// `put` input, so large objects never need to be buffered in memory whole.
pub type ObjectStream = Pin<Box<dyn AsyncRead + Send + Unpin>>;
/// Metadata about a stored object.
#[derive(Debug, Clone)]
pub struct ObjectMeta {
/// Size in bytes.
pub size: u64,
/// A backend-provided content hash/version tag, if available.
pub etag: Option<String>,
/// Last-modified timestamp, if available.
pub last_modified: Option<SystemTime>,
}
/// A unified interface for storing, reading, and deleting files regardless of
/// the physical backend (local disk, S3/MinIO, Google Cloud Storage, ...).
///
/// All operations are stream-based: [`StorageDriver::put`] takes an
/// [`ObjectStream`] and [`StorageDriver::get`] returns one, so a caller
/// forwarding a large upload/download never has to hold the whole object in
/// memory.
#[async_trait]
pub trait StorageDriver: Send + Sync {
/// Opens `path` for streaming read.
async fn get(&self, path: &str) -> Result<ObjectStream, StorageError>;
/// Writes `data` to `path`, returning the number of bytes written.
async fn put(&self, path: &str, data: ObjectStream) -> Result<u64, StorageError>;
/// Removes the object at `path`.
async fn delete(&self, path: &str) -> Result<(), StorageError>;
/// Whether an object exists at `path`.
async fn exists(&self, path: &str) -> Result<bool, StorageError>;
/// Fetches metadata about the object at `path` without downloading it.
async fn stat(&self, path: &str) -> Result<ObjectMeta, StorageError>;
}
/// A minimal, dependency-free in-memory [`AsyncRead`] over an owned buffer.
struct VecCursor {
data: Vec<u8>,
pos: usize,
}
impl AsyncRead for VecCursor {
fn poll_read(
mut self: Pin<&mut Self>,
_cx: &mut Context<'_>,
buf: &mut ReadBuf<'_>,
) -> Poll<std::io::Result<()>> {
let remaining = &self.data[self.pos..];
let n = remaining.len().min(buf.remaining());
buf.put_slice(&remaining[..n]);
self.pos += n;
Poll::Ready(Ok(()))
}
}
/// Convenience: wraps an owned `Vec<u8>` as an [`ObjectStream`] for
/// [`StorageDriver::put`].
pub fn bytes_stream(data: Vec<u8>) -> ObjectStream {
Box::pin(VecCursor { data, pos: 0 })
}
/// Convenience: drains an [`ObjectStream`] into an owned `Vec<u8>` (defeats
/// the point of streaming for huge objects — intended for tests/small files).
pub async fn read_to_vec(mut stream: ObjectStream) -> Result<Vec<u8>, StorageError> {
use tokio::io::AsyncReadExt;
let mut buf = Vec::new();
stream
.read_to_end(&mut buf)
.await
.map_err(|e| StorageError::Io(e.to_string()))?;
Ok(buf)
}