use std::collections::HashMap;
use crate::error::{KrafkaError, ProtocolErrorKind, Result};
use crate::protocol::{
ApiKey, ApiVersionsRequest, ApiVersionsResponse, FeatureUpdateKey, FinalizedFeature,
SupportedFeature, UpdateFeaturesRequest, UpdateFeaturesResponse, versions,
};
use super::AdminClient;
use super::driver::{Mode, Target, answer, exchange, negotiate};
#[non_exhaustive]
#[derive(Debug, Clone)]
pub struct FeatureMetadata {
pub supported: Vec<SupportedFeature>,
pub finalized: Vec<FinalizedFeature>,
pub finalized_epoch: Option<i64>,
}
admin_options! {
DescribeFeaturesOptions {}
optional {
node_id: i32,
}
}
admin_options! {
UpdateFeaturesOptions {
validate_only: bool,
}
}
impl AdminClient {
pub async fn describe_features(
&self,
options: DescribeFeaturesOptions,
) -> Result<FeatureMetadata> {
let call = self.call("DescribeFeatures", Mode::Read, options.timeout)?;
let target = options.node_id.map_or(Target::AnyBroker, Target::Broker);
call.single(target, |conn| async move {
let request =
ApiVersionsRequest::new().with_client_software("krafka", env!("CARGO_PKG_VERSION"));
let version = negotiate(&conn, ApiKey::ApiVersions, 3, versions::API_VERSIONS_MAX)?;
let mut response = conn
.send_request(ApiKey::ApiVersions, version, |buf| {
if version >= 5 {
request.encode_v5(buf)
} else {
request.encode_v3(buf)
}
})
.await?;
let response = ApiVersionsResponse::decode_v3(&mut response)?;
answer(crate::error::ErrorCode::from(response.error_code), None)?;
Ok(FeatureMetadata {
supported: response.supported_features,
finalized: response.finalized_features,
finalized_epoch: (response.finalized_features_epoch >= 0)
.then_some(response.finalized_features_epoch),
})
})
.await
}
pub async fn update_features(
&self,
updates: Vec<FeatureUpdateKey>,
options: UpdateFeaturesOptions,
) -> Result<HashMap<String, Result<()>>> {
let call = self.call("UpdateFeatures", Mode::Write, options.timeout)?;
let validate_only = options.validate_only;
let updates_ref = &updates;
let call_ref = &call;
call.single(Target::Controller, |conn| async move {
let version = negotiate(
&conn,
ApiKey::UpdateFeatures,
versions::UPDATE_FEATURES_MIN,
versions::UPDATE_FEATURES_MAX,
)?;
if validate_only && version < 1 {
return Err(KrafkaError::protocol_kind(
ProtocolErrorKind::UnknownApiVersion,
"validate_only needs UpdateFeatures v1; the controller supports v0 only",
));
}
let mut request =
UpdateFeaturesRequest::new(updates_ref.clone()).with_validate_only(validate_only);
request.timeout_ms = call_ref.remaining_ms();
let response: UpdateFeaturesResponse =
exchange(&conn, ApiKey::UpdateFeatures, version, &request).await?;
answer(response.error_code, response.error_message)?;
let mut per_feature: HashMap<String, Result<()>> = updates_ref
.iter()
.map(|u| (u.feature.clone(), Ok(())))
.collect();
for r in response.results {
per_feature.insert(r.feature, answer(r.error_code, r.error_message));
}
Ok(per_feature)
})
.await
}
}