hubcio commented on code in PR #4318:
URL: https://github.com/apache/iggy/pull/4318#discussion_r4136702515
##########
core/server/src/sysinfo_probe.rs:
##########
@@ -105,32 +117,52 @@ pub fn stats_disk_space() -> (u64, u64) {
})
}
+/// Count open descriptors, scanning where the kernel keeps no count, and
+/// publish the result for [`probe_system_stats`].
+///
+/// Only the sysinfo printer calls this, once per interval. The scan cost
+/// grows with the count, so `GetStats` reads the published value instead of
+/// scanning on the request path. The value is as old as the printer interval,
+/// and stays 0 while the printer is disabled.
+pub fn publish_open_files_count() {
+ PUBLISHED_OPEN_FILES_COUNT.store(count_open_files().unwrap_or(0),
Ordering::Relaxed);
+}
+
+/// Probe through this thread's [`SYSINFO`], the `GetStats` path.
pub fn probe_system_stats() -> SystemStats {
- let host = HOST_IDENTITY.get_or_init(HostIdentity::probe);
- let probe = SYSINFO.with_borrow_mut(|slot| {
- let sys = slot.get_or_insert_with(SysinfoSystem::new);
- SystemProbe::capture(sys)
- });
-
- SystemStats {
- process_id: probe.process_id,
- cpu_usage: probe.cpu_usage,
- total_cpu_usage: probe.total_cpu_usage,
- memory_usage: probe.memory_usage,
- total_memory: probe.total_memory,
- available_memory: probe.available_memory,
- // sysinfo reports whole seconds; the wire fields are micros (the
- // SDK decodes them via `IggyDuration` / `IggyTimestamp::from`, both
- // micro-based).
- run_time: probe.run_time_secs.saturating_mul(1_000_000),
- start_time: probe.start_time_secs.saturating_mul(1_000_000),
- read_bytes: probe.read_bytes,
- written_bytes: probe.written_bytes,
- threads_count: probe.threads_count,
- hostname: host.hostname.clone(),
- os_name: host.os_name.clone(),
- os_version: host.os_version.clone(),
- kernel_version: host.kernel_version.clone(),
+ SYSINFO
+ .with_borrow_mut(|slot|
SystemStats::capture(slot.get_or_insert_with(SysinfoSystem::new)))
+}
+
+impl SystemStats {
+ /// Probe through `sys`. A caller that keeps its own `sys` gets CPU deltas
+ /// over its own interval, and does not reset the window of `GetStats`.
+ pub fn capture(sys: &mut SysinfoSystem) -> Self {
+ let host = HOST_IDENTITY.get_or_init(HostIdentity::probe);
+ let probe = SystemProbe::capture(sys);
+ Self {
+ process_id: probe.process_id,
+ cpu_usage: probe.cpu_usage,
+ total_cpu_usage: probe.total_cpu_usage,
+ memory_usage: probe.memory_usage,
+ total_memory: probe.total_memory,
+ available_memory: probe.available_memory,
+ // sysinfo reports whole seconds; the wire fields are micros (the
+ // SDK decodes them via `IggyDuration` / `IggyTimestamp::from`,
both
+ // micro-based).
+ run_time: probe.run_time_secs.saturating_mul(1_000_000),
+ start_time: probe.start_time_secs.saturating_mul(1_000_000),
+ read_bytes: probe.read_bytes,
+ written_bytes: probe.written_bytes,
+ threads_count: probe.threads_count,
+ open_files_count: count_open_files_without_scan()
+ .unwrap_or_else(||
PUBLISHED_OPEN_FILES_COUNT.load(Ordering::Relaxed)),
+ open_files_limit: getrlimit(Resource::RLIMIT_NOFILE).map_or(0,
|(soft, _)| soft),
Review Comment:
warning: remove `open_files_limit` from stats - it's set once at boot and
static for the whole runtime, so there's no sense in having it there. drop it
before the SDKs ship the 16-byte tail, or removing it later breaks clients.
##########
core/metadata/src/impls/metadata.rs:
##########
@@ -2220,6 +2236,61 @@ where
}
}
+ /// Node-wide `[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, and creates in flight
+ /// together can overshoot it. 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> {
+ let partitions_max = self.partitions_max.get();
+ if partitions_max == 0 {
+ return Ok(());
+ }
+ let header = message.header();
+ let body =
&message.as_slice()[size_of::<RoutedRequestHeader>()..header.size as usize];
+ let (requested, stream_id, topic_id) = match header.operation {
+ Operation::CreateTopic => match
WireCreateTopicRequest::decode_from(body) {
+ Ok(request) => (request.partitions_count, request.stream_id,
None),
+ Err(_) => return Ok(()),
+ },
+ Operation::CreatePartitions => match
WireCreatePartitionsRequest::decode_from(body) {
+ Ok(request) => (
+ request.partitions_count,
+ request.stream_id,
+ Some(request.topic_id),
+ ),
+ Err(_) => return Ok(()),
+ },
+ _ => return Ok(()),
+ };
+ let admitted = validate_partitions_limit(partitions_max, requested, ||
{
Review Comment:
warning: this compares each create against committed partitions only, so
in-flight creates don't add up - 96 queued creates of 1000 partitions each can
all pass. count queued creates too, or state the real bound in `config.toml`
instead of "slightly".
##########
core/server/config.toml:
##########
@@ -407,6 +407,11 @@ rotation_check_interval = "1 h"
# Time to retain log files before deletion. Avoid less than 1s, too.
retention = "7 days"
+# Interval for printing process and host usage (CPU, memory, disk, clients,
Review Comment:
warning: on kernels before 6.2 and on macOS only the printer publishes
`open_files_count`, so `sysinfo_print_interval = "0 s"` silently pins it at 0
in `GetStats`. document that here and in the SDK field docs, or publish the
count on its own timer.
##########
core/server/src/sysinfo_probe.rs:
##########
@@ -105,32 +117,52 @@ pub fn stats_disk_space() -> (u64, u64) {
})
}
+/// Count open descriptors, scanning where the kernel keeps no count, and
+/// publish the result for [`probe_system_stats`].
+///
+/// Only the sysinfo printer calls this, once per interval. The scan cost
+/// grows with the count, so `GetStats` reads the published value instead of
+/// scanning on the request path. The value is as old as the printer interval,
+/// and stays 0 while the printer is disabled.
+pub fn publish_open_files_count() {
+ PUBLISHED_OPEN_FILES_COUNT.store(count_open_files().unwrap_or(0),
Ordering::Relaxed);
Review Comment:
nit: on kernels before 6.2 the scan needs a free fd, so right at exhaustion
it fails and this stores 0, which reads as "unknown". keep the previous count
when `count_open_files()` returns `None`.
##########
core/server_common/src/fatal.rs:
##########
@@ -0,0 +1,163 @@
+// 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};
+
+/// 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,
+ /// A storage open for a write or a sync failed because the process
+ /// (`EMFILE`) or the host (`ENFILE`) has no free file descriptor. The
shard
+ /// cannot persist that write, and serving on only fails each later write
in
+ /// turn.
+ DescriptorsExhausted = 4,
+}
+
+impl FatalReason {
+ #[must_use]
+ pub const fn exit_status(self) -> u8 {
+ self as u8
+ }
+}
+
+/// 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: "iggy.fatal.diag",
+ 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 stop the process when no file descriptor is free.
+///
+/// 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. A failure there leaves the
+/// write half done. Other reads keep their own error path: a failed read fails
+/// only its request, and the descriptors that fill the table can be client
+/// sockets or other reads that close again soon.
+/// Accept loops never stop here, see `message_bus::accept`.
+pub trait ExitOnDescriptorExhaustion: Sized {
+ /// Pass the result through, unless it failed with `EMFILE` or `ENFILE`:
then
+ /// stop the process with [`FatalReason::DescriptorsExhausted`].
`operation`
+ /// names what was being opened, for the fatal message.
+ #[must_use]
+ fn exit_on_descriptor_exhaustion(self, operation: impl FnOnce() -> String)
-> Self;
+}
+
+impl<T> ExitOnDescriptorExhaustion for io::Result<T> {
+ fn exit_on_descriptor_exhaustion(self, operation: impl FnOnce() -> String)
-> Self {
+ if let Err(error) = &self
+ && is_descriptor_exhaustion(error)
+ {
+ exit_descriptors_exhausted(error, &operation());
Review Comment:
critical: exiting here skips the shutdown flush `FatalCommit` used to run -
`EMFILE` doesn't stop writes to open fds, so under the default
replicated/replicated policy we drop acked messages that could be flushed.
return the error and exit 4 after the shutdown flush.
##########
core/metadata/src/impls/metadata.rs:
##########
@@ -4260,19 +4343,44 @@ fn unreplayable_secret_refusal(
)
}
+/// `requested` more partitions on top of the `committed` count against the
+/// `partitions_max` cap. Zero is no cap, and `committed` is read only when a
+/// cap is set. A create of zero partitions adds nothing, so it passes also
+/// on a node already past the cap, for example after the cap was lowered.
+fn validate_partitions_limit(
+ partitions_max: u32,
+ requested: u32,
+ committed: impl FnOnce() -> usize,
+) -> Result<(), IggyError> {
+ if partitions_max == 0 || requested == 0 {
Review Comment:
simplification: the `partitions_max == 0` return repeats line 2250, and the
lazy closure buys nothing on this cold path. keep that check at 2250 only and
pass the count as a plain `usize`.
##########
core/message_bus/src/connection_cap.rs:
##########
@@ -0,0 +1,166 @@
+// 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),
+ }
+ }
+
+ #[must_use]
+ pub const fn max(&self) -> Option<usize> {
Review Comment:
simplification: `ConnectionCap::max()` has no callers, tests included. drop
it.
##########
core/server_common/src/segment_storage/messages_writer.rs:
##########
@@ -56,6 +57,7 @@ impl MessagesWriter {
let file = opts
.open(file_path)
.await
+ .exit_on_descriptor_exhaustion(|| format!("opening {file_path}"))
Review Comment:
warning: a create that crosses the fd limit commits on every replica, then
this open exits all of them and the cluster loses quorum. derive a nonzero
`partitions_max` default from the soft limit, or keep the old reconciler error
here.
##########
core/server_common/src/fatal.rs:
##########
@@ -0,0 +1,163 @@
+// 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};
+
+/// 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,
+ /// A storage open for a write or a sync failed because the process
+ /// (`EMFILE`) or the host (`ENFILE`) has no free file descriptor. The
shard
+ /// cannot persist that write, and serving on only fails each later write
in
+ /// turn.
+ DescriptorsExhausted = 4,
+}
+
+impl FatalReason {
+ #[must_use]
+ pub const fn exit_status(self) -> u8 {
+ self as u8
+ }
+}
+
+/// 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: "iggy.fatal.diag",
Review Comment:
nit: this renames the fatal log target from `iggy.consensus.diag` (what
0.9.0 ships) to `iggy.fatal.diag`, so existing log filters stop matching. keep
the old target or call out the rename in the release notes.
##########
core/server/src/sysinfo_printer.rs:
##########
@@ -0,0 +1,205 @@
+// 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, publish_open_files_count,
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) {
+ // `GetStats` reads this count where the kernel keeps none, so it is
+ // published also while the line itself is filtered out.
+ publish_open_files_count();
Review Comment:
warning: on kernels before 6.2 this scans `/proc/self/fd` on shard 0 every
tick - a readdir over 40k fds took 13-28 ms on my box. move the scan off the
shard thread, or report 0 where the kernel has no count.
##########
core/server_common/src/fatal.rs:
##########
@@ -0,0 +1,163 @@
+// 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};
+
+/// 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,
+ /// A storage open for a write or a sync failed because the process
+ /// (`EMFILE`) or the host (`ENFILE`) has no free file descriptor. The
shard
+ /// cannot persist that write, and serving on only fails each later write
in
+ /// turn.
+ DescriptorsExhausted = 4,
+}
+
+impl FatalReason {
+ #[must_use]
+ pub const fn exit_status(self) -> u8 {
+ self as u8
+ }
+}
+
+/// 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: "iggy.fatal.diag",
+ 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 stop the process when no file descriptor is free.
+///
+/// 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. A failure there leaves the
+/// write half done. Other reads keep their own error path: a failed read fails
+/// only its request, and the descriptors that fill the table can be client
+/// sockets or other reads that close again soon.
+/// Accept loops never stop here, see `message_bus::accept`.
+pub trait ExitOnDescriptorExhaustion: Sized {
+ /// Pass the result through, unless it failed with `EMFILE` or `ENFILE`:
then
+ /// stop the process with [`FatalReason::DescriptorsExhausted`].
`operation`
+ /// names what was being opened, for the fatal message.
+ #[must_use]
+ fn exit_on_descriptor_exhaustion(self, operation: impl FnOnce() -> String)
-> Self;
+}
+
+impl<T> ExitOnDescriptorExhaustion for io::Result<T> {
+ fn exit_on_descriptor_exhaustion(self, operation: impl FnOnce() -> String)
-> Self {
+ if let Err(error) = &self
+ && is_descriptor_exhaustion(error)
+ {
+ exit_descriptors_exhausted(error, &operation());
+ }
+ self
+ }
+}
+
+/// `EMFILE` (this process is at its `RLIMIT_NOFILE`) or `ENFILE` (the host is
+/// at its system-wide file table limit).
+#[must_use]
+pub fn is_descriptor_exhaustion(error: &io::Error) -> bool {
+ error
+ .raw_os_error()
+ .is_some_and(|code| code == Errno::EMFILE as i32 || code ==
Errno::ENFILE as i32)
+}
+
+fn exit_descriptors_exhausted(error: &io::Error, operation: &str) -> ! {
+ let limits = getrlimit(Resource::RLIMIT_NOFILE).map_or_else(
+ |errno| format!("RLIMIT_NOFILE unreadable: {errno}"),
+ |(soft, hard)| format!("RLIMIT_NOFILE soft={soft} hard={hard}"),
+ );
+ fatal(
+ FatalReason::DescriptorsExhausted,
+ &format!("no free file descriptor while {operation}: {error}
({limits})"),
+ )
+}
+
+#[cfg(test)]
Review Comment:
warning: the tests only cover errno classification and pass-through -
nothing proves the exit code 4 or that a restart after it recovers the data.
add an integration test that lowers the hard `RLIMIT_NOFILE` and asserts both.
##########
core/partitions/src/install_backup.rs:
##########
@@ -135,9 +136,11 @@ async fn link_tree<S: DurableStorage>(
// Transfer unlinks or atomically replaces these frozen files.
// Hard links retain the old bytes without copying segment
data.
storage.hard_link(&source.join(&name), &destination).await?;
+ // Read-only, but only to sync, so it stops like a write open.
storage
.open(&destination, OpenMode::Read)
- .await?
+ .await
+ .exit_on_descriptor_exhaustion(|| format!("opening {}",
destination.display()))?
Review Comment:
nit: this read open exits on `EMFILE` like a write open, but `copy_file`
returns `EMFILE` from its source open and exits on the destination open, so the
outcome depends on which open fails first. pick one rule for read-then-write
pairs.
##########
core/message_bus/src/replica/listener.rs:
##########
@@ -109,6 +110,7 @@ pub async fn run(listener: TcpListener, token:
ShutdownToken, on_accepted: Accep
}
Err(e) => {
error!("Replica listener accept failed: {e}");
+ pause_after_accept_error(&e).await;
Review Comment:
nit: the doc on `run` (line 92) says the loop never awaits anything but
`accept()`, but it now sleeps here after `EMFILE`/`ENFILE`. mention the pause
there.
##########
core/server/config.toml:
##########
@@ -855,6 +860,12 @@ journal_slots = 1024
# it lifts both. Must be between 2 and 65536.
clients_table_max = 8192
+# Cap on partitions across all streams and topics of the node. A request to
Review Comment:
nit: this reads like a per-node cap, but only the metadata primary checks
it, with its own value, for creates across the whole cluster. say that the
primary enforces it, so every node needs the same value.
##########
core/configs/src/common/system.rs:
##########
@@ -44,6 +45,10 @@ pub struct LoggingConfig {
#[config_env(leaf)]
#[serde_as(as = "DisplayFromStr")]
pub retention: IggyDuration,
+ #[config_env(leaf)]
+ #[serde_as(as = "DisplayFromStr")]
+ #[serde(default = "default_sysinfo_print_interval")]
+ pub sysinfo_print_interval: IggyDuration,
Review Comment:
nit: validation floors `retention` and `rotation_check_interval` at 1 s but
lets this through at any nonzero value, so `1 ms` rescans `/proc/self/fd` every
millisecond on older kernels. reject nonzero values below 1 s in
`LoggingConfig::validate` too.
##########
.github/actions/go/pre-merge/action.yml:
##########
@@ -85,7 +85,7 @@ runs:
if: inputs.task == 'lint'
uses: golangci/golangci-lint-action@v9
with:
- version: v2.11.3
+ version: v2.13.2
Review Comment:
nit: the PR body mentions the golangci-lint bump to v2.13.2 (six pins) but
not why. if it's unrelated to fd limits, move it to its own chore commit.
##########
core/system_stats/src/lib.rs:
##########
@@ -106,6 +106,24 @@ impl SystemProbe {
}
}
+/// Descriptors the calling process holds open, or `None` where sysinfo
+/// cannot count them.
+///
+/// Kept out of [`SystemProbe::capture`] because the fallback cost grows with
+/// the count: Linux before 6.2 lists every entry of `/proc/self/fd`, and
+/// macOS copies the whole descriptor table. The Linux listing also counts
+/// the descriptor of the scan itself.
+pub fn count_open_files() -> Option<u64> {
+ count_open_files_without_scan().or_else(scan_open_files)
+}
+
+/// The open-descriptor count where the kernel keeps it, or `None` where only
+/// the scan in [`count_open_files`] can find it. Cheap enough for a request
+/// path: Linux 6.2 and later count the descriptor bitmap for one `stat`.
+pub fn count_open_files_without_scan() -> Option<u64> {
Review Comment:
simplification: `count_open_files_without_scan` only forwards to the private
`kernel_open_files_count`. give the two `cfg` variants the public name and drop
the wrapper.
##########
core/partitions/src/state_transfer.rs:
##########
@@ -3875,6 +3908,53 @@ impl fmt::Display for SegmentLoadError {
impl std::error::Error for SegmentLoadError {}
+/// A failed read inside a served segment, with where it failed.
+///
+/// It travels as an `io::Error` of the read's own kind and keeps the OS error
+/// as its source, so [`SegmentLoadError::classify`] still tells a dying disk
+/// from a stale offer. A message-only rewrap hid `EIO` behind `Other`.
+#[derive(Debug)]
+struct SegmentReadError {
Review Comment:
simplification: `SegmentReadError` only adds the read position to the OS
error, at the cost of a struct and three trait impls. return the raw
`io::Error` from the read at line 3801 and log the position there.
--
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]