Skip to content

Commit f0e6cbb

Browse files
committed
refactor: switch to v2 JavaCommonConfig & remove framework module
1 parent c3a263e commit f0e6cbb

14 files changed

Lines changed: 482 additions & 750 deletions

File tree

Cargo.lock

Lines changed: 127 additions & 139 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

Cargo.nix

Lines changed: 167 additions & 223 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.
Lines changed: 42 additions & 73 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,7 @@
1-
use serde::Serialize;
21
use snafu::{OptionExt, ResultExt, Snafu};
32
use stackable_operator::{
43
memory::{BinaryMultiple, MemoryQuantity},
5-
role_utils::{self, GenericRoleConfig, JavaCommonConfig, JvmArgumentOverrides, Role},
6-
schemars::JsonSchema,
4+
v2::jvm_argument_overrides::JvmArgumentOverrides,
75
};
86

97
use crate::crd::{ConfigFileName, METRICS_PORT, STACKABLE_CONFIG_DIR, role::AnyConfig};
@@ -19,20 +17,14 @@ pub enum Error {
1917
InvalidMemoryConfig {
2018
source: stackable_operator::memory::Error,
2119
},
22-
23-
#[snafu(display("failed to merge jvm argument overrides"))]
24-
MergeJvmArgumentOverrides { source: role_utils::Error },
2520
}
2621

27-
/// All JVM arguments.
28-
fn construct_jvm_args<ConfigFragment, ConfigOverrides>(
22+
/// All JVM arguments: operator-generated base args with the already-merged
23+
/// (role <- role group) `jvmArgumentOverrides` applied on top.
24+
fn construct_jvm_args(
2925
merged_config: &AnyConfig,
30-
role: &Role<ConfigFragment, ConfigOverrides, GenericRoleConfig, JavaCommonConfig>,
31-
role_group: &str,
32-
) -> Result<Vec<String>, Error>
33-
where
34-
ConfigOverrides: Default + JsonSchema + Serialize,
35-
{
26+
jvm_argument_overrides: &JvmArgumentOverrides,
27+
) -> Result<Vec<String>, Error> {
3628
let heap_size = MemoryQuantity::try_from(
3729
merged_config
3830
.resources()
@@ -49,7 +41,6 @@ where
4941
.context(InvalidMemoryConfigSnafu)?;
5042

5143
let jvm_args = vec![
52-
// Heap settings
5344
format!("-Xmx{java_heap}"),
5445
format!("-Xms{java_heap}"),
5546
format!(
@@ -61,43 +52,26 @@ where
6152
),
6253
];
6354

64-
let operator_generated = JvmArgumentOverrides::new_with_only_additions(jvm_args);
65-
let merged = role
66-
.get_merged_jvm_argument_overrides(role_group, &operator_generated)
67-
.context(MergeJvmArgumentOverridesSnafu)?;
68-
Ok(merged
69-
.effective_jvm_config_after_merging()
70-
// Sorry for the clone, that's how operator-rs is currently modelled :P
71-
.clone())
55+
Ok(jvm_argument_overrides.apply_to(jvm_args))
7256
}
7357

74-
/// Arguments that go into `EXTRA_ARGS`, so *not* the heap settings (which you can get using
75-
/// [`construct_heap_jvm_args`]).
76-
pub fn construct_non_heap_jvm_args<ConfigFragment, ConfigOverrides>(
58+
/// Arguments that go into `EXTRA_ARGS` (everything except heap settings).
59+
pub fn construct_non_heap_jvm_args(
7760
merged_config: &AnyConfig,
78-
role: &Role<ConfigFragment, ConfigOverrides, GenericRoleConfig, JavaCommonConfig>,
79-
role_group: &str,
80-
) -> Result<String, Error>
81-
where
82-
ConfigOverrides: Default + JsonSchema + Serialize,
83-
{
84-
let mut jvm_args = construct_jvm_args(merged_config, role, role_group)?;
61+
jvm_argument_overrides: &JvmArgumentOverrides,
62+
) -> Result<String, Error> {
63+
let mut jvm_args = construct_jvm_args(merged_config, jvm_argument_overrides)?;
8564
jvm_args.retain(|arg| !is_heap_jvm_argument(arg));
8665

8766
Ok(jvm_args.join(" "))
8867
}
8968

90-
/// Arguments that go into `KAFKA_HEAP_OPTS`.
91-
/// You can get the normal JVM arguments using [`construct_non_heap_jvm_args`].
92-
pub fn construct_heap_jvm_args<ConfigFragment, ConfigOverrides>(
69+
/// Arguments that go into `KAFKA_HEAP_OPTS` (only the heap settings).
70+
pub fn construct_heap_jvm_args(
9371
merged_config: &AnyConfig,
94-
role: &Role<ConfigFragment, ConfigOverrides, GenericRoleConfig, JavaCommonConfig>,
95-
role_group: &str,
96-
) -> Result<String, Error>
97-
where
98-
ConfigOverrides: Default + JsonSchema + Serialize,
99-
{
100-
let mut jvm_args = construct_jvm_args(merged_config, role, role_group)?;
72+
jvm_argument_overrides: &JvmArgumentOverrides,
73+
) -> Result<String, Error> {
74+
let mut jvm_args = construct_jvm_args(merged_config, jvm_argument_overrides)?;
10175
jvm_args.retain(|arg| is_heap_jvm_argument(arg));
10276

10377
Ok(jvm_args.join(" "))
@@ -111,25 +85,34 @@ fn is_heap_jvm_argument(jvm_argument: &str) -> bool {
11185

11286
#[cfg(test)]
11387
mod tests {
114-
use stackable_operator::kube::ResourceExt;
115-
11688
use super::*;
11789
use crate::{
118-
crd::{
119-
BrokerRole,
120-
role::{KafkaRole, broker::BrokerConfig},
121-
v1alpha1,
122-
},
123-
framework::role_utils::with_validated_config,
90+
controller::test_support::{minimal_kafka, validated_cluster},
91+
crd::role::KafkaRole,
12492
};
12593

94+
/// Pulls the broker `default` role group's merged config + merged JVM overrides out of
95+
/// a validated cluster built from the given YAML.
96+
fn broker_default(yaml: &str) -> (AnyConfig, JvmArgumentOverrides) {
97+
let kafka = minimal_kafka(yaml);
98+
let validated = validated_cluster(&kafka);
99+
let rg = validated
100+
.role_group_configs
101+
.get(&KafkaRole::Broker)
102+
.and_then(|groups| groups.get("default"))
103+
.expect("broker default role group should exist");
104+
(rg.config.clone(), rg.jvm_argument_overrides.clone())
105+
}
106+
126107
#[test]
127108
fn test_construct_jvm_arguments_defaults() {
128109
let input = r#"
129110
apiVersion: kafka.stackable.tech/v1alpha1
130111
kind: KafkaCluster
131112
metadata:
132113
name: simple-kafka
114+
namespace: default
115+
uid: 12345678-1234-1234-1234-123456789012
133116
spec:
134117
image:
135118
productVersion: 3.9.2
@@ -140,10 +123,9 @@ mod tests {
140123
default:
141124
replicas: 1
142125
"#;
143-
let (kafka_role, role, merged_config) = construct_boilerplate(input);
144-
let non_heap_jvm_args =
145-
construct_non_heap_jvm_args(&kafka_role, &role, &merged_config).unwrap();
146-
let heap_jvm_args = construct_heap_jvm_args(&kafka_role, &role, &merged_config).unwrap();
126+
let (merged_config, jvm) = broker_default(input);
127+
let non_heap_jvm_args = construct_non_heap_jvm_args(&merged_config, &jvm).unwrap();
128+
let heap_jvm_args = construct_heap_jvm_args(&merged_config, &jvm).unwrap();
147129

148130
assert_eq!(
149131
non_heap_jvm_args,
@@ -160,6 +142,8 @@ mod tests {
160142
kind: KafkaCluster
161143
metadata:
162144
name: simple-kafka
145+
namespace: default
146+
uid: 12345678-1234-1234-1234-123456789012
163147
spec:
164148
image:
165149
productVersion: 3.9.2
@@ -187,10 +171,9 @@ mod tests {
187171
- -Xmx40000m
188172
- -Dhttps.proxyPort=1234
189173
"#;
190-
let (merged_config, role, role_group) = construct_boilerplate(input);
191-
let non_heap_jvm_args =
192-
construct_non_heap_jvm_args(&merged_config, &role, &role_group).unwrap();
193-
let heap_jvm_args = construct_heap_jvm_args(&merged_config, &role, &role_group).unwrap();
174+
let (merged_config, jvm) = broker_default(input);
175+
let non_heap_jvm_args = construct_non_heap_jvm_args(&merged_config, &jvm).unwrap();
176+
let heap_jvm_args = construct_heap_jvm_args(&merged_config, &jvm).unwrap();
194177

195178
assert_eq!(
196179
non_heap_jvm_args,
@@ -202,18 +185,4 @@ mod tests {
202185
);
203186
assert_eq!(heap_jvm_args, "-Xms34406m -Xmx40000m");
204187
}
205-
206-
fn construct_boilerplate(kafka_cluster: &str) -> (AnyConfig, BrokerRole, String) {
207-
let kafka: v1alpha1::KafkaCluster =
208-
serde_yaml::from_str(kafka_cluster).expect("illegal test input");
209-
210-
let role = kafka.spec.brokers.clone().unwrap();
211-
let role_group = role.role_groups.get("default").unwrap();
212-
let default_config =
213-
BrokerConfig::default_config(&kafka.name_any(), &KafkaRole::Broker.to_string());
214-
let validated = with_validated_config(role_group, &role, &default_config).unwrap();
215-
let merged_config = AnyConfig::Broker(validated.config);
216-
217-
(merged_config, role, "default".to_owned())
218-
}
219188
}

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

Lines changed: 65 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -20,15 +20,12 @@ pub(crate) mod build;
2020
pub(crate) mod dereference;
2121
pub(crate) mod validate;
2222

23-
use crate::{
24-
crd::{
25-
KafkaPodDescriptor, MetadataManager,
26-
authorization::KafkaAuthorizationConfig,
27-
role::{AnyConfig, AnyConfigOverrides, KafkaRole},
28-
security::KafkaTlsSecurity,
29-
v1alpha1,
30-
},
31-
framework::role_utils::RoleGroupConfig,
23+
use crate::crd::{
24+
KafkaPodDescriptor, MetadataManager,
25+
authorization::KafkaAuthorizationConfig,
26+
role::{AnyConfig, AnyConfigOverrides, KafkaRole},
27+
security::KafkaTlsSecurity,
28+
v1alpha1,
3229
};
3330

3431
pub type RoleGroupName = String;
@@ -140,12 +137,62 @@ impl Resource for ValidatedCluster {
140137
/// A validated, merged Kafka role-group config.
141138
///
142139
/// The merged config fragment is wrapped in [`AnyConfig`] and the merged
143-
/// `configOverrides` in [`AnyConfigOverrides`], so a single role-agnostic type
144-
/// carries both broker and controller role groups (their concrete config and
145-
/// override types differ). Produced via the local-`framework`
146-
/// [`with_validated_config`](crate::framework::role_utils::with_validated_config).
147-
pub type ValidatedRoleGroupConfig = RoleGroupConfig<
148-
AnyConfig,
149-
stackable_operator::role_utils::JavaCommonConfig,
150-
AnyConfigOverrides,
151-
>;
140+
/// `configOverrides` in [`AnyConfigOverrides`], so a single role-agnostic type carries
141+
/// both broker and controller role groups. Produced from the upstream
142+
/// [`stackable_operator::v2::role_utils::with_validated_config`] result in
143+
/// [`validate`](crate::controller::validate). `jvm_argument_overrides` is already merged
144+
/// (role <- role group) at validation time and applied as-is during build.
145+
#[derive(Clone, Debug, PartialEq)]
146+
pub struct ValidatedRoleGroupConfig {
147+
pub replicas: u16,
148+
pub config: AnyConfig,
149+
pub config_overrides: AnyConfigOverrides,
150+
pub env_overrides: stackable_operator::v2::builder::pod::container::EnvVarSet,
151+
pub pod_overrides: stackable_operator::k8s_openapi::api::core::v1::PodTemplateSpec,
152+
pub jvm_argument_overrides:
153+
stackable_operator::v2::jvm_argument_overrides::JvmArgumentOverrides,
154+
}
155+
156+
#[cfg(test)]
157+
pub(crate) mod test_support {
158+
use stackable_operator::{
159+
cli::OperatorEnvironmentOptions,
160+
commons::networking::DomainName,
161+
utils::{cluster_info::KubernetesClusterInfo, yaml_from_str_singleton_map},
162+
};
163+
164+
use super::{ValidatedCluster, dereference::DereferencedObjects, validate::validate};
165+
use crate::crd::{authentication::ResolvedAuthenticationClasses, v1alpha1};
166+
167+
pub fn minimal_kafka(yaml: &str) -> v1alpha1::KafkaCluster {
168+
yaml_from_str_singleton_map(yaml).expect("invalid test KafkaCluster YAML")
169+
}
170+
171+
fn cluster_info() -> KubernetesClusterInfo {
172+
KubernetesClusterInfo {
173+
cluster_domain: DomainName::try_from("cluster.local").expect("valid domain"),
174+
}
175+
}
176+
177+
fn operator_environment() -> OperatorEnvironmentOptions {
178+
OperatorEnvironmentOptions {
179+
operator_namespace: "stackable-operators".to_owned(),
180+
operator_service_name: "kafka-operator".to_owned(),
181+
image_repository: "oci.example.org".to_owned(),
182+
}
183+
}
184+
185+
/// Runs the real validate step against a minimal (auth/OPA-free) fixture.
186+
pub fn validated_cluster(kafka: &v1alpha1::KafkaCluster) -> ValidatedCluster {
187+
validate(
188+
kafka,
189+
DereferencedObjects {
190+
authentication_classes: ResolvedAuthenticationClasses::new(Vec::new()),
191+
authorization_config: None,
192+
kubernetes_cluster_info: cluster_info(),
193+
},
194+
&operator_environment(),
195+
)
196+
.expect("validate should succeed for the test fixture")
197+
}
198+
}

0 commit comments

Comments
 (0)