Skip to content

Commit e916710

Browse files
committed
refactor crd to use spec.clusterConfig.clientProtocol.spooling
1 parent 78f2e64 commit e916710

7 files changed

Lines changed: 216 additions & 201 deletions

File tree

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

Lines changed: 144 additions & 138 deletions
Large diffs are not rendered by default.

docs/modules/trino/pages/usage-guide/client-spooling-protocol.adoc

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -16,12 +16,13 @@ To enable it, you need to set the `spec.clusterConfig.clientSpoolingProtocol` co
1616
----
1717
spec:
1818
clusterConfig:
19-
clientSpoolingProtocol:
20-
location: "s3://spooling-bucket/trino/" # <1>
21-
filesystem:
22-
s3: # <2>
23-
connection:
24-
reference: "minio"
19+
clientProtocol:
20+
spooling:
21+
location: "s3://spooling-bucket/trino/" # <1>
22+
filesystem:
23+
s3: # <2>
24+
connection:
25+
reference: "minio"
2526
----
2627
<1> Specifies the location where spooled data will be stored. This example uses an S3 bucket.
2728
<2> Configures the filesystem type for spooling. Only S3 is supported currently via the custom resource definition.

rust/operator-binary/src/command.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -24,7 +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>,
27+
resolved_spooling_config: &Option<client_protocol::ResolvedClientProtocolConfig>,
2828
) -> Vec<String> {
2929
let mut args = vec![];
3030

@@ -98,7 +98,7 @@ pub fn container_trino_args(
9898
authentication_config: &TrinoAuthenticationConfig,
9999
catalogs: &[CatalogConfig],
100100
resolved_fte_config: &Option<ResolvedFaultTolerantExecutionConfig>,
101-
resolved_spooling_config: &Option<client_protocol::ResolvedSpoolingProtocolConfig>,
101+
resolved_spooling_config: &Option<client_protocol::ResolvedClientProtocolConfig>,
102102
) -> Vec<String> {
103103
let mut args = vec![
104104
// copy config files to a writeable empty folder

rust/operator-binary/src/controller.rs

Lines changed: 7 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -464,10 +464,9 @@ pub async fn reconcile_trino(
464464
};
465465

466466
// Resolve client spooling protocol configuration with S3 connections if needed
467-
let resolved_spooling_config = match trino.spec.cluster_config.client_spooling_protocol.as_ref()
468-
{
467+
let resolved_client_protocol_config = match trino.spec.cluster_config.client_protocol.as_ref() {
469468
Some(spooling_config) => Some(
470-
client_protocol::ResolvedSpoolingProtocolConfig::from_config(
469+
client_protocol::ResolvedClientProtocolConfig::from_config(
471470
spooling_config,
472471
Some(client),
473472
&namespace,
@@ -599,7 +598,7 @@ pub async fn reconcile_trino(
599598
&trino_opa_config,
600599
&client.kubernetes_cluster_info,
601600
&resolved_fte_config,
602-
&resolved_spooling_config,
601+
&resolved_client_protocol_config,
603602
)?;
604603
let rg_catalog_configmap = build_rolegroup_catalog_config_map(
605604
trino,
@@ -618,7 +617,7 @@ pub async fn reconcile_trino(
618617
&catalogs,
619618
&rbac_sa.name_any(),
620619
&resolved_fte_config,
621-
&resolved_spooling_config,
620+
&resolved_client_protocol_config,
622621
)?;
623622

624623
cluster_resources
@@ -728,7 +727,7 @@ fn build_rolegroup_config_map(
728727
trino_opa_config: &Option<TrinoOpaConfig>,
729728
cluster_info: &KubernetesClusterInfo,
730729
resolved_fte_config: &Option<ResolvedFaultTolerantExecutionConfig>,
731-
resolved_spooling_config: &Option<client_protocol::ResolvedSpoolingProtocolConfig>,
730+
resolved_spooling_config: &Option<client_protocol::ResolvedClientProtocolConfig>,
732731
) -> Result<ConfigMap> {
733732
let mut cm_conf_data = BTreeMap::new();
734733

@@ -1013,7 +1012,7 @@ fn build_rolegroup_statefulset(
10131012
catalogs: &[CatalogConfig],
10141013
sa_name: &str,
10151014
resolved_fte_config: &Option<ResolvedFaultTolerantExecutionConfig>,
1016-
resolved_spooling_config: &Option<client_protocol::ResolvedSpoolingProtocolConfig>,
1015+
resolved_spooling_config: &Option<client_protocol::ResolvedClientProtocolConfig>,
10171016
) -> Result<StatefulSet> {
10181017
let role = trino
10191018
.role(trino_role)
@@ -1689,7 +1688,7 @@ fn tls_volume_mounts(
16891688
catalogs: &[CatalogConfig],
16901689
requested_secret_lifetime: &Duration,
16911690
resolved_fte_config: &Option<ResolvedFaultTolerantExecutionConfig>,
1692-
resolved_spooling_config: &Option<client_protocol::ResolvedSpoolingProtocolConfig>,
1691+
resolved_spooling_config: &Option<client_protocol::ResolvedClientProtocolConfig>,
16931692
) -> Result<()> {
16941693
if let Some(server_tls) = trino.get_server_tls() {
16951694
cb_prepare

rust/operator-binary/src/crd/client_protocol.rs

Lines changed: 48 additions & 39 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,11 @@ use crate::{
2323
const SPOOLING_S3_AWS_ACCESS_KEY: &str = "SPOOLING_S3_AWS_ACCESS_KEY";
2424
const SPOOLING_S3_AWS_SECRET_KEY: &str = "SPOOLING_S3_AWS_SECRET_KEY";
2525

26+
#[derive(Clone, Debug, Deserialize, JsonSchema, PartialEq, Serialize)]
27+
#[serde(rename_all = "camelCase")]
28+
pub enum ClientProtocolConfig {
29+
Spooling(ClientSpoolingProtocolConfig),
30+
}
2631
#[derive(Clone, Debug, Deserialize, JsonSchema, PartialEq, Serialize)]
2732
#[serde(rename_all = "camelCase")]
2833
pub struct ClientSpoolingProtocolConfig {
@@ -70,7 +75,7 @@ pub struct S3SpoolingConfig {
7075
pub upload_part_size: Option<Quantity>,
7176
}
7277

73-
pub struct ResolvedSpoolingProtocolConfig {
78+
pub struct ResolvedClientProtocolConfig {
7479
/// Properties to add to config.properties
7580
pub config_properties: BTreeMap<String, String>,
7681

@@ -92,11 +97,11 @@ pub struct ResolvedSpoolingProtocolConfig {
9297
pub init_container_extra_start_commands: Vec<String>,
9398
}
9499

95-
impl ResolvedSpoolingProtocolConfig {
100+
impl ResolvedClientProtocolConfig {
96101
/// Resolve S3 connection properties from Kubernetes resources
97102
/// and prepare spooling filesystem configuration.
98103
pub async fn from_config(
99-
config: &ClientSpoolingProtocolConfig,
104+
config: &ClientProtocolConfig,
100105
client: Option<&Client>,
101106
namespace: &str,
102107
) -> Result<Self, Error> {
@@ -109,39 +114,43 @@ impl ResolvedSpoolingProtocolConfig {
109114
init_container_extra_start_commands: Vec::new(),
110115
};
111116

112-
// Resolve external resources if Kubernetes client is available
113-
// This should always be the case, except for when this function is called during unit tests
114-
if let Some(client) = client {
115-
match &config.filesystem {
116-
SpoolingFileSystemConfig::S3(s3_config) => {
117-
resolved_config
118-
.resolve_s3_backend(s3_config, client, namespace)
119-
.await?;
117+
match config {
118+
ClientProtocolConfig::Spooling(spooling_config) => {
119+
// Resolve external resources if Kubernetes client is available
120+
// This should always be the case, except for when this function is called during unit tests
121+
if let Some(client) = client {
122+
match &spooling_config.filesystem {
123+
SpoolingFileSystemConfig::S3(s3_config) => {
124+
resolved_config
125+
.resolve_s3_backend(s3_config, client, namespace)
126+
.await?;
127+
}
128+
}
120129
}
121-
}
122-
}
123130

124-
resolved_config.spooling_manager_properties.extend([
125-
("fs.location".to_string(), config.location.clone()),
126-
(
127-
"spooling-manager.name".to_string(),
128-
"filesystem".to_string(),
129-
),
130-
]);
131-
132-
// Enable spooling protocol
133-
resolved_config.config_properties.extend([
134-
("protocol.spooling.enabled".to_string(), "true".to_string()),
135-
(
136-
"protocol.spooling.shared-secret-key".to_string(),
137-
format!("${{ENV:{secret}}}", secret = ENV_SPOOLING_SECRET),
138-
),
139-
]);
131+
resolved_config.spooling_manager_properties.extend([
132+
("fs.location".to_string(), spooling_config.location.clone()),
133+
(
134+
"spooling-manager.name".to_string(),
135+
"filesystem".to_string(),
136+
),
137+
]);
138+
139+
// Enable spooling protocol
140+
resolved_config.config_properties.extend([
141+
("protocol.spooling.enabled".to_string(), "true".to_string()),
142+
(
143+
"protocol.spooling.shared-secret-key".to_string(),
144+
format!("${{ENV:{secret}}}", secret = ENV_SPOOLING_SECRET),
145+
),
146+
]);
140147

141-
// Finally, extend the spooling manager properties with any user configuration
142-
resolved_config
143-
.spooling_manager_properties
144-
.extend(config.config_overrides.clone());
148+
// Finally, extend the spooling manager properties with any user configuration
149+
resolved_config
150+
.spooling_manager_properties
151+
.extend(spooling_config.config_overrides.clone());
152+
}
153+
}
145154

146155
Ok(resolved_config)
147156
}
@@ -252,7 +261,7 @@ mod tests {
252261

253262
#[tokio::test]
254263
async fn test_spooling_config() {
255-
let config = ClientSpoolingProtocolConfig {
264+
let config = ClientProtocolConfig::Spooling(ClientSpoolingProtocolConfig {
256265
location: "s3://my-bucket/spooling".to_string(),
257266
filesystem: SpoolingFileSystemConfig::S3(S3SpoolingConfig {
258267
connection:
@@ -265,9 +274,9 @@ mod tests {
265274
upload_part_size: None,
266275
}),
267276
config_overrides: HashMap::new(),
268-
};
277+
});
269278

270-
let resolved_spooling_config = ResolvedSpoolingProtocolConfig::from_config(
279+
let resolved_spooling_config = ResolvedClientProtocolConfig::from_config(
271280
&config, None, // No client, so no external resolution
272281
"default",
273282
)
@@ -292,7 +301,7 @@ mod tests {
292301

293302
#[tokio::test]
294303
async fn test_spooling_config_overrides() {
295-
let config = ClientSpoolingProtocolConfig {
304+
let config = ClientProtocolConfig::Spooling(ClientSpoolingProtocolConfig {
296305
location: "s3://my-bucket/spooling".to_string(),
297306
filesystem: SpoolingFileSystemConfig::S3(S3SpoolingConfig {
298307
connection:
@@ -308,9 +317,9 @@ mod tests {
308317
"protocol.spooling.retrieval-mode".to_string(),
309318
"STORAGE".to_string(),
310319
)]),
311-
};
320+
});
312321

313-
let resolved_spooling_config = ResolvedSpoolingProtocolConfig::from_config(
322+
let resolved_spooling_config = ResolvedClientProtocolConfig::from_config(
314323
&config, None, // No client, so no external resolution
315324
"default",
316325
)

rust/operator-binary/src/crd/mod.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -305,7 +305,7 @@ pub mod versioned {
305305

306306
/// Client spooling protocol configuration.
307307
#[serde(skip_serializing_if = "Option::is_none")]
308-
pub client_spooling_protocol: Option<client_protocol::ClientSpoolingProtocolConfig>,
308+
pub client_protocol: Option<client_protocol::ClientProtocolConfig>,
309309

310310
/// Name of the Vector aggregator [discovery ConfigMap](DOCS_BASE_URL_PLACEHOLDER/concepts/service_discovery).
311311
/// It must contain the key `ADDRESS` with the address of the Vector aggregator.

tests/templates/kuttl/client-spooling/02-install-trino.yaml.j2

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -16,13 +16,13 @@ spec:
1616
catalogLabelSelector:
1717
matchLabels:
1818
trino: trino
19-
clientSpoolingProtocol:
20-
enabled: true
21-
location: "s3://spooling-bucket/trino/"
22-
filesystem:
23-
s3:
24-
connection:
25-
reference: "minio"
19+
clientProtocol:
20+
spooling:
21+
location: "s3://spooling-bucket/trino/"
22+
filesystem:
23+
s3:
24+
connection:
25+
reference: "minio"
2626
{% if lookup('env', 'VECTOR_AGGREGATOR') %}
2727
vectorAggregatorConfigMapName: vector-aggregator-discovery
2828
{% endif %}

0 commit comments

Comments
 (0)