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
11 changes: 11 additions & 0 deletions docs-mintlify/reference/configuration/environment-variables.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -1953,6 +1953,17 @@ The name of a bucket in AWS S3. Required when using AWS S3.
| -------------------------------------- | ---------------------- | --------------------- |
| [A valid AWS region][aws-docs-regions] | N/A | N/A |

## `CUBESTORE_S3_SSE`

When set, Cube Store sends the `x-amz-server-side-encryption` header with this
value on every AWS S3 request. Optional. Useful when an AWS Organizations
service control policy denies `s3:PutObject` requests that lack this header,
even if the bucket has default encryption enabled.

| Possible Values | Default in Development | Default in Production |
| --------------------------------- | ---------------------- | --------------------- |
| `AES256`, `aws:kms`, `aws:kms:dsse` | N/A | N/A |

## `CUBESTORE_S3_SUB_PATH`

The path in a AWS S3 bucket to store pre-aggregations. Optional.
Expand Down
12 changes: 6 additions & 6 deletions packages/cubejs-backend-native/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

12 changes: 6 additions & 6 deletions rust/cubesql/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion rust/cubesql/cubesql/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ homepage = "https://cube.dev"

[dependencies]
arc-swap = "1"
datafusion = { git = 'https://github.com/cube-js/arrow-datafusion.git', rev = "1c073dfa64c10949df15064aa7eb71ce513e725a", default-features = false, features = [
datafusion = { git = 'https://github.com/cube-js/arrow-datafusion.git', rev = "ed911e3d05215a4c97e6c977d5c3bd7f9cc33b8b", default-features = false, features = [
"regex_expressions",
"unicode_expressions",
] }
Expand Down
55 changes: 55 additions & 0 deletions rust/cubesql/cubesql/src/compile/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19006,6 +19006,61 @@ LIMIT {{ limit }}{% endif %}"#.to_string(),
.contains("must appear in the GROUP BY clause or be used in an aggregate function"));
}

#[tokio::test]
async fn test_cte_order_by_unprojected_column() {
init_testing_logger();

// The CTE's ORDER BY references `cnt`, which is not in the CTE's projection;
// the sort column must be resolved through the CTE's aliased projection
// during DF post-processing planning
let query_plan = convert_select_to_query_plan(
// language=PostgreSQL
r#"
WITH t1 AS (
SELECT
customer_gender,
taxful_total_price AS price,
MEASURE(count) AS cnt
FROM KibanaSampleDataEcommerce
GROUP BY 1, 2
),
t2 AS (
SELECT customer_gender, price
FROM t1
ORDER BY customer_gender, cnt DESC
)
SELECT * FROM t2
LIMIT 100
"#
.to_string(),
DatabaseProtocol::PostgreSQL,
)
.await;

let logical_plan = query_plan.as_logical_plan();
let request = logical_plan.find_cube_scan().request;
assert_eq!(
request.dimensions,
Some(vec![
"KibanaSampleDataEcommerce.customer_gender".to_string(),
"KibanaSampleDataEcommerce.taxful_total_price".to_string(),
])
);
assert_eq!(
request.order,
Some(vec![
vec![
"KibanaSampleDataEcommerce.customer_gender".to_string(),
"asc".to_string(),
],
vec![
"KibanaSampleDataEcommerce.count".to_string(),
"desc".to_string(),
],
])
);
}

#[tokio::test]
async fn test_set_cache_mode() -> Result<(), CubeError> {
if !Rewriter::sql_push_down_enabled() {
Expand Down
113 changes: 111 additions & 2 deletions rust/cubestore/cubestore/src/remotefs/s3.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,10 @@ pub struct S3RemoteFs {
/// STS AssumeRoleWithWebIdentity with the JWT inside it.
web_identity_token_file: Option<String>,
web_identity_role_arn: Option<String>,
/// When set, every request carries `x-amz-server-side-encryption` with this
/// value. Some AWS Organizations SCPs deny `s3:PutObject` unless the header
/// is present, even when the bucket has default encryption.
server_side_encryption: Option<String>,
}

impl fmt::Debug for S3RemoteFs {
Expand Down Expand Up @@ -89,21 +93,61 @@ impl S3RemoteFs {
let region = region.parse::<Region>().map_err(|e| {
CubeError::internal(format!("Failed to parse Region '{}': {}", region, e))
})?;
let bucket = Bucket::new(&bucket_name, region.clone(), credentials)?;
let server_side_encryption = server_side_encryption_from_env()?;
let bucket = new_bucket(
&bucket_name,
region.clone(),
credentials,
&server_side_encryption,
)?;
let fs = Arc::new(Self {
dir,
bucket: arc_swap::ArcSwap::new(Arc::new(bucket)),
sub_path,
delete_mut: Mutex::new(()),
web_identity_token_file: token_file,
web_identity_role_arn: role_arn,
server_side_encryption,
});
spawn_creds_refresh_loop(access_key, secret_key, bucket_name, region, &fs);

Ok(fs)
}
}

/// Values S3 accepts in `x-amz-server-side-encryption`.
const ALLOWED_SSE_VALUES: &[&str] = &["AES256", "aws:kms", "aws:kms:dsse"];

fn server_side_encryption_from_env() -> Result<Option<String>, CubeError> {
parse_sse_value(env::var("CUBESTORE_S3_SSE").ok())
}

fn parse_sse_value(value: Option<String>) -> Result<Option<String>, CubeError> {
match value {
None => Ok(None),
Some(v) if v.is_empty() => Ok(None),
Some(v) if ALLOWED_SSE_VALUES.contains(&v.as_str()) => Ok(Some(v)),
Some(v) => Err(CubeError::user(format!(
"Invalid CUBESTORE_S3_SSE value '{}'. Expected one of: {}",
v,
ALLOWED_SSE_VALUES.join(", ")
))),
}
}

fn new_bucket(
bucket_name: &str,
region: Region,
credentials: Credentials,
server_side_encryption: &Option<String>,
) -> Result<Bucket, CubeError> {
let mut bucket = Bucket::new(bucket_name, region, credentials)?;
if let Some(sse) = server_side_encryption {
bucket.add_header("x-amz-server-side-encryption", sse);
}
Ok(bucket)
}

fn spawn_creds_refresh_loop(
access_key: Option<String>,
secret_key: Option<String>,
Expand All @@ -113,6 +157,7 @@ fn spawn_creds_refresh_loop(
) {
let token_file = fs.web_identity_token_file.clone();
let role_arn = fs.web_identity_role_arn.clone();
let server_side_encryption = fs.server_side_encryption.clone();
let is_web_identity = token_file.is_some() && role_arn.is_some();

// Web identity STS credentials expire in ~1 hour, so poll the token file
Expand Down Expand Up @@ -186,7 +231,7 @@ fn spawn_creds_refresh_loop(
continue;
}
};
let b = match Bucket::new(&bucket_name, region.clone(), c) {
let b = match new_bucket(&bucket_name, region.clone(), c, &server_side_encryption) {
Ok(b) => b,
Err(e) => {
log::error!("Failed to refresh S3 credentials: {}", e);
Expand Down Expand Up @@ -513,3 +558,67 @@ impl S3RemoteFs {
)
}
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn parse_sse_value_accepts_allowed_values() {
assert_eq!(parse_sse_value(None).unwrap(), None);
assert_eq!(parse_sse_value(Some("".to_string())).unwrap(), None);
assert_eq!(
parse_sse_value(Some("AES256".to_string())).unwrap(),
Some("AES256".to_string())
);
assert_eq!(
parse_sse_value(Some("aws:kms".to_string())).unwrap(),
Some("aws:kms".to_string())
);
assert_eq!(
parse_sse_value(Some("aws:kms:dsse".to_string())).unwrap(),
Some("aws:kms:dsse".to_string())
);
}

#[test]
fn parse_sse_value_rejects_unknown_values() {
assert!(parse_sse_value(Some("aes256".to_string())).is_err());
assert!(parse_sse_value(Some("true".to_string())).is_err());
}

#[test]
fn new_bucket_applies_sse_header() {
let credentials = Credentials::new(Some("key"), Some("secret"), None, None, None).unwrap();
let bucket = new_bucket(
"test-bucket",
"us-east-1".parse().unwrap(),
credentials,
&Some("AES256".to_string()),
)
.unwrap();
assert_eq!(
bucket
.extra_headers()
.get("x-amz-server-side-encryption")
.map(|v| v.to_str().unwrap()),
Some("AES256")
);
}

#[test]
fn new_bucket_without_sse_has_no_header() {
let credentials = Credentials::new(Some("key"), Some("secret"), None, None, None).unwrap();
let bucket = new_bucket(
"test-bucket",
"us-east-1".parse().unwrap(),
credentials,
&None,
)
.unwrap();
assert!(bucket
.extra_headers()
.get("x-amz-server-side-encryption")
.is_none());
}
}
Loading