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
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Changed

- A listener with SASL configured closes a connection that sends a frame larger than 512KiB before the client authenticates, matching the Apache Kafka default for `sasl.server.max.receive.size`. The same limit applies while a client re-authenticates. The broker logs this rejection as `PreAuthenticationFrameTooBig`, and counts it in `nisshi_frames_rejected`.
- The broker stops at startup when its storage URL has a query option that the storage engine does not read, such as a misspelt option or `vacuum_into` on a `postgres://` URL. The error names the option and the engine. Previously the broker ignored the option.
Comment thread
solace-aross marked this conversation as resolved.
- The broker stops at startup when `maintenance_interval` or `transaction_maintenance_interval` has an invalid value: unparsable, zero, longer than 365 days, or a bare number without a unit. Previously the broker ignored an invalid value and used the default interval. Give a bare number its unit, for example `600s` or `10m` instead of `600`. Compound values such as `1h30m` and `5min` still work.
- A `parquet`, `iceberg` or `delta` broker without `--schema-registry` stops at startup with an error that names the missing option, instead of panicking. `--schema-registry` is accepted before or after the subcommand.
- When the broker closes the connection of a client that sends a request other than ApiVersions, SaslHandshake or SaslAuthenticate before it authenticates, it logs an ERROR line that names the client's address.

### Security

Expand Down
1 change: 1 addition & 0 deletions Cargo.lock

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

1 change: 1 addition & 0 deletions nisshi-broker/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ console.workspace = true
futures.workspace = true
glob.workspace = true
http-body-util.workspace = true
human-units.workspace = true
humantime.workspace = true
hyper-util.workspace = true
hyper.workspace = true
Expand Down
80 changes: 21 additions & 59 deletions nisshi-broker/src/broker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
// limitations under the License.

pub mod group;
mod maintenance;

use crate::{
CancelKind, Error, Result,
Expand Down Expand Up @@ -455,6 +456,18 @@ where
// not an anomaly worth an `error!` on every occurrence.
|| io.kind() == ErrorKind::TimedOut => {}

// The broker closes the connection of a client that sends a
// request other than ApiVersions or SASL before it
// authenticates. This line names the peer at ERROR, the level
// of the default filter, because that filter disables the
// INFO `peer` span that carries the address.
Err(Error::KafkaProtocol(nisshi_sans_io::Error::NotAuthenticated)) => {
error!(
%addr,
"closed connection: client sent a request before it authenticated"
);
}

Err(error) => {
error!(?error);
},
Expand Down Expand Up @@ -556,8 +569,6 @@ pub struct Builder<N, C, I, A, S, L> {
authentication: bool,
tls_server_config: Option<ServerConfig>,
silent: bool,
maintenance_interval: Option<Duration>,
transaction_maintenance_interval: Option<Duration>,

cancellation: CancellationToken,
}
Expand All @@ -572,9 +583,6 @@ type PhantomBuilder = Builder<
>;

impl<N, C, I, A, S, L> Builder<N, C, I, A, S, L> {
const MAINTENANCE_INTERVAL: &str = "maintenance_interval";
const TRANSACTION_MAINTENANCE_INTERVAL: &str = "transaction_maintenance_interval";

pub fn node_id(self, node_id: i32) -> Builder<i32, C, I, A, S, L> {
Builder {
node_id,
Expand All @@ -589,8 +597,6 @@ impl<N, C, I, A, S, L> Builder<N, C, I, A, S, L> {
authentication: self.authentication,
tls_server_config: self.tls_server_config,
silent: self.silent,
maintenance_interval: self.maintenance_interval,
transaction_maintenance_interval: self.transaction_maintenance_interval,
cancellation: self.cancellation,
}
}
Expand All @@ -609,8 +615,6 @@ impl<N, C, I, A, S, L> Builder<N, C, I, A, S, L> {
authentication: self.authentication,
tls_server_config: self.tls_server_config,
silent: self.silent,
maintenance_interval: self.maintenance_interval,
transaction_maintenance_interval: self.transaction_maintenance_interval,

cancellation: self.cancellation,
}
Expand All @@ -630,8 +634,6 @@ impl<N, C, I, A, S, L> Builder<N, C, I, A, S, L> {
authentication: self.authentication,
tls_server_config: self.tls_server_config,
silent: self.silent,
maintenance_interval: self.maintenance_interval,
transaction_maintenance_interval: self.transaction_maintenance_interval,

cancellation: self.cancellation,
}
Expand All @@ -654,52 +656,13 @@ impl<N, C, I, A, S, L> Builder<N, C, I, A, S, L> {
authentication: self.authentication,
tls_server_config: self.tls_server_config,
silent: self.silent,
maintenance_interval: self.maintenance_interval,
transaction_maintenance_interval: self.transaction_maintenance_interval,

cancellation: self.cancellation,
}
}

pub fn storage(self, mut storage: Url) -> Builder<N, C, I, A, Url, L> {
let maintenance_interval = storage.query_pairs().find_map(|(k, v)| {
if k == Self::MAINTENANCE_INTERVAL {
v.parse::<humantime::Duration>().map(Into::into).ok()
} else {
None
}
});

let transaction_maintenance_interval = storage.query_pairs().find_map(|(k, v)| {
if k == Self::TRANSACTION_MAINTENANCE_INTERVAL {
v.parse::<humantime::Duration>().map(Into::into).ok()
} else {
None
}
});

let pairs = storage
.query_pairs()
.filter_map(|(k, v)| {
if k == Self::MAINTENANCE_INTERVAL || k == Self::TRANSACTION_MAINTENANCE_INTERVAL {
None
} else {
Some((k.to_string(), v.to_string()))
}
})
.collect::<Vec<_>>();

if pairs.is_empty() {
storage.set_query(None);
} else {
_ = storage.query_pairs_mut().clear().extend_pairs(pairs);
}

debug!(
?maintenance_interval,
?transaction_maintenance_interval,
storage = %redact_url(&storage)
);
pub fn storage(self, storage: Url) -> Builder<N, C, I, A, Url, L> {
debug!(storage = %redact_url(&storage));

Builder {
node_id: self.node_id,
Expand All @@ -714,8 +677,6 @@ impl<N, C, I, A, S, L> Builder<N, C, I, A, S, L> {
authentication: self.authentication,
tls_server_config: self.tls_server_config,
silent: self.silent,
maintenance_interval,
transaction_maintenance_interval,

cancellation: self.cancellation,
}
Expand All @@ -737,8 +698,6 @@ impl<N, C, I, A, S, L> Builder<N, C, I, A, S, L> {
authentication: self.authentication,
tls_server_config: self.tls_server_config,
silent: self.silent,
maintenance_interval: self.maintenance_interval,
transaction_maintenance_interval: self.transaction_maintenance_interval,

cancellation: self.cancellation,
}
Expand Down Expand Up @@ -793,6 +752,9 @@ impl Builder<i32, String, Uuid, Url, Url, Url> {
.map(otel::metric_exporter)
.transpose()?;

let (storage, intervals) = maintenance::take_intervals(self.storage.clone())?;
debug!(?intervals);

let builder = {
let mut builder = StorageContainer::builder();

Expand Down Expand Up @@ -832,7 +794,7 @@ impl Builder<i32, String, Uuid, Url, Url, Url> {
.advertised_listener(self.advertised_listener.clone())
.schema_registry(self.schema_registry.clone())
.lake_house(self.lake_house.clone())
.storage(self.storage.clone())
.storage(storage)
.cancellation(self.cancellation.clone())
.silent(self.silent)
.build()
Expand All @@ -859,8 +821,8 @@ impl Builder<i32, String, Uuid, Url, Url, Url> {
tls_server_config: self.tls_server_config.map(Arc::new),

silent: self.silent,
maintenance_interval: self.maintenance_interval,
transaction_maintenance_interval: self.transaction_maintenance_interval,
maintenance_interval: intervals.maintenance,
Comment thread
solace-aross marked this conversation as resolved.
transaction_maintenance_interval: intervals.transaction_maintenance,
cancellation: self.cancellation,
meter_provider,
})
Expand Down
Loading
Loading