hubcio commented on code in PR #4318:
URL: https://github.com/apache/iggy/pull/4318#discussion_r4138942574
##########
core/server/src/partition_helpers.rs:
##########
@@ -340,6 +341,15 @@ pub async fn load_partition_or_fence(
.await
.map(Some)
}
+ // No free file descriptor says nothing about the record, so the load
+ // fails like any other: the reconciler retries it with backoff, and a
Review Comment:
nit: at boot this `Err` stops the node through `boot/recovery.rs:240` with
exit 1, because the superblock read is not noted, and no reconciler retries it.
document the boot stop and note this read so the exit is 4.
##########
core/server_common/src/fatal.rs:
##########
@@ -0,0 +1,193 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Stopping the process on an environmental failure that has no in-process
answer.
+
+use nix::errno::Errno;
+use nix::sys::resource::{Resource, getrlimit};
+use std::io::{self, Write};
+use std::sync::atomic::{AtomicBool, Ordering};
+
+/// Why the process is stopping. The discriminant is the exit status, one per
+/// condition; `1` stays the binary's generic startup failure.
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+#[repr(u8)]
+pub enum FatalReason {
+ /// A prepare's WAL append failed and the op it had claimed could not be
handed
+ /// back. The durable log is intact up to the previous op, so recovery
re-derives
+ /// the frontier and restarting is the repair.
+ UnreconcilableLogFrontier = 2,
+ /// The superblock stayed unwritable past the configured fail-stop window.
+ /// The replica was already fenced quorum-invisible, so exiting hands the
+ /// wedge to a supervisor instead of a log reader.
+ SuperblockWedged = 3,
+ /// The server stopped on an error after a storage open for a write or a
+ /// sync failed because the process (`EMFILE`) or the host (`ENFILE`) had
no
+ /// free file descriptor. The exit comes after the ordinary shutdown, so
the
+ /// shards flushed what they could. See [`NoteDescriptorExhaustion`].
+ DescriptorsExhausted = 4,
+}
+
+impl FatalReason {
+ #[must_use]
+ pub const fn exit_status(self) -> u8 {
+ self as u8
+ }
+}
+
+/// The target that 0.9.0 shipped for the fatal log line, so existing log
+/// filters keep matching.
+const FATAL_LOG_TARGET: &str = "iggy.consensus.diag";
+
+/// Log `message` and terminate the process.
+///
+/// For an environmental failure where stopping IS the answer, rather than an
error
+/// threaded up a stack whose top knows less than this leaf does. Not for bugs
in
+/// this process, which are `assert!` / `panic!` and say so.
+///
+/// `exit`, not `panic!`: a panic unwinds one shard of a thread-per-core
runtime and
+/// leaves its siblings serving, which is the half-alive state this exists to
avoid.
+/// Skipping destructors is wanted here, since the reason for stopping is that
+/// further writes cannot be trusted.
+///
+/// The reason also goes straight to stderr. The tracing appenders are
+/// non-blocking workers, and `exit` stops them before they flush, so without
+/// this write a supervisor's journal never learns why the process stopped. The
+/// write result is ignored because `eprintln!` would panic on a broken stderr.
+pub fn fatal(reason: FatalReason, message: &str) -> ! {
+ tracing::error!(
+ target: FATAL_LOG_TARGET,
+ reason = ?reason,
+ exit_status = reason.exit_status(),
+ "{message}"
+ );
+ let _ = writeln!(
+ std::io::stderr().lock(),
+ "iggy fatal: reason={reason:?} exit_status={}: {message}",
+ reason.exit_status()
+ );
+ std::process::exit(i32::from(reason.exit_status()));
+}
+
+/// Storage results that record a missing file descriptor for the exit status.
+///
+/// Only for opens that write, create, truncate or sync, and for read-only
+/// opens that are one step of such a write: the existence probes of segment
+/// setup, or an open that exists only to sync. The error still goes to the
+/// caller, whose own handling decides what happens. A failed partition write
+/// fences the partition and stops the server through the shutdown flush, and a
+/// failed partition build retries. Stopping at the open instead would skip
that
+/// flush and lose committed messages that are not yet in a segment.
+///
+/// If the server then stops on an error, it exits with
+/// [`FatalReason::DescriptorsExhausted`] instead of 1. The record lasts for
the
+/// life of the process. A failed read fails only its own request, so reads do
+/// not record it, and accept loops do not either, see `message_bus::accept`.
+pub trait NoteDescriptorExhaustion: Sized {
+ /// Pass the result through. If it failed with `EMFILE` or `ENFILE`, record
+ /// that for [`descriptors_exhausted`], and log the first one with the
+ /// limits. `operation` names what was being opened, for that log line.
+ #[must_use]
+ fn note_descriptor_exhaustion(self, operation: impl FnOnce() -> String) ->
Self;
+}
+
+static DESCRIPTORS_EXHAUSTED: AtomicBool = AtomicBool::new(false);
+
+impl<T> NoteDescriptorExhaustion for io::Result<T> {
+ fn note_descriptor_exhaustion(self, operation: impl FnOnce() -> String) ->
Self {
+ if let Err(error) = &self
+ && is_descriptor_exhaustion(error)
+ && !DESCRIPTORS_EXHAUSTED.swap(true, Ordering::Relaxed)
Review Comment:
warning: nothing clears this flag, so after a recovered `EMFILE`, for
example a superblock persist that succeeds on retry, a later `ENOSPC` fence or
any other error that reaches `main` exits 4. carry the exhaustion in the stop
cause instead.
##########
core/metadata/src/stm/authz.rs:
##########
@@ -345,6 +345,38 @@ pub(crate) fn authorize(
}
}
+/// Whether `authorize` lets a `CreateTopic`, or a `CreatePartitions` when
+/// `topic_id` is set, reach its apply: the ids resolve and `user_id` holds the
+/// grant. The primary asks before it denies a create for `[metadata]
+/// partitions_max`, so apply answers `NotFound` or `Unauthorized` first and a
+/// user who cannot create learns nothing about the cap.
+pub(crate) fn admits_partitions_create(
+ users: &Users,
+ streams: &Streams,
+ user_id: u32,
+ stream_id: &WireIdentifier,
+ topic_id: Option<&WireIdentifier>,
+) -> bool {
+ let Some(topic_id) = topic_id else {
+ return streams
Review Comment:
warning: when a `CreateTopic` for an existing name exceeds `partitions_max`,
clients get `PartitionsLimitReached`, not `TopicNameAlreadyExists`, and Kafka
clients get POLICY_VIOLATION. if the stream holds the name, skip the cap.
##########
core/partitions/src/state_transfer.rs:
##########
@@ -3054,6 +3076,11 @@ where
source,
})?,
};
+ // The open fsyncs the installed files, and then a sealed segment
+ // gives its writer descriptors back.
+ if index + 1 < staged.len() {
Review Comment:
nit: no test asserts that install closes sealed writers, so a leak here
stays silent, and every unit install test stages at most one segment. add a
multi-segment install test that asserts every storage but the last has no
`index_writer`.
##########
core/server/src/sysinfo_printer.rs:
##########
@@ -0,0 +1,202 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Periodic one-line log of process and host usage.
+//!
+//! Spawned on shard 0 only: the numbers describe the whole process, so one
+//! line per node is enough, and the client count is already a cross-shard
+//! gather.
+
+use crate::responses::{StatsTotals, stats_totals};
+use crate::shell::ServerShard;
+use crate::sysinfo_probe::{SystemStats, stats_disk_space};
+use iggy_common::IggyByteSize;
+use shard::Receiver;
+use std::fmt;
+use std::rc::Rc;
+use std::time::Duration;
+use sysinfo::System as SysinfoSystem;
+use tracing::level_filters::LevelFilter;
+use tracing::{error, info, trace};
+
+/// Run the printer until `stop` fires, logging one line every `interval`.
+pub async fn run_sysinfo_printer(shard: Rc<ServerShard>, stop: Receiver<()>,
interval: Duration) {
+ info!("System info logger is enabled, OS info will be printed every:
{interval:?}");
+ // Not the shared `GetStats` sampler: a tick would reset its CPU window,
and
+ // a `GetStats` on shard 0 right after a tick would show CPU near zero.
+ // Sampled once now, so the first line covers a full interval.
+ let mut system = SysinfoSystem::new();
+ SystemStats::capture(&mut system);
+ loop {
+ // `Ok(_)`: stop signalled -> exit. `Err(_)`: interval elapsed ->
print.
+ match compio::time::timeout(interval, stop.recv()).await {
+ Ok(_) => break,
+ Err(_) => print_sysinfo(&shard, &mut system).await,
+ }
+ }
+ trace!(shard = shard.id, "sysinfo printer exited");
+}
+
+async fn print_sysinfo(shard: &Rc<ServerShard>, system: &mut SysinfoSystem) {
+ // The global max level, not `tracing::enabled!`: the logger's idle
+ // OpenTelemetry layers veto every `enabled` query, so that is always
false.
+ if LevelFilter::current() < LevelFilter::INFO {
+ return;
+ }
+ let clients_count = shard.list_all_clients().await.len();
+ let totals = match stats_totals(shard) {
Review Comment:
simplification: the line reads only two `StatsTotals` fields, the message
sums, but each tick runs the full walk with one `StatsRegistry` lock per
partition. summing `stream.stats` here gives the same numbers and makes the
`stats_totals` extraction unnecessary.
##########
core/message_bus/src/connection_cap.rs:
##########
@@ -0,0 +1,161 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Node-wide cap on open client sockets.
+//!
+//! Shard 0 accepts every client socket, but the owning shard closes it. Each
+//! socket therefore carries a [`ConnectionPermit`] to its owning shard, and
+//! the permit frees its slot on drop.
+//!
+//! The count is a process-wide static, because one node runs per process. A
+//! permit that pointed at its count would add 16 bytes to the client setup
+//! frames, and every slot of every shard inbox pays for the largest frame.
+
+use std::cell::Cell;
+use std::sync::atomic::{AtomicUsize, Ordering};
+use std::time::{Duration, Instant};
+use tracing::warn;
+
+/// Shortest interval between two refusal warnings. A client that reconnects
+/// in a loop would otherwise log one line per attempt.
+const REFUSAL_LOG_INTERVAL: Duration = Duration::from_secs(10);
+
+static LIVE_CONNECTIONS: AtomicUsize = AtomicUsize::new(0);
+
+/// The open client sockets of the node, capped at `max`.
+///
+/// Owned by shard 0, where every client accept runs.
+#[derive(Debug)]
+pub struct ConnectionCap {
+ max: Option<usize>,
+ last_refusal_log: Cell<Option<Instant>>,
+ refused_since_log: Cell<u64>,
+}
+
+impl ConnectionCap {
+ /// A cap of `max` open sockets. `None` counts sockets and refuses none.
+ #[must_use]
+ pub const fn new(max: Option<usize>) -> Self {
Review Comment:
nit: `ConnectionCap::new(Some(0))` refuses every socket while
`connections_max = 0` means no cap, so zero means opposite things in the
configuration and here. take `Option<NonZeroUsize>` - `resolve_connections_max`
already maps zero to `None`.
##########
core/metadata/src/impls/metadata.rs:
##########
@@ -2220,6 +2237,66 @@ where
}
}
+ /// `[metadata] partitions_max` admission for a `CreateTopic` or
+ /// `CreatePartitions`. Other operations pass.
+ ///
+ /// A soft cap: the apply must not branch on node config, so this counts
+ /// the partitions this primary has committed. The creates in flight, up to
+ /// a full prepare queue and request queue of them, each pass it on their
+ /// own, so together they can overshoot it by up to 1000 partitions each.
+ ///
+ /// A body that does not decode passes here, and `prepare_request` evicts
+ /// the session for it. A create that the gated apply refuses, for a
+ /// missing target or a missing grant, also passes, so it gets that error
+ /// and not the cap's.
+ fn admit_partitions(&self, message: &Message<RoutedRequestHeader>) ->
Result<(), IggyError> {
Review Comment:
nit: a binary-protocol create that `admit_partitions` denies leaves no
server log line or metric, and the HTTP `error!` carries no counts. log a
warning with `partitions_max`, the committed count and the requested count.
##########
foreign/node/src/wire/system/get-stats.command.ts:
##########
@@ -91,6 +101,26 @@ const deserializeGetStats = (b: Buffer) => {
position + 4,
position + 4 + kernelVersionLength
).toString();
+ position += 4 + kernelVersionLength;
+
+ // iggy_server_version, iggy_server_semver
+ const iggyServerVersionLength = b.readUInt32LE(position);
+ position += 4 + iggyServerVersionLength + 4;
+
+ const cacheMetricsCount = b.readUInt32LE(position);
+ position += 4 + cacheMetricsCount * CACHE_METRIC_SIZE +
THREADS_AND_DISK_SIZE;
Review Comment:
nit: the parser now walks past the server version and semver, cache metrics,
thread count and disk fields, but `Stats` returns none of them. read them at
the positions the parser already tracks and add them to `Stats`.
##########
core/server_common/src/fatal.rs:
##########
@@ -0,0 +1,193 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Stopping the process on an environmental failure that has no in-process
answer.
+
+use nix::errno::Errno;
+use nix::sys::resource::{Resource, getrlimit};
+use std::io::{self, Write};
+use std::sync::atomic::{AtomicBool, Ordering};
+
+/// Why the process is stopping. The discriminant is the exit status, one per
+/// condition; `1` stays the binary's generic startup failure.
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+#[repr(u8)]
+pub enum FatalReason {
+ /// A prepare's WAL append failed and the op it had claimed could not be
handed
+ /// back. The durable log is intact up to the previous op, so recovery
re-derives
+ /// the frontier and restarting is the repair.
+ UnreconcilableLogFrontier = 2,
+ /// The superblock stayed unwritable past the configured fail-stop window.
+ /// The replica was already fenced quorum-invisible, so exiting hands the
+ /// wedge to a supervisor instead of a log reader.
+ SuperblockWedged = 3,
+ /// The server stopped on an error after a storage open for a write or a
+ /// sync failed because the process (`EMFILE`) or the host (`ENFILE`) had
no
+ /// free file descriptor. The exit comes after the ordinary shutdown, so
the
+ /// shards flushed what they could. See [`NoteDescriptorExhaustion`].
+ DescriptorsExhausted = 4,
+}
+
+impl FatalReason {
+ #[must_use]
+ pub const fn exit_status(self) -> u8 {
+ self as u8
+ }
+}
+
+/// The target that 0.9.0 shipped for the fatal log line, so existing log
+/// filters keep matching.
+const FATAL_LOG_TARGET: &str = "iggy.consensus.diag";
+
+/// Log `message` and terminate the process.
+///
+/// For an environmental failure where stopping IS the answer, rather than an
error
+/// threaded up a stack whose top knows less than this leaf does. Not for bugs
in
+/// this process, which are `assert!` / `panic!` and say so.
+///
+/// `exit`, not `panic!`: a panic unwinds one shard of a thread-per-core
runtime and
+/// leaves its siblings serving, which is the half-alive state this exists to
avoid.
+/// Skipping destructors is wanted here, since the reason for stopping is that
+/// further writes cannot be trusted.
+///
+/// The reason also goes straight to stderr. The tracing appenders are
+/// non-blocking workers, and `exit` stops them before they flush, so without
+/// this write a supervisor's journal never learns why the process stopped. The
+/// write result is ignored because `eprintln!` would panic on a broken stderr.
+pub fn fatal(reason: FatalReason, message: &str) -> ! {
+ tracing::error!(
+ target: FATAL_LOG_TARGET,
+ reason = ?reason,
+ exit_status = reason.exit_status(),
+ "{message}"
+ );
+ let _ = writeln!(
+ std::io::stderr().lock(),
+ "iggy fatal: reason={reason:?} exit_status={}: {message}",
+ reason.exit_status()
+ );
+ std::process::exit(i32::from(reason.exit_status()));
+}
+
+/// Storage results that record a missing file descriptor for the exit status.
+///
+/// Only for opens that write, create, truncate or sync, and for read-only
+/// opens that are one step of such a write: the existence probes of segment
+/// setup, or an open that exists only to sync. The error still goes to the
+/// caller, whose own handling decides what happens. A failed partition write
+/// fences the partition and stops the server through the shutdown flush, and a
+/// failed partition build retries. Stopping at the open instead would skip
that
+/// flush and lose committed messages that are not yet in a segment.
+///
+/// If the server then stops on an error, it exits with
+/// [`FatalReason::DescriptorsExhausted`] instead of 1. The record lasts for
the
+/// life of the process. A failed read fails only its own request, so reads do
+/// not record it, and accept loops do not either, see `message_bus::accept`.
+pub trait NoteDescriptorExhaustion: Sized {
+ /// Pass the result through. If it failed with `EMFILE` or `ENFILE`, record
+ /// that for [`descriptors_exhausted`], and log the first one with the
+ /// limits. `operation` names what was being opened, for that log line.
+ #[must_use]
+ fn note_descriptor_exhaustion(self, operation: impl FnOnce() -> String) ->
Self;
+}
+
+static DESCRIPTORS_EXHAUSTED: AtomicBool = AtomicBool::new(false);
+
+impl<T> NoteDescriptorExhaustion for io::Result<T> {
+ fn note_descriptor_exhaustion(self, operation: impl FnOnce() -> String) ->
Self {
+ if let Err(error) = &self
+ && is_descriptor_exhaustion(error)
+ && !DESCRIPTORS_EXHAUSTED.swap(true, Ordering::Relaxed)
+ {
+ tracing::error!(
+ target: FATAL_LOG_TARGET,
+ "no free file descriptor while {}: {error} ({}); if the server
stops \
Review Comment:
nit: this log line promises exit 4, but a wedged superblock create still
exits 3 through `fatal` at `shard/src/lib.rs:7752`. say here and in the doc at
line 96 that `fatal` exits keep their own status.
##########
core/journal/src/prepare_journal.rs:
##########
@@ -819,7 +820,9 @@ impl Journal for PrepareJournal {
let tmp_path = wal_path.with_extension("wal.tmp");
Review Comment:
simplification: `truncate_from` and `drain` at line 987 repeat the same tmp
write, rename, dir sync and reopen sequence, so every new note lands twice.
pull it into one `rewrite_live(&live, label)` helper both call - about 35 lines
less.
##########
core/server/src/http.rs:
##########
@@ -282,12 +292,60 @@ pub fn start(
/// `ConnectInfo<ClientAddr>` extension per connection instead. Handlers read
/// it through the `Identity` extractor's `client_ip`, which picks the
/// advertised address a client is told about - never authorization.
-#[derive(Debug, Clone, Copy)]
-pub struct ClientAddr(pub SocketAddr);
+///
+/// It also holds the connection's slot in the node's connection cap. The
+/// per-connection service owns this value, so the slot frees when hyper drops
+/// the connection.
+#[derive(Debug, Clone)]
+pub struct ClientAddr {
+ pub addr: SocketAddr,
+ _permit: Option<Arc<ConnectionPermit>>,
+}
+
+impl Connected<cyper_axum::IncomingStream<'_, CappedListener>> for ClientAddr {
+ fn connect_info(stream: cyper_axum::IncomingStream<'_, CappedListener>) ->
Self {
+ let peer = stream.remote_addr();
+ Self {
+ addr: peer.addr,
+ _permit: peer.permit.clone(),
+ }
+ }
+}
+
+/// The plain HTTP listener behind the node's connection cap. A socket past
+/// the cap is closed at accept.
+pub struct CappedListener {
+ listener: TcpListener,
+ connections: Rc<ConnectionCap>,
+}
+
+/// The peer of an accepted plain HTTP socket and its slot in the connection
+/// cap. `permit` is `None` only for the listener's own local address.
+#[derive(Debug, Clone)]
+pub struct CappedPeer {
Review Comment:
simplification: `CappedPeer` repeats the two fields of `ClientAddr` and
`connect_info` just copies them across. use `type Addr = ClientAddr` and return
`stream.remote_addr().clone()` - drops the type and about 10 lines.
##########
core/message_bus/src/connection_cap.rs:
##########
@@ -0,0 +1,161 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Node-wide cap on open client sockets.
+//!
+//! Shard 0 accepts every client socket, but the owning shard closes it. Each
+//! socket therefore carries a [`ConnectionPermit`] to its owning shard, and
+//! the permit frees its slot on drop.
+//!
+//! The count is a process-wide static, because one node runs per process. A
+//! permit that pointed at its count would add 16 bytes to the client setup
+//! frames, and every slot of every shard inbox pays for the largest frame.
+
+use std::cell::Cell;
+use std::sync::atomic::{AtomicUsize, Ordering};
+use std::time::{Duration, Instant};
+use tracing::warn;
+
+/// Shortest interval between two refusal warnings. A client that reconnects
+/// in a loop would otherwise log one line per attempt.
+const REFUSAL_LOG_INTERVAL: Duration = Duration::from_secs(10);
+
+static LIVE_CONNECTIONS: AtomicUsize = AtomicUsize::new(0);
+
+/// The open client sockets of the node, capped at `max`.
+///
+/// Owned by shard 0, where every client accept runs.
+#[derive(Debug)]
+pub struct ConnectionCap {
+ max: Option<usize>,
+ last_refusal_log: Cell<Option<Instant>>,
+ refused_since_log: Cell<u64>,
+}
+
+impl ConnectionCap {
+ /// A cap of `max` open sockets. `None` counts sockets and refuses none.
+ #[must_use]
+ pub const fn new(max: Option<usize>) -> Self {
+ Self {
+ max,
+ last_refusal_log: Cell::new(None),
+ refused_since_log: Cell::new(0),
+ }
+ }
+
+ /// Sockets of this process that hold a permit now.
+ #[must_use]
+ pub fn live(&self) -> usize {
+ LIVE_CONNECTIONS.load(Ordering::Relaxed)
+ }
+
+ /// Take a slot for a socket that was just accepted, or `None` at the
+ /// cap. The caller closes a refused socket by dropping it.
+ pub fn try_acquire(&self) -> Option<ConnectionPermit> {
+ if let Some(max) = self.max {
+ let admitted =
+ LIVE_CONNECTIONS.fetch_update(Ordering::Relaxed,
Ordering::Relaxed, |live| {
+ (live < max).then_some(live + 1)
+ });
+ if let Err(live) = admitted {
+ self.log_refusal(live, max);
+ return None;
+ }
+ } else {
+ LIVE_CONNECTIONS.fetch_add(1, Ordering::Relaxed);
+ }
+ Some(ConnectionPermit { _private: () })
+ }
+
+ fn log_refusal(&self, live: usize, max: usize) {
+ let refused = self.refused_since_log.get() + 1;
+ let now = Instant::now();
+ if self
+ .last_refusal_log
+ .get()
+ .is_some_and(|last| now.duration_since(last) <
REFUSAL_LOG_INTERVAL)
+ {
+ self.refused_since_log.set(refused);
+ return;
+ }
+ self.last_refusal_log.set(Some(now));
+ self.refused_since_log.set(0);
+ warn!(
Review Comment:
nit: this warning does not name `message_bus.connections_max`, so the
operator cannot tell which setting to raise. name the key, as most boot lines
in `boot/fd_limit.rs` do.
##########
core/message_bus/src/accept.rs:
##########
@@ -0,0 +1,45 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Accept-loop handling of descriptor exhaustion.
+
+use server_common::fatal::is_descriptor_exhaustion;
+use std::io;
+use std::time::Duration;
+
+/// How long an accept loop waits after `EMFILE` or `ENFILE`.
+///
+/// The kernel keeps the connection in the listen backlog, so an immediate
+/// retry fails again and spins the shard at full CPU while the descriptor
+/// table stays full.
+pub const DESCRIPTOR_EXHAUSTION_BACKOFF: Duration = Duration::from_secs(1);
+
+/// Wait [`DESCRIPTOR_EXHAUSTION_BACKOFF`] if `error` says no file descriptor
+/// is free, and return at once for any other `accept()` error.
+///
+/// An accept loop only waits. It does not record the exhaustion for the exit
+/// status as the storage write paths do, because no write failed.
+/// `[message_bus] connections_max` keeps clients below the descriptor limit,
+/// so an accept reaches `EMFILE` only when that cap is off or set too high, or
+/// when storage holds the descriptors. A process exit here would let clients
+/// stop the node in the first two cases.
+#[allow(clippy::future_not_send)]
+pub async fn pause_after_accept_error(error: &io::Error) {
+ if is_descriptor_exhaustion(error) {
Review Comment:
nit: as before this PR, a persistent `ENOMEM` or `ENOBUFS` from accept still
spins shard 0 with one error line per retry. pause on those too, the memory
errors that `accept(2)` documents.
##########
core/server/src/main.rs:
##########
@@ -99,6 +116,13 @@ fn main() -> Result<(), ServerError> {
if let Err(error) = &joined {
server::boot::systemd::notify_shutdown_failure(error);
}
+ if let Err(error) = &joined
+ && descriptors_exhausted()
+ {
+ // `fatal` skips destructors, and the log appenders flush on drop.
+ drop(logging);
Review Comment:
nit: `drop(logging)` runs before `fatal` emits its event, so the exit 4
reason never reaches the log file or stdout. emit the event first, then drop
logging, write stderr and exit.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]