hubcio commented on code in PR #4318:
URL: https://github.com/apache/iggy/pull/4318#discussion_r4144655584


##########
foreign/swift/Sources/Iggy/Models/Resources.swift:
##########
@@ -474,13 +474,17 @@ public struct Stats: Sendable, Hashable {
     public var threadsCount: UInt32
     public var freeDiskSpace: UInt64
     public var totalDiskSpace: UInt64
+    /// The number of file descriptors the server process holds open, 0 when 
unknown.
+    public var openFilesCount: UInt64
+    /// The soft limit on open file descriptors (`RLIMIT_NOFILE`) of the 
server process, 0 when unknown.
+    public var openFilesLimit: UInt64
 
     public init(
         processID: UInt32, cpuUsage: Float, totalCPUUsage: Float, memoryUsage: 
UInt64, totalMemory: UInt64, availableMemory: UInt64, runTime: Duration,
         startTime: IggyTimestamp, readBytes: UInt64, writtenBytes: UInt64, 
messagesSizeBytes: UInt64, streamsCount: UInt32, topicsCount: UInt32,
         partitionsCount: UInt32, segmentsCount: UInt32, messagesCount: UInt64, 
clientsCount: UInt32, consumerGroupsCount: UInt32, hostname: String,
         osName: String, osVersion: String, kernelVersion: String, 
iggyServerVersion: String, iggyServerSemver: UInt32?, cacheMetrics: 
[CacheMetrics],
-        threadsCount: UInt32, freeDiskSpace: UInt64, totalDiskSpace: UInt64
+        threadsCount: UInt32, freeDiskSpace: UInt64, totalDiskSpace: UInt64, 
openFilesCount: UInt64 = 0, openFilesLimit: UInt64 = 0

Review Comment:
   nit: only `openFilesCount` and `openFilesLimit` get a `= 0` default, so a 
future decoder can skip the tail and silently report 0. drop the two defaults 
to match the other 28 parameters.



##########
core/message_bus/src/installer/mod.rs:
##########
@@ -99,8 +100,16 @@ pub trait ConnectionInstaller {
     /// Same for an SDK client connection. The owning shard is already
     /// encoded in the top 16 bits of `meta.client_id`. `meta` is stored
     /// on the bus and exposed via [`IggyMessageBus::client_meta`] for
-    /// the lifetime of the connection.
-    fn install_client_fd(&self, fd: DupedFd, meta: ClientConnMeta, on_request: 
RequestHandler);
+    /// the lifetime of the connection. `permit` is the socket's slot in
+    /// the node's connection cap, and the install holds it until the socket
+    /// closes.
+    fn install_client_fd(
+        &self,
+        fd: DupedFd,
+        meta: ClientConnMeta,
+        permit: Option<ConnectionPermit>,

Review Comment:
   nit: `Option<ConnectionPermit>` lets a new caller pass `None` and install an 
uncounted client socket, though production always passes `Some`. take 
`ConnectionPermit` in the four client fd methods and build test permits like 
`test_permit` in `shard/src/coordinator.rs`.



##########
core/binary_protocol/src/responses/system/get_stats.rs:
##########
@@ -55,6 +55,8 @@ impl CacheMetricEntry {
 /// [threads_count:4]
 /// [free_disk_space:8]
 /// [total_disk_space:8]
+/// [open_files_count:8]

Review Comment:
   nit: this lists `open_files_count` and `open_files_limit` as always present, 
but `decode` treats them as an optional tail. say that an older server ends the 
reply at `total_disk_space` and a partial tail fails.



##########
core/server_common/src/fatal.rs:
##########
@@ -0,0 +1,243 @@
+// 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 iggy_common::IggyTimestamp;
+use nix::errno::Errno;
+use nix::sys::resource::{Resource, getrlimit};
+use std::io::{self, Write};
+use std::sync::atomic::{AtomicU64, 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`].
+    ///
+    /// The other reasons keep their own status after such a failure, for
+    /// example a superblock that stays unwritable for lack of a descriptor
+    /// exits 3. Their exit line then gives the time of the first one.
+    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.
+///
+/// If a storage open found no free file descriptor earlier, the line gives the
+/// time of the first one, whatever the reason. See 
[`NoteDescriptorExhaustion`].
+pub fn fatal(reason: FatalReason, message: &str) -> ! {
+    fatal_with_log_flush(reason, message, || {});
+}
+
+/// [`fatal`] for the owner of the log appenders. `flush_logs` runs after the
+/// log event and before the exit, so the event reaches the log files too.
+pub fn fatal_with_log_flush(reason: FatalReason, message: &str, flush_logs: 
impl FnOnce()) -> ! {
+    let exhaustion = first_descriptor_exhaustion()
+        .map(|at| format!("; a storage open first found no free file 
descriptor at {at} UTC"))
+        .unwrap_or_default();
+    tracing::error!(
+        target: FATAL_LOG_TARGET,
+        reason = ?reason,
+        exit_status = reason.exit_status(),
+        "{message}{exhaustion}"
+    );
+    flush_logs();
+    let _ = writeln!(
+        std::io::stderr().lock(),
+        "iggy fatal: reason={reason:?} exit_status={}: {message}{exhaustion}",
+        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

Review Comment:
   nit: this says a failed partition build retries, but at boot 
`boot/recovery.rs:240` returns the error and stops the server - only 
`PartitionOffsetReservationClaim` goes to the reconciler. reword it to match.



##########
core/server_common/src/segment_storage/mod.rs:
##########
@@ -101,31 +106,22 @@ impl SegmentStorage {
         indexes_size: u64,
         file_exists: bool,
     ) -> Result<Self, IggyError> {
-        let size = Rc::new(std::sync::atomic::AtomicU64::new(messages_size));
-        let indexes_size = 
Rc::new(std::sync::atomic::AtomicU64::new(indexes_size));
-        let messages_writer = Rc::new(MessagesWriter::new(messages_path, size, 
file_exists).await?);
-
-        let index_writer = Rc::new(IndexWriter::new(index_path, indexes_size, 
file_exists).await?);
-
-        if file_exists {
-            messages_writer.fsync().await?;
-            index_writer.fsync().await?;
-        }
-
+        prepare_for_writes(messages_path, "messages", messages_size, 
file_exists).await?;
+        prepare_for_writes(index_path, "index", indexes_size, 
file_exists).await?;
         let messages_reader = 
Rc::new(MessagesReader::new(messages_path).await?);

Review Comment:
   simplification: `prepare_for_writes` just opened both files, so these reader 
opens only prove they exist. add `read(true)` at line 147, give `IndexReader` a 
`from_validated_path` like `MessagesReader`, and skip the second open.
   
   also at line 98.



##########
core/server_common/src/fatal.rs:
##########
@@ -0,0 +1,243 @@
+// 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 iggy_common::IggyTimestamp;
+use nix::errno::Errno;
+use nix::sys::resource::{Resource, getrlimit};
+use std::io::{self, Write};
+use std::sync::atomic::{AtomicU64, 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`].
+    ///
+    /// The other reasons keep their own status after such a failure, for
+    /// example a superblock that stays unwritable for lack of a descriptor
+    /// exits 3. Their exit line then gives the time of the first one.
+    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.
+///
+/// If a storage open found no free file descriptor earlier, the line gives the
+/// time of the first one, whatever the reason. See 
[`NoteDescriptorExhaustion`].
+pub fn fatal(reason: FatalReason, message: &str) -> ! {
+    fatal_with_log_flush(reason, message, || {});
+}
+
+/// [`fatal`] for the owner of the log appenders. `flush_logs` runs after the
+/// log event and before the exit, so the event reaches the log files too.
+pub fn fatal_with_log_flush(reason: FatalReason, message: &str, flush_logs: 
impl FnOnce()) -> ! {
+    let exhaustion = first_descriptor_exhaustion()
+        .map(|at| format!("; a storage open first found no free file 
descriptor at {at} UTC"))
+        .unwrap_or_default();
+    tracing::error!(
+        target: FATAL_LOG_TARGET,
+        reason = ?reason,
+        exit_status = reason.exit_status(),
+        "{message}{exhaustion}"
+    );
+    flush_logs();
+    let _ = writeln!(
+        std::io::stderr().lock(),
+        "iggy fatal: reason={reason:?} exit_status={}: {message}{exhaustion}",
+        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. A stop through
+/// [`fatal`] keeps the status of its own reason. A failed read fails only its
+/// own request, so reads do not record it, and accept loops do not either, see
+/// `message_bus::accept`.
+///
+/// The record lasts for the life of the process. An exhaustion that the server
+/// recovers from, such as a superblock write that succeeds on a retry, 
therefore
+/// also sets the status of a later stop on an unrelated error. The exit line
+/// gives the time of the first exhaustion, so the two can be told apart.
+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;
+}
+
+/// Microseconds since the Unix epoch of the first noted exhaustion, or 0 for
+/// none.
+static FIRST_DESCRIPTOR_EXHAUSTION: AtomicU64 = AtomicU64::new(0);

Review Comment:
   simplification: a `OnceLock<IggyTimestamp>` does the same job without the 0 
sentinel, `.max(1)` and `compare_exchange` - 
`set(IggyTimestamp::now()).is_ok()` keeps the first-writer rule and 
`get().copied()` replaces the load.



##########
core/partitions/src/state_transfer.rs:
##########
@@ -645,6 +646,24 @@ const fn validate_consumer_offset_transfer_count(
 mod tests {
     use super::*;
 
+    #[test]
+    fn 
given_segment_read_error_when_classifying_should_tell_local_faults_from_stale_offers()
 {

Review Comment:
   nit: this only calls `classify`, which the PR did not change, so it passes 
without the `hash_segment_range` fix at line 3791. drive it through 
`load_verified_segment_artifact` with a directory path and assert `Stale` with 
`raw_os_error().is_some()`.



-- 
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]

Reply via email to