Skip to content

Commit e732da4

Browse files
committed
add spooling config to STS and some unit tests
1 parent 32dfb3a commit e732da4

5 files changed

Lines changed: 220 additions & 56 deletions

File tree

deploy/helm/trino-operator/crds/crds.yaml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -115,8 +115,8 @@ spec:
115115
configOverrides:
116116
additionalProperties:
117117
type: string
118-
default: {}
119118
description: The `configOverrides` allow overriding arbitrary client protocol properties.
119+
nullable: true
120120
type: object
121121
enabled:
122122
description: Enable spooling protocol.

rust/operator-binary/src/command.rs

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,7 @@ use crate::{
1414
CONFIG_DIR_NAME, Container, LOG_PROPERTIES, RW_CONFIG_DIR_NAME, STACKABLE_CLIENT_TLS_DIR,
1515
STACKABLE_INTERNAL_TLS_DIR, STACKABLE_MOUNT_INTERNAL_TLS_DIR,
1616
STACKABLE_MOUNT_SERVER_TLS_DIR, STACKABLE_SERVER_TLS_DIR, STACKABLE_TLS_STORE_PASSWORD,
17-
SYSTEM_TRUST_STORE, SYSTEM_TRUST_STORE_PASSWORD, TrinoRole,
17+
SYSTEM_TRUST_STORE, SYSTEM_TRUST_STORE_PASSWORD, TrinoRole, client_protocol,
1818
fault_tolerant_execution::ResolvedFaultTolerantExecutionConfig, v1alpha1,
1919
},
2020
};
@@ -24,6 +24,7 @@ pub fn container_prepare_args(
2424
catalogs: &[CatalogConfig],
2525
merged_config: &v1alpha1::TrinoConfig,
2626
resolved_fte_config: &Option<ResolvedFaultTolerantExecutionConfig>,
27+
resolved_spooling_config: &Option<client_protocol::ResolvedSpoolingProtocolConfig>,
2728
) -> Vec<String> {
2829
let mut args = vec![];
2930

@@ -85,13 +86,19 @@ pub fn container_prepare_args(
8586
args.extend_from_slice(&resolved_fte.init_container_extra_start_commands);
8687
}
8788

89+
// Add the commands that are needed for the client spooling protocol (e.g., TLS certificates for S3)
90+
if let Some(resolved_spooling) = resolved_spooling_config {
91+
args.extend_from_slice(&resolved_spooling.init_container_extra_start_commands);
92+
}
93+
8894
args
8995
}
9096

9197
pub fn container_trino_args(
9298
authentication_config: &TrinoAuthenticationConfig,
9399
catalogs: &[CatalogConfig],
94100
resolved_fte_config: &Option<ResolvedFaultTolerantExecutionConfig>,
101+
resolved_spooling_config: &Option<client_protocol::ResolvedSpoolingProtocolConfig>,
95102
) -> Vec<String> {
96103
let mut args = vec![
97104
// copy config files to a writeable empty folder
@@ -126,6 +133,13 @@ pub fn container_trino_args(
126133
}
127134
}
128135

136+
// Add client spooling environment variables from files
137+
if let Some(resolved_spooling) = resolved_spooling_config {
138+
for (env_name, file) in &resolved_spooling.load_env_from_files {
139+
args.push(format!("export {env_name}=\"$(cat {file})\""));
140+
}
141+
}
142+
129143
args.push("set -x".to_string());
130144

131145
// Start command

rust/operator-binary/src/controller.rs

Lines changed: 76 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -81,11 +81,11 @@ use crate::{
8181
command, config,
8282
crd::{
8383
ACCESS_CONTROL_PROPERTIES, APP_NAME, CONFIG_DIR_NAME, CONFIG_PROPERTIES, Container,
84-
DISCOVERY_URI, ENV_INTERNAL_SECRET, EXCHANGE_MANAGER_PROPERTIES, HTTP_PORT, HTTP_PORT_NAME,
85-
HTTPS_PORT, HTTPS_PORT_NAME, JVM_CONFIG, JVM_SECURITY_PROPERTIES, LOG_PROPERTIES,
86-
MAX_TRINO_LOG_FILES_SIZE, METRICS_PORT, METRICS_PORT_NAME, NODE_PROPERTIES,
87-
RW_CONFIG_DIR_NAME, SPOOLING_MANAGER_PROPERTIES, STACKABLE_CLIENT_TLS_DIR,
88-
STACKABLE_INTERNAL_TLS_DIR, STACKABLE_MOUNT_INTERNAL_TLS_DIR,
84+
DISCOVERY_URI, ENV_INTERNAL_SECRET, ENV_SPOOLING_SECRET, EXCHANGE_MANAGER_PROPERTIES,
85+
HTTP_PORT, HTTP_PORT_NAME, HTTPS_PORT, HTTPS_PORT_NAME, JVM_CONFIG,
86+
JVM_SECURITY_PROPERTIES, LOG_PROPERTIES, MAX_TRINO_LOG_FILES_SIZE, METRICS_PORT,
87+
METRICS_PORT_NAME, NODE_PROPERTIES, RW_CONFIG_DIR_NAME, SPOOLING_MANAGER_PROPERTIES,
88+
STACKABLE_CLIENT_TLS_DIR, STACKABLE_INTERNAL_TLS_DIR, STACKABLE_MOUNT_INTERNAL_TLS_DIR,
8989
STACKABLE_MOUNT_SERVER_TLS_DIR, STACKABLE_SERVER_TLS_DIR, TrinoRole,
9090
authentication::resolve_authentication_classes,
9191
catalog, client_protocol,
@@ -466,9 +466,9 @@ pub async fn reconcile_trino(
466466
// Resolve client spooling protocol configuration with S3 connections if needed
467467
let resolved_spooling_config = match trino.spec.cluster_config.client_spooling_protocol.as_ref()
468468
{
469-
Some(client_protocol_config) => Some(
469+
Some(spooling_config) => Some(
470470
client_protocol::ResolvedSpoolingProtocolConfig::from_config(
471-
client_protocol_config,
471+
spooling_config,
472472
Some(client),
473473
&namespace,
474474
)
@@ -524,7 +524,21 @@ pub async fn reconcile_trino(
524524
None => None,
525525
};
526526

527-
create_shared_internal_secret(trino, client).await?;
527+
create_random_secret(
528+
&shared_internal_secret_name(trino),
529+
ENV_INTERNAL_SECRET,
530+
trino,
531+
client,
532+
)
533+
.await?;
534+
535+
create_random_secret(
536+
&shared_spooling_secret_name(trino),
537+
ENV_SPOOLING_SECRET,
538+
trino,
539+
client,
540+
)
541+
.await?;
528542

529543
let mut sts_cond_builder = StatefulSetConditionBuilder::default();
530544

@@ -600,6 +614,7 @@ pub async fn reconcile_trino(
600614
&catalogs,
601615
&rbac_sa.name_any(),
602616
&resolved_fte_config,
617+
&resolved_spooling_config,
603618
)?;
604619

605620
cluster_resources
@@ -709,7 +724,7 @@ fn build_rolegroup_config_map(
709724
trino_opa_config: &Option<TrinoOpaConfig>,
710725
cluster_info: &KubernetesClusterInfo,
711726
resolved_fte_config: &Option<ResolvedFaultTolerantExecutionConfig>,
712-
resolved_spooling_protocol_config: &Option<client_protocol::ResolvedSpoolingProtocolConfig>,
727+
resolved_spooling_config: &Option<client_protocol::ResolvedSpoolingProtocolConfig>,
713728
) -> Result<ConfigMap> {
714729
let mut cm_conf_data = BTreeMap::new();
715730

@@ -862,20 +877,17 @@ fn build_rolegroup_config_map(
862877
}
863878

864879
// Add client protocol properties (especially spooling properties)
865-
if let Some(resolved_client_protocol) = resolved_spooling_protocol_config {
866-
if resolved_client_protocol.is_enabled() {
867-
let spooling_props_with_options: BTreeMap<String, Option<String>> =
868-
resolved_client_protocol
869-
.spooling_manager_properties
870-
.iter()
871-
.map(|(k, v)| (k.clone(), Some(v.clone())))
872-
.collect();
873-
cm_conf_data.insert(
874-
SPOOLING_MANAGER_PROPERTIES.to_string(),
875-
to_java_properties_string(spooling_props_with_options.iter())
876-
.with_context(|_| FailedToWriteJavaPropertiesSnafu)?,
877-
);
878-
}
880+
if let Some(spooling_config) = resolved_spooling_config {
881+
let spooling_props_with_options: BTreeMap<String, Option<String>> = spooling_config
882+
.spooling_manager_properties
883+
.iter()
884+
.map(|(k, v)| (k.clone(), Some(v.clone())))
885+
.collect();
886+
cm_conf_data.insert(
887+
SPOOLING_MANAGER_PROPERTIES.to_string(),
888+
to_java_properties_string(spooling_props_with_options.iter())
889+
.with_context(|_| FailedToWriteJavaPropertiesSnafu)?,
890+
);
879891
}
880892

881893
let jvm_sec_props: BTreeMap<String, Option<String>> = config
@@ -987,6 +999,7 @@ fn build_rolegroup_statefulset(
987999
catalogs: &[CatalogConfig],
9881000
sa_name: &str,
9891001
resolved_fte_config: &Option<ResolvedFaultTolerantExecutionConfig>,
1002+
resolved_spooling_config: &Option<client_protocol::ResolvedSpoolingProtocolConfig>,
9901003
) -> Result<StatefulSet> {
9911004
let role = trino
9921005
.role(trino_role)
@@ -1014,7 +1027,7 @@ fn build_rolegroup_statefulset(
10141027
// additional authentication env vars
10151028
let mut env = trino_authentication_config.env_vars(trino_role, &Container::Trino);
10161029

1017-
let secret_name = build_shared_internal_secret_name(trino);
1030+
let secret_name = shared_internal_secret_name(trino);
10181031
env.push(env_var_from_secret(&secret_name, None, ENV_INTERNAL_SECRET));
10191032

10201033
trino_authentication_config
@@ -1078,6 +1091,7 @@ fn build_rolegroup_statefulset(
10781091
catalogs,
10791092
&requested_secret_lifetime,
10801093
resolved_fte_config,
1094+
resolved_spooling_config,
10811095
)?;
10821096

10831097
let mut prepare_args = vec![];
@@ -1097,6 +1111,7 @@ fn build_rolegroup_statefulset(
10971111
catalogs,
10981112
merged_config,
10991113
resolved_fte_config,
1114+
resolved_spooling_config,
11001115
));
11011116

11021117
prepare_args
@@ -1165,6 +1180,7 @@ fn build_rolegroup_statefulset(
11651180
trino_authentication_config,
11661181
catalogs,
11671182
resolved_fte_config,
1183+
resolved_spooling_config,
11681184
)
11691185
.join("\n"),
11701186
])
@@ -1447,11 +1463,27 @@ fn build_recommended_labels<'a>(
14471463
}
14481464
}
14491465

1450-
async fn create_shared_internal_secret(
1466+
async fn create_random_secret(
1467+
secret_name: &str,
1468+
secret_key: &str,
14511469
trino: &v1alpha1::TrinoCluster,
14521470
client: &Client,
14531471
) -> Result<()> {
1454-
let secret = build_shared_internal_secret(trino)?;
1472+
let mut internal_secret = BTreeMap::new();
1473+
internal_secret.insert(secret_key.to_string(), get_random_base64());
1474+
1475+
let secret = Secret {
1476+
immutable: Some(true),
1477+
metadata: ObjectMetaBuilder::new()
1478+
.name(secret_name)
1479+
.namespace_opt(trino.namespace())
1480+
.ownerreference_from_resource(trino, None, Some(true))
1481+
.context(ObjectMissingMetadataForOwnerRefSnafu)?
1482+
.build(),
1483+
string_data: Some(internal_secret),
1484+
..Secret::default()
1485+
};
1486+
14551487
if client
14561488
.get_opt::<Secret>(
14571489
&secret.name_any(),
@@ -1473,25 +1505,12 @@ async fn create_shared_internal_secret(
14731505
Ok(())
14741506
}
14751507

1476-
fn build_shared_internal_secret(trino: &v1alpha1::TrinoCluster) -> Result<Secret> {
1477-
let mut internal_secret = BTreeMap::new();
1478-
internal_secret.insert(ENV_INTERNAL_SECRET.to_string(), get_random_base64());
1479-
1480-
Ok(Secret {
1481-
immutable: Some(true),
1482-
metadata: ObjectMetaBuilder::new()
1483-
.name(build_shared_internal_secret_name(trino))
1484-
.namespace_opt(trino.namespace())
1485-
.ownerreference_from_resource(trino, None, Some(true))
1486-
.context(ObjectMissingMetadataForOwnerRefSnafu)?
1487-
.build(),
1488-
string_data: Some(internal_secret),
1489-
..Secret::default()
1490-
})
1508+
fn shared_internal_secret_name(trino: &v1alpha1::TrinoCluster) -> String {
1509+
format!("{}-internal-secret", trino.name_any())
14911510
}
14921511

1493-
fn build_shared_internal_secret_name(trino: &v1alpha1::TrinoCluster) -> String {
1494-
format!("{}-internal-secret", trino.name_any())
1512+
fn shared_spooling_secret_name(trino: &v1alpha1::TrinoCluster) -> String {
1513+
format!("{}-spooling-secret", trino.name_any())
14951514
}
14961515

14971516
fn get_random_base64() -> String {
@@ -1644,6 +1663,7 @@ fn tls_volume_mounts(
16441663
catalogs: &[CatalogConfig],
16451664
requested_secret_lifetime: &Duration,
16461665
resolved_fte_config: &Option<ResolvedFaultTolerantExecutionConfig>,
1666+
resolved_spooling_config: &Option<client_protocol::ResolvedSpoolingProtocolConfig>,
16471667
) -> Result<()> {
16481668
if let Some(server_tls) = trino.get_server_tls() {
16491669
cb_prepare
@@ -1736,6 +1756,19 @@ fn tls_volume_mounts(
17361756
.context(AddVolumeSnafu)?;
17371757
}
17381758

1759+
// client spooling S3 credentials and other resources
1760+
if let Some(resolved_spooling) = resolved_spooling_config {
1761+
cb_prepare
1762+
.add_volume_mounts(resolved_spooling.volume_mounts.clone())
1763+
.context(AddVolumeMountSnafu)?;
1764+
cb_trino
1765+
.add_volume_mounts(resolved_spooling.volume_mounts.clone())
1766+
.context(AddVolumeMountSnafu)?;
1767+
pod_builder
1768+
.add_volumes(resolved_spooling.volumes.clone())
1769+
.context(AddVolumeSnafu)?;
1770+
}
1771+
17391772
Ok(())
17401773
}
17411774

0 commit comments

Comments
 (0)