use std::fmt;
use std::sync::Arc;
use bytes::Bytes;
use futures::StreamExt;
use object_store::aws::{AmazonS3Builder, AwsCredential};
use object_store::local::LocalFileSystem;
use object_store::path::Path as StorePath;
use object_store::{
Attribute, AttributeValue, Attributes, CredentialProvider, ObjectStore, ObjectStoreExt,
PutOptions, PutPayload,
};
use super::error::DrainError;
use super::uri::DestinationUri;
pub const LIST_LIMIT: usize = 10_000;
#[derive(Debug, Clone, Default, PartialEq, Eq)]
#[non_exhaustive]
pub struct PutMeta {
pub content_type: Option<String>,
pub content_encoding: Option<String>,
}
impl PutMeta {
pub fn gzipped_text() -> Self {
Self {
content_type: Some("text/plain".to_string()),
content_encoding: Some("gzip".to_string()),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct ObjectMeta {
pub key: String,
pub size: u64,
pub last_modified_unix: i64,
}
#[async_trait::async_trait]
pub trait LogDestination: Send + Sync + fmt::Debug {
async fn put(&self, key: &str, body: Bytes, meta: PutMeta) -> Result<(), DrainError>;
async fn head(&self, key: &str) -> Result<Option<ObjectMeta>, DrainError>;
async fn get(&self, key: &str) -> Result<Option<Bytes>, DrainError>;
async fn list(&self, prefix: &str) -> Result<Vec<ObjectMeta>, DrainError>;
}
pub struct ObjectStoreDestination {
store: Arc<dyn ObjectStore>,
prefix: String,
label: String,
supports_attributes: bool,
}
impl fmt::Debug for ObjectStoreDestination {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("ObjectStoreDestination")
.field("destination", &self.label)
.finish()
}
}
impl ObjectStoreDestination {
pub async fn connect(uri: &DestinationUri) -> Result<Self, DrainError> {
match uri {
DestinationUri::File { path } => {
std::fs::create_dir_all(path).map_err(|source| DrainError::Io {
path: path.clone(),
source,
})?;
let store = LocalFileSystem::new_with_prefix(path).map_err(|e| DrainError::Io {
path: path.clone(),
source: std::io::Error::other(e),
})?;
Ok(Self {
store: Arc::new(store),
prefix: String::new(),
label: format!("file://{}", path.display()),
supports_attributes: false,
})
}
DestinationUri::S3 {
bucket,
prefix,
region,
} => Self::connect_s3(bucket, prefix, region.as_deref()).await,
}
}
async fn connect_s3(
bucket: &str,
prefix: &str,
region_override: Option<&str>,
) -> Result<Self, DrainError> {
let label = format!("s3://{bucket}/{prefix}");
let sdk_config = aws_config::defaults(aws_config::BehaviorVersion::latest())
.load()
.await;
let region = region_override
.map(str::to_string)
.or_else(|| sdk_config.region().map(|r| r.as_ref().to_string()))
.ok_or_else(|| DrainError::Credentials {
uri: label.clone(),
source: "no AWS region: none in the URI's `?region=`, none from the \
credential chain (set AWS_REGION or add `?region=` to the URI)"
.into(),
})?;
let provider =
sdk_config
.credentials_provider()
.ok_or_else(|| DrainError::Credentials {
uri: label.clone(),
source: "the AWS default provider chain supplied no credentials provider"
.into(),
})?;
let store = AmazonS3Builder::new()
.with_bucket_name(bucket)
.with_region(®ion)
.with_credentials(Arc::new(AwsChainCredentials { inner: provider }))
.build()
.map_err(|source| DrainError::Transport {
op: "connect",
key: label.clone(),
source,
})?;
Ok(Self {
store: Arc::new(store),
prefix: prefix.to_string(),
label,
supports_attributes: true,
})
}
fn absolute(&self, key: &str) -> String {
if self.prefix.is_empty() {
key.to_string()
} else {
format!("{}/{}", self.prefix, key)
}
}
fn convert(meta: object_store::ObjectMeta) -> ObjectMeta {
ObjectMeta {
key: meta.location.as_ref().to_string(),
size: meta.size,
last_modified_unix: meta.last_modified.timestamp(),
}
}
}
#[async_trait::async_trait]
impl LogDestination for ObjectStoreDestination {
async fn put(&self, key: &str, body: Bytes, meta: PutMeta) -> Result<(), DrainError> {
let absolute = self.absolute(key);
let path = StorePath::from(absolute.as_str());
let mut attributes = Attributes::new();
if self.supports_attributes {
if let Some(ct) = meta.content_type {
attributes.insert(Attribute::ContentType, AttributeValue::from(ct));
}
if let Some(ce) = meta.content_encoding {
attributes.insert(Attribute::ContentEncoding, AttributeValue::from(ce));
}
}
let options = PutOptions {
attributes,
..PutOptions::default()
};
self.store
.put_opts(&path, PutPayload::from_bytes(body), options)
.await
.map_err(|source| DrainError::Transport {
op: "put",
key: absolute,
source,
})?;
Ok(())
}
async fn head(&self, key: &str) -> Result<Option<ObjectMeta>, DrainError> {
let absolute = self.absolute(key);
let path = StorePath::from(absolute.as_str());
match self.store.head(&path).await {
Ok(meta) => Ok(Some(Self::convert(meta))),
Err(object_store::Error::NotFound { .. }) => Ok(None),
Err(source) => Err(DrainError::Transport {
op: "head",
key: absolute,
source,
}),
}
}
async fn get(&self, key: &str) -> Result<Option<Bytes>, DrainError> {
let absolute = self.absolute(key);
let path = StorePath::from(absolute.as_str());
let result = match self.store.get(&path).await {
Ok(result) => result,
Err(object_store::Error::NotFound { .. }) => return Ok(None),
Err(source) => {
return Err(DrainError::Transport {
op: "get",
key: absolute,
source,
});
}
};
let bytes = result
.bytes()
.await
.map_err(|source| DrainError::Transport {
op: "get",
key: absolute,
source,
})?;
Ok(Some(bytes))
}
async fn list(&self, prefix: &str) -> Result<Vec<ObjectMeta>, DrainError> {
let absolute = self.absolute(prefix);
let path = StorePath::from(absolute.as_str());
let mut stream = self.store.list(Some(&path));
let mut out = Vec::new();
while let Some(next) = stream.next().await {
let meta = next.map_err(|source| DrainError::Transport {
op: "list",
key: absolute.clone(),
source,
})?;
out.push(Self::convert(meta));
if out.len() >= LIST_LIMIT {
tracing::warn!(
prefix = %absolute,
limit = LIST_LIMIT,
"log-drain list hit its entry cap; results truncated"
);
break;
}
}
Ok(out)
}
}
#[derive(Debug)]
struct AwsChainCredentials {
inner: aws_types::sdk_config::SharedCredentialsProvider,
}
#[async_trait::async_trait]
impl CredentialProvider for AwsChainCredentials {
type Credential = AwsCredential;
async fn get_credential(&self) -> object_store::Result<Arc<AwsCredential>> {
use aws_credential_types::provider::ProvideCredentials;
let creds = self.inner.provide_credentials().await.map_err(|e| {
object_store::Error::Unauthenticated {
path: "aws-credential-chain".to_string(),
source: Box::new(e),
}
})?;
Ok(Arc::new(AwsCredential {
key_id: creds.access_key_id().to_string(),
secret_key: creds.secret_access_key().to_string(),
token: creds.session_token().map(str::to_string),
}))
}
}