This is an automated email from the ASF dual-hosted git repository. mmodzelewski pushed a commit to branch sysinfo-improvements in repository https://gitbox.apache.org/repos/asf/iggy.git
commit 1a6c289cdef97de950ac09d97e0a0ab665772fb9 Author: Maciej Modzelewski <[email protected]> AuthorDate: Mon Sep 28 16:52:52 2026 +0200 feat(server): guard against running out of file descriptors The server did not raise its open-file limit. If a storage open failed with EMFILE or ENFILE, the shard kept serving and failed each later write. Accept loops retried at once and spun the shard at full CPU, because the kernel keeps the connection in the backlog. At startup, the server now raises the soft RLIMIT_NOFILE to the hard limit. On macOS, it uses OPEN_MAX if the hard limit is higher. If a storage open fails with EMFILE or ENFILE, the process now stops with exit status 4. Every later open fails too, so a supervisor restart is the repair. fatal() moves to server_common and writes to stderr, because exit stops the log appenders before they flush. Accept loops wait one second after EMFILE or ENFILE. They do not stop the process, because a client with many sockets can then stop the node. The new [metadata] partitions_max limits the partitions of a node. CreateTopic and CreatePartitions past the limit fail before consensus with PartitionsLimitReached (2022). The limit is soft, because concurrent creates can go over it. 0 means no limit. Shard 0 logs process and host usage, with open descriptors, every logging.sysinfo_print_interval (10 s by default, 0 disables it). Partition groups log the restored-view line at debug, because one INFO line per partition flooded the boot log. --- .github/actions/go/pre-merge/action.yml | 6 +- .pre-commit-config.yaml | 6 +- core/common/src/error/iggy_error.rs | 13 ++ core/configs/src/common/defaults.rs | 11 ++ core/configs/src/common/displays.rs | 5 +- core/configs/src/common/system.rs | 7 + core/configs/src/server_config/defaults.rs | 1 + core/configs/src/server_config/displays.rs | 7 +- core/configs/src/server_config/metadata.rs | 18 ++ core/consensus/src/fatal.rs | 60 ------- core/consensus/src/impls.rs | 23 ++- core/consensus/src/lib.rs | 3 +- core/integration/tests/server/mod.rs | 1 + .../tests/server/partition_view_durability_vsr.rs | 8 +- .../tests/server/partitions_limit_vsr.rs | 121 +++++++++++++ core/journal/src/durable_storage.rs | 29 ++- core/journal/src/file_storage.rs | 15 +- core/journal/src/prepare_journal.rs | 23 ++- core/journal/src/superblock.rs | 14 +- core/message_bus/src/accept.rs | 42 +++++ core/message_bus/src/client_listener/tcp.rs | 2 + core/message_bus/src/client_listener/tcp_tls.rs | 2 + core/message_bus/src/client_listener/ws.rs | 2 + core/message_bus/src/client_listener/wss.rs | 2 + core/message_bus/src/lib.rs | 1 + core/message_bus/src/replica/listener.rs | 2 + core/metadata/src/impls/metadata.rs | 25 ++- core/metadata/src/stm/stream.rs | 14 ++ core/partitions/src/iggy_index_reader.rs | 2 + core/partitions/src/iggy_index_writer.rs | 2 + core/partitions/src/iggy_partition.rs | 2 +- core/partitions/src/messages_writer.rs | 2 + core/partitions/src/offset_storage.rs | 9 +- core/partitions/src/poll_plan.rs | 11 +- core/partitions/src/segment_anchor.rs | 10 +- core/partitions/src/segment_recovery.rs | 105 ++++++----- core/partitions/src/state_transfer.rs | 24 ++- core/server/config.toml | 11 ++ core/server/src/boot/fd_limit.rs | 116 ++++++++++++ core/server/src/boot/mod.rs | 19 ++ core/server/src/boot/threads.rs | 2 + core/server/src/dispatch/mod.rs | 1 + core/server/src/http/handlers.rs | 25 ++- core/server/src/http/tls.rs | 6 +- core/server/src/lib.rs | 1 + core/server/src/main.rs | 22 ++- core/server/src/responses.rs | 53 ++++-- core/server/src/rewrite.rs | 63 ++++++- core/server/src/sysinfo_printer.rs | 197 +++++++++++++++++++++ core/server/src/sysinfo_probe.rs | 4 + core/server_common/src/fatal.rs | 160 +++++++++++++++++ core/server_common/src/fs_utils.rs | 6 +- core/server_common/src/lib.rs | 1 + .../src/segment_storage/index_reader.rs | 2 + .../src/segment_storage/index_writer.rs | 2 + .../src/segment_storage/messages_reader.rs | 2 + .../src/segment_storage/messages_writer.rs | 2 + core/server_common/src/segment_storage/mod.rs | 2 + core/system_stats/src/lib.rs | 11 ++ foreign/go/errors/errors.yaml | 4 + foreign/go/errors/errors_gen.go | 17 ++ .../org/apache/iggy/exception/IggyErrorCode.java | 1 + .../apache/iggy/exception/IggyErrorCodeTest.java | 1 + foreign/node/src/wire/error.code.test.ts | 4 + foreign/node/src/wire/error.code.ts | 1 + .../swift/Sources/Iggy/Errors/IggyErrorCode.swift | 2 + 66 files changed, 1188 insertions(+), 180 deletions(-) diff --git a/.github/actions/go/pre-merge/action.yml b/.github/actions/go/pre-merge/action.yml index 1d06e1da9..b06f680a5 100644 --- a/.github/actions/go/pre-merge/action.yml +++ b/.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 working-directory: foreign/go args: --timeout=5m skip-cache: true @@ -94,7 +94,7 @@ runs: if: inputs.task == 'lint' && hashFiles('bdd/go/go.mod') != '' uses: golangci/golangci-lint-action@v9 with: - version: v2.11.3 + version: v2.13.2 working-directory: bdd/go args: --timeout=5m skip-cache: true @@ -103,7 +103,7 @@ runs: if: inputs.task == 'lint' && hashFiles('examples/go/go.mod') != '' uses: golangci/golangci-lint-action@v9 with: - version: v2.11.3 + version: v2.13.2 working-directory: examples/go args: --timeout=5m skip-cache: true diff --git a/.pre-commit-config.yaml b/.pre-commit-config.yaml index 3cb1a7450..cfd967de8 100644 --- a/.pre-commit-config.yaml +++ b/.pre-commit-config.yaml @@ -233,7 +233,7 @@ repos: entry: bash -c 'cd foreign/go && golangci-lint run --timeout=5m ./...' language: golang additional_dependencies: - ["github.com/golangci/golangci-lint/v2/cmd/[email protected]"] + ["github.com/golangci/golangci-lint/v2/cmd/[email protected]"] files: ^foreign/go/.*\.go$ pass_filenames: false stages: [pre-push] @@ -243,7 +243,7 @@ repos: entry: bash -c 'cd bdd/go && golangci-lint run --timeout=5m ./...' language: golang additional_dependencies: - ["github.com/golangci/golangci-lint/v2/cmd/[email protected]"] + ["github.com/golangci/golangci-lint/v2/cmd/[email protected]"] files: ^bdd/go/.*\.go$ pass_filenames: false stages: [pre-push] @@ -253,7 +253,7 @@ repos: entry: bash -c 'cd examples/go && golangci-lint run --timeout=5m ./...' language: golang additional_dependencies: - ["github.com/golangci/golangci-lint/v2/cmd/[email protected]"] + ["github.com/golangci/golangci-lint/v2/cmd/[email protected]"] files: ^examples/go/.*\.go$ pass_filenames: false stages: [pre-push] diff --git a/core/common/src/error/iggy_error.rs b/core/common/src/error/iggy_error.rs index 6a31bcb48..754dd1d76 100644 --- a/core/common/src/error/iggy_error.rs +++ b/core/common/src/error/iggy_error.rs @@ -285,6 +285,11 @@ pub enum IggyError { TopicDirectoryNotFound(String) = 2020, #[error("Too many topics")] TooManyTopics = 2021, + /// Distinct from [`Self::TooManyPartitions`] by remedy: that one is cleared + /// by a smaller batch, this one by deleting partitions or raising the + /// node-wide cap. + #[error("Partitions limit reached, raise [metadata] partitions_max")] + PartitionsLimitReached = 2022, #[error("Cannot create partition with ID: {0} for stream with ID: {1} and topic with ID: {2}")] CannotCreatePartition(usize, usize, usize) = 3000, #[error( @@ -630,6 +635,14 @@ mod tests { ) } + #[test] + fn partitions_limit_reached_round_trips_by_code() { + let error = IggyError::PartitionsLimitReached; + assert_eq!(error.as_code(), 2022); + assert_eq!(IggyError::from_code(2022), error); + assert_eq!(IggyError::from_code_as_string(2022), error.as_string()); + } + #[test] fn too_many_consumer_offsets_round_trips_by_code() { let error = IggyError::TooManyConsumerOffsets; diff --git a/core/configs/src/common/defaults.rs b/core/configs/src/common/defaults.rs index 18e588c54..2ede3d6d2 100644 --- a/core/configs/src/common/defaults.rs +++ b/core/configs/src/common/defaults.rs @@ -15,6 +15,8 @@ // specific language governing permissions and limitations // under the License. +use iggy_common::IggyDuration; + use super::http::{HttpConfig, HttpCorsConfig, HttpJwtConfig, HttpMetricsConfig, HttpTlsConfig}; use super::server::{ ConsumerGroupConfig, HeartbeatConfig, MemoryPoolConfig, MessagesMaintenanceConfig, @@ -221,10 +223,19 @@ impl Default for LoggingConfig { .parse() .unwrap(), retention: SERVER_CONFIG.logging.retention.parse().unwrap(), + sysinfo_print_interval: default_sysinfo_print_interval(), } } } +pub(crate) fn default_sysinfo_print_interval() -> IggyDuration { + SERVER_CONFIG + .logging + .sysinfo_print_interval + .parse() + .unwrap() +} + impl Default for EncryptionConfig { fn default() -> EncryptionConfig { EncryptionConfig { diff --git a/core/configs/src/common/displays.rs b/core/configs/src/common/displays.rs index f2535cbd6..c06d0d1cf 100644 --- a/core/configs/src/common/displays.rs +++ b/core/configs/src/common/displays.rs @@ -134,14 +134,15 @@ impl Display for LoggingConfig { fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { write!( f, - "{{ path: {}, level: {}, file_enabled: {}, max_file_size: {}, max_total_size: {}, rotation_check_interval: {}, retention: {} }}", + "{{ path: {}, level: {}, file_enabled: {}, max_file_size: {}, max_total_size: {}, rotation_check_interval: {}, retention: {}, sysinfo_print_interval: {} }}", self.path, self.level, self.file_enabled, self.max_file_size.as_human_string_with_zero_as_unlimited(), self.max_total_size.as_human_string_with_zero_as_unlimited(), self.rotation_check_interval, - self.retention + self.retention, + self.sysinfo_print_interval ) } } diff --git a/core/configs/src/common/system.rs b/core/configs/src/common/system.rs index 0b42e493a..5fbba2a02 100644 --- a/core/configs/src/common/system.rs +++ b/core/configs/src/common/system.rs @@ -15,6 +15,7 @@ // specific language governing permissions and limitations // under the License. +use super::defaults::default_sysinfo_print_interval; use configs::ConfigEnv; use iggy_common::IggyByteSize; use iggy_common::IggyDuration; @@ -44,6 +45,12 @@ pub struct LoggingConfig { #[config_env(leaf)] #[serde_as(as = "DisplayFromStr")] pub retention: IggyDuration, + /// How often shard 0 logs one line of process and host usage. Zero + /// disables the line. + #[config_env(leaf)] + #[serde_as(as = "DisplayFromStr")] + #[serde(default = "default_sysinfo_print_interval")] + pub sysinfo_print_interval: IggyDuration, } impl From<&LoggingConfig> for LoggingSettings { diff --git a/core/configs/src/server_config/defaults.rs b/core/configs/src/server_config/defaults.rs index d96b274ac..7e8ecd0d2 100644 --- a/core/configs/src/server_config/defaults.rs +++ b/core/configs/src/server_config/defaults.rs @@ -174,6 +174,7 @@ impl Default for MetadataConfig { prepare_queue_depth: metadata.prepare_queue_depth as usize, journal_slots: metadata.journal_slots as usize, clients_table_max: metadata.clients_table_max as usize, + partitions_max: metadata.partitions_max as u32, } } } diff --git a/core/configs/src/server_config/displays.rs b/core/configs/src/server_config/displays.rs index 06e2b50ad..c3d2874fb 100644 --- a/core/configs/src/server_config/displays.rs +++ b/core/configs/src/server_config/displays.rs @@ -86,8 +86,11 @@ impl Display for MetadataConfig { fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { write!( f, - "{{ prepare_queue_depth: {}, journal_slots: {}, clients_table_max: {} }}", - self.prepare_queue_depth, self.journal_slots, self.clients_table_max, + "{{ prepare_queue_depth: {}, journal_slots: {}, clients_table_max: {}, partitions_max: {} }}", + self.prepare_queue_depth, + self.journal_slots, + self.clients_table_max, + self.partitions_max, ) } } diff --git a/core/configs/src/server_config/metadata.rs b/core/configs/src/server_config/metadata.rs index fbec9853f..f0a8b2abb 100644 --- a/core/configs/src/server_config/metadata.rs +++ b/core/configs/src/server_config/metadata.rs @@ -111,6 +111,15 @@ pub struct MetadataConfig { /// HTTP session cap tracks this at half, so raising it lifts /// both. pub clients_table_max: usize, + + /// Node-wide cap on partitions across all streams and topics. A + /// CreateTopic or CreatePartitions that would exceed it is rejected with + /// `PartitionsLimitReached` before it enters consensus. Zero is no cap. + /// + /// A soft cap: it counts committed partitions only, so creates in flight + /// at the same time can overshoot it. + #[serde(default)] + pub partitions_max: u32, } impl MetadataConfig { @@ -190,6 +199,7 @@ mod tests { prepare_queue_depth: DEFAULT_METADATA_PREPARE_QUEUE_DEPTH, journal_slots: DEFAULT_METADATA_JOURNAL_SLOTS, clients_table_max: DEFAULT_METADATA_CLIENTS_TABLE_MAX, + partitions_max: 0, }; assert!(config.validate().is_ok()); assert_eq!(config.checkpoint_margin(), METADATA_CHECKPOINT_MARGIN_FLOOR); @@ -201,6 +211,7 @@ mod tests { prepare_queue_depth: MAX_METADATA_PREPARE_QUEUE_DEPTH, journal_slots: 4096, clients_table_max: DEFAULT_METADATA_CLIENTS_TABLE_MAX, + partitions_max: 0, }; assert!(config.validate().is_ok()); assert_eq!(config.checkpoint_margin(), MAX_METADATA_PREPARE_QUEUE_DEPTH); @@ -215,6 +226,7 @@ mod tests { prepare_queue_depth: MAX_METADATA_PREPARE_QUEUE_DEPTH, journal_slots: min_slots, clients_table_max: DEFAULT_METADATA_CLIENTS_TABLE_MAX, + partitions_max: 0, }; assert!(boundary.validate().is_ok()); // ...one slot fewer is refused. @@ -222,6 +234,7 @@ mod tests { prepare_queue_depth: MAX_METADATA_PREPARE_QUEUE_DEPTH, journal_slots: min_slots - 1, clients_table_max: DEFAULT_METADATA_CLIENTS_TABLE_MAX, + partitions_max: 0, }; assert!(starved.validate().is_err()); } @@ -235,6 +248,7 @@ mod tests { prepare_queue_depth: MAX_METADATA_PREPARE_QUEUE_DEPTH + 1, journal_slots: MAX_METADATA_JOURNAL_SLOTS, clients_table_max: DEFAULT_METADATA_CLIENTS_TABLE_MAX, + partitions_max: 0, }; assert!(over.validate().is_err()); assert_eq!( @@ -250,6 +264,7 @@ mod tests { prepare_queue_depth: 0, journal_slots: DEFAULT_METADATA_JOURNAL_SLOTS, clients_table_max: DEFAULT_METADATA_CLIENTS_TABLE_MAX, + partitions_max: 0, }; assert!(config.validate().is_err()); } @@ -270,6 +285,7 @@ mod tests { prepare_queue_depth: DEFAULT_METADATA_PREPARE_QUEUE_DEPTH, journal_slots: DEFAULT_METADATA_JOURNAL_SLOTS, clients_table_max: MIN_METADATA_CLIENTS_TABLE_MAX - 1, + partitions_max: 0, }; assert!(config.validate().is_err()); } @@ -280,6 +296,7 @@ mod tests { prepare_queue_depth: DEFAULT_METADATA_PREPARE_QUEUE_DEPTH, journal_slots: DEFAULT_METADATA_JOURNAL_SLOTS, clients_table_max: MIN_METADATA_CLIENTS_TABLE_MAX, + partitions_max: 0, }; assert!(config.validate().is_ok()); } @@ -290,6 +307,7 @@ mod tests { prepare_queue_depth: DEFAULT_METADATA_PREPARE_QUEUE_DEPTH, journal_slots: DEFAULT_METADATA_JOURNAL_SLOTS, clients_table_max: MAX_METADATA_CLIENTS_TABLE_MAX + 1, + partitions_max: 0, }; assert!(config.validate().is_err()); } diff --git a/core/consensus/src/fatal.rs b/core/consensus/src/fatal.rs deleted file mode 100644 index 341bda389..000000000 --- a/core/consensus/src/fatal.rs +++ /dev/null @@ -1,60 +0,0 @@ -// 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. - -/// 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, -} - -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. -pub fn fatal(reason: FatalReason, message: &str) -> ! { - tracing::error!( - target: "iggy.consensus.diag", - reason = ?reason, - exit_status = reason.exit_status(), - "{message}" - ); - std::process::exit(i32::from(reason.exit_status())); -} diff --git a/core/consensus/src/impls.rs b/core/consensus/src/impls.rs index e31680eec..e462c4651 100644 --- a/core/consensus/src/impls.rs +++ b/core/consensus/src/impls.rs @@ -1383,12 +1383,23 @@ impl<B: MessageBus, P: Pipeline<Entry = PipelineEntry>> VsrConsensus<B, P> { // The one line proving the durable record was READ BACK, not merely // written: a replica that came back at view 0 is otherwise // indistinguishable from one that resumed correctly until it votes. - tracing::info!( - group, - view, - log_view, - "restored group view from its superblock" - ); + // Only the metadata group says so at INFO. Every partition group + // restores through here too, and one line each floods the boot log. + if group == METADATA_GROUP { + tracing::info!( + group, + view, + log_view, + "restored group view from its superblock" + ); + } else { + tracing::debug!( + group, + view, + log_view, + "restored group view from its superblock" + ); + } consensus.set_view(view); consensus.set_log_view(log_view); consensus.mark_superblock_durable(view, log_view); diff --git a/core/consensus/src/lib.rs b/core/consensus/src/lib.rs index 0abad9657..b5ff7d510 100644 --- a/core/consensus/src/lib.rs +++ b/core/consensus/src/lib.rs @@ -188,8 +188,7 @@ pub use state_transfer::{ pub(crate) mod oneshot; pub use oneshot::{Canceled, Receiver, Sender, channel as oneshot_channel}; -mod fatal; -pub use fatal::{FatalReason, fatal}; +pub use server_common::fatal::{FatalReason, fatal}; mod impls; pub use impls::*; diff --git a/core/integration/tests/server/mod.rs b/core/integration/tests/server/mod.rs index 159b9217b..84d4a51f2 100644 --- a/core/integration/tests/server/mod.rs +++ b/core/integration/tests/server/mod.rs @@ -80,6 +80,7 @@ mod partition_view_durability_vsr; mod concurrent_addition; mod consumer_offset_quota_vsr; mod general; +mod partitions_limit_vsr; // The per-shard segment cleaner deletes expired / oversize segments from disk. mod message_cleanup; mod message_retrieval; diff --git a/core/integration/tests/server/partition_view_durability_vsr.rs b/core/integration/tests/server/partition_view_durability_vsr.rs index 61d25c9fe..4d64b58d5 100644 --- a/core/integration/tests/server/partition_view_durability_vsr.rs +++ b/core/integration/tests/server/partition_view_durability_vsr.rs @@ -55,11 +55,15 @@ const CONVERGE_TIMEOUT: Duration = Duration::from_secs(60); /// Boot line reporting the `(view, log_view)` a consensus group restored from /// its superblock: the recovered replica's own account of what it read back. /// One line per group since the unified restore constructor, so the parser -/// filters the metadata group's line out by its `group` field. +/// filters the metadata group's line out by its `group` field. Partition +/// groups log it at `debug`, hence the `logging.level` override below. const RESTORED_VIEW_MARKER: &str = "restored group view from its superblock"; const POLL_INTERVAL: Duration = Duration::from_millis(250); -#[iggy_harness(cluster_nodes = 3, server(sharding.cpu_allocation = "0..1"))] +#[iggy_harness( + cluster_nodes = 3, + server(sharding.cpu_allocation = "0..1", logging.level = "info,consensus=debug") +)] async fn given_advanced_partition_view_when_survivor_restarts_should_recover_view_from_superblock( harness: &mut TestHarness, ) { diff --git a/core/integration/tests/server/partitions_limit_vsr.rs b/core/integration/tests/server/partitions_limit_vsr.rs new file mode 100644 index 000000000..414d84842 --- /dev/null +++ b/core/integration/tests/server/partitions_limit_vsr.rs @@ -0,0 +1,121 @@ +// 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. + +//! The node-wide `[metadata] partitions_max` cap. A create-topic or +//! create-partitions that would push the committed partition count past it +//! denies typed with `PartitionsLimitReached` over TCP and HTTP. Creates up to +//! the cap succeed, and deleting partitions frees room again. + +use iggy::prelude::*; +use iggy_common::create_topic::CreateTopic; +use integration::iggy_harness; +use reqwest::StatusCode; +use serde_json::json; + +use super::http_client::HttpClient; + +const STREAM_NAME: &str = "partitions-limit-stream"; + +async fn create_topic( + client: &IggyClient, + stream_id: &Identifier, + name: &str, + partitions_count: u32, +) -> Result<TopicDetails, IggyError> { + client + .create_topic( + stream_id, + name, + &TopicCreateOptions { + partitions_count: Some(partitions_count), + message_expiry: Some(IggyExpiry::NeverExpire), + ..TopicCreateOptions::default() + }, + ) + .await +} + +fn assert_limit_reached<T: std::fmt::Debug>(result: &Result<T, IggyError>, context: &str) { + let limit_reached = IggyError::PartitionsLimitReached.as_code(); + assert!( + matches!(result, Err(error) if error.as_code() == limit_reached), + "{context} must deny with PartitionsLimitReached, got {result:?}" + ); +} + +#[iggy_harness( + cluster_nodes = 1, + server(metadata.partitions_max = "4") +)] +async fn given_partitions_cap_when_creating_past_it_should_reject_typed(harness: &TestHarness) { + let client = harness.tcp_root_client().await.expect("TCP root client"); + client + .create_stream(STREAM_NAME) + .await + .expect("create stream"); + let stream_id = Identifier::named(STREAM_NAME).expect("stream identifier"); + let topic_id = Identifier::named("first").expect("topic identifier"); + + create_topic(&client, &stream_id, "first", 3) + .await + .expect("3 of 4 partitions fit under the cap"); + assert_limit_reached( + &create_topic(&client, &stream_id, "second", 2).await, + "a topic of 2 partitions on top of 3", + ); + assert_limit_reached( + &client.create_partitions(&stream_id, &topic_id, 2).await, + "adding 2 partitions on top of 3", + ); + client + .create_partitions(&stream_id, &topic_id, 1) + .await + .expect("the cap itself is admissible"); + assert_limit_reached( + &create_topic(&client, &stream_id, "second", 1).await, + "a topic of 1 partition at the cap", + ); + + let http = HttpClient::login_root(harness).await; + let topic_body = serde_json::to_value(CreateTopic { + name: "second".to_owned(), + partitions_count: 1, + ..CreateTopic::default() + }) + .expect("serialize create topic"); + for (path, body) in [ + (format!("/streams/{STREAM_NAME}/topics"), topic_body), + ( + format!("/streams/{STREAM_NAME}/topics/first/partitions"), + json!({ "partitions_count": 1 }), + ), + ] { + let response = http.post_json(&path, &body).await; + assert_eq!(response.status(), StatusCode::BAD_REQUEST, "POST {path}"); + let body: serde_json::Value = response.json().await.expect("HTTP error body"); + assert_eq!(body["id"], 2022, "POST {path}"); + assert_eq!(body["code"], "partitions_limit_reached", "POST {path}"); + } + + client + .delete_partitions(&stream_id, &topic_id, 2) + .await + .expect("delete 2 partitions"); + create_topic(&client, &stream_id, "second", 2) + .await + .expect("deleted partitions free room under the cap"); +} diff --git a/core/journal/src/durable_storage.rs b/core/journal/src/durable_storage.rs index fc6f148ed..2f4f32832 100644 --- a/core/journal/src/durable_storage.rs +++ b/core/journal/src/durable_storage.rs @@ -23,6 +23,7 @@ use compio::io::{AsyncReadAtExt, AsyncWriteAtExt}; use futures::channel::oneshot; use futures::lock::Mutex; use futures::{Stream, stream}; +use server_common::fatal::ExitOnDescriptorExhaustion; use server_common::iobuf::{Frozen, Owned}; use std::ffi::OsString; use std::io; @@ -249,7 +250,10 @@ impl DurableStorage for DiskStorage { .create(true) .truncate(matches!(mode, OpenMode::Create | OpenMode::CreateWriteOnly)); } - options.open(path).await + options + .open(path) + .await + .exit_on_descriptor_exhaustion(|| format!("opening {}", path.display())) } async fn create_directories(&self, path: &Path) -> io::Result<()> { @@ -257,7 +261,11 @@ impl DurableStorage for DiskStorage { } async fn sync_directory(&self, path: &Path) -> io::Result<()> { - File::open(path).await?.sync_all().await + File::open(path) + .await + .exit_on_descriptor_exhaustion(|| format!("opening directory {}", path.display()))? + .sync_all() + .await } async fn rename(&self, source: &Path, target: &Path) -> io::Result<()> { @@ -291,7 +299,11 @@ impl DurableStorage for DiskStorage { async fn entries(&self, path: &Path) -> io::Result<Vec<StorageEntry>> { // getdents has no io_uring operation, and shard fallback pools are disabled. let path = path.to_path_buf(); - run_blocking("iggy-directory-scan", move || directory_entries(&path)).await + run_blocking("iggy-directory-scan", move || { + directory_entries(&path) + .exit_on_descriptor_exhaustion(|| format!("listing directory {}", path.display())) + }) + .await } async fn regular_files(&self, path: &Path) -> io::Result<RegularFiles> { @@ -311,7 +323,10 @@ impl DurableStorage for DiskStorage { let result = (|| { // An unreadable entry must not hide the remaining files. // Only opening the directory fails the scan. - for entry in std::fs::read_dir(&directory)? { + let entries = std::fs::read_dir(&directory).exit_on_descriptor_exhaustion( + || format!("listing directory {}", directory.display()), + )?; + for entry in entries { let entry = match entry { Ok(entry) => entry, Err(error) => { @@ -443,7 +458,11 @@ impl DurableFile for File { async fn truncate(&self, length: u64) -> io::Result<()> { // Older kernels lack IORING_OP_FTRUNCATE and shard fallback pools are // disabled. Own the inode until the worker completes, even on cancellation. - let descriptor = std::os::fd::AsFd::as_fd(self).try_clone_to_owned()?; + let descriptor = std::os::fd::AsFd::as_fd(self) + .try_clone_to_owned() + .exit_on_descriptor_exhaustion(|| { + "duplicating a file descriptor to truncate".to_owned() + })?; run_blocking("iggy-file-truncate", move || { std::fs::File::from(descriptor).set_len(length) }) diff --git a/core/journal/src/file_storage.rs b/core/journal/src/file_storage.rs index b8f07fb08..1af706c50 100644 --- a/core/journal/src/file_storage.rs +++ b/core/journal/src/file_storage.rs @@ -18,6 +18,7 @@ use crate::Storage; use compio::buf::IoBuf; use compio::io::{AsyncReadAtExt, AsyncWriteAtExt}; +use server_common::fatal::ExitOnDescriptorExhaustion; use std::cell::{Cell, UnsafeCell}; use std::fs; use std::io; @@ -45,7 +46,8 @@ impl FileStorage { .create(true) .truncate(false) .open(path) - .await?; + .await + .exit_on_descriptor_exhaustion(|| format!("opening {}", path.display()))?; let len = file.metadata().await?.len(); Ok(Self { file: UnsafeCell::new(file), @@ -79,7 +81,13 @@ impl FileStorage { pub(crate) fn truncate(&self, len: u64) -> io::Result<()> { // SAFETY: single-threaded compio runtime, no concurrent access to the file. let file = unsafe { &*self.file.get() }; - let file = fs::File::from(file.as_fd().try_clone_to_owned()?); + let file = fs::File::from( + file.as_fd() + .try_clone_to_owned() + .exit_on_descriptor_exhaustion(|| { + format!("duplicating the descriptor of {}", self.path.display()) + })?, + ); file.set_len(len)?; self.write_offset.set(len); file.sync_all() @@ -162,7 +170,8 @@ impl FileStorage { .read(true) .write(true) .open(&self.path) - .await?; + .await + .exit_on_descriptor_exhaustion(|| format!("reopening {}", self.path.display()))?; let len = file.metadata().await?.len(); // SAFETY: single-threaded compio runtime, no concurrent access to the file. unsafe { *self.file.get() = file }; diff --git a/core/journal/src/prepare_journal.rs b/core/journal/src/prepare_journal.rs index d7de23875..a6c108f49 100644 --- a/core/journal/src/prepare_journal.rs +++ b/core/journal/src/prepare_journal.rs @@ -19,6 +19,7 @@ use crate::file_storage::FileStorage; use crate::{Journal, JournalHandle}; use compio::io::AsyncWriteAtExt; use iggy_binary_protocol::consensus::{CHECKSUM_UNSEALED, Command, PrepareHeader}; +use server_common::fatal::ExitOnDescriptorExhaustion; use server_common::{MESSAGE_ALIGN, Message, iobuf::Owned}; use std::cell::{Cell, OnceCell, Ref, RefCell}; use std::fmt; @@ -819,7 +820,9 @@ impl Journal for PrepareJournal { let tmp_path = wal_path.with_extension("wal.tmp"); let tmp_guard = TmpFileGuard::new(tmp_path.clone()); { - let mut tmp = compio::fs::File::create(&tmp_path).await?; + let mut tmp = compio::fs::File::create(&tmp_path) + .await + .exit_on_descriptor_exhaustion(|| format!("creating {}", tmp_path.display()))?; let mut write_pos: u64 = 0; for (header, old_offset) in &live { let size = header.size as usize; @@ -839,7 +842,12 @@ impl Journal for PrepareJournal { tmp_guard.defuse(); if let Some(parent) = wal_path.parent() { - let dir = match compio::fs::File::open(parent).await { + let opened = compio::fs::File::open(parent) + .await + .exit_on_descriptor_exhaustion(|| { + format!("opening directory {}", parent.display()) + }); + let dir = match opened { Ok(dir) => dir, Err(error) => { return Err(self.poison("truncate_from: open parent dir for fsync", error)); @@ -981,7 +989,9 @@ impl Journal for PrepareJournal { let tmp_path = wal_path.with_extension("wal.tmp"); let tmp_guard = TmpFileGuard::new(tmp_path.clone()); { - let mut tmp = compio::fs::File::create(&tmp_path).await?; + let mut tmp = compio::fs::File::create(&tmp_path) + .await + .exit_on_descriptor_exhaustion(|| format!("creating {}", tmp_path.display()))?; let mut write_pos: u64 = 0; for (header, old_offset) in &live { let size = header.size as usize; @@ -1013,7 +1023,12 @@ impl Journal for PrepareJournal { // poisoned the caller learns the drain is not durable instead // of silently proceeding. if let Some(parent) = wal_path.parent() { - let dir = match compio::fs::File::open(parent).await { + let opened = compio::fs::File::open(parent) + .await + .exit_on_descriptor_exhaustion(|| { + format!("opening directory {}", parent.display()) + }); + let dir = match opened { Ok(d) => d, Err(e) => { return Err(self.poison("drain: open parent dir for fsync", e)); diff --git a/core/journal/src/superblock.rs b/core/journal/src/superblock.rs index 77f49ffe0..c94d3f28d 100644 --- a/core/journal/src/superblock.rs +++ b/core/journal/src/superblock.rs @@ -50,6 +50,7 @@ use std::io; use std::path::{Path, PathBuf}; use compio::io::{AsyncReadAtExt, AsyncWriteAtExt}; +use server_common::fatal::ExitOnDescriptorExhaustion; use crate::prepare_journal::TmpFileGuard; use twox_hash::XxHash3_64; @@ -491,7 +492,10 @@ const fn has_unreadable_sequence(slot: &SlotClass) -> bool { } async fn read_slot(path: &Path) -> io::Result<SlotClass> { - let file = match compio::fs::File::open(path).await { + let opened = compio::fs::File::open(path) + .await + .exit_on_descriptor_exhaustion(|| format!("opening superblock slot {}", path.display())); + let file = match opened { Ok(file) => file, Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok(SlotClass::Absent), Err(e) => return Err(e), @@ -528,7 +532,9 @@ async fn atomic_replace(dir: &Path, file_name: &str, bytes: Vec<u8>) -> io::Resu // un-advanced so a retry re-targets the same slot; it keeps a failing disk from // littering `superblock.{a,b}.tmp` next to the slots an operator is inspecting. let guard = TmpFileGuard::new(tmp_path.clone()); - let mut tmp = compio::fs::File::create(&tmp_path).await?; + let mut tmp = compio::fs::File::create(&tmp_path) + .await + .exit_on_descriptor_exhaustion(|| format!("creating {}", tmp_path.display()))?; let (result, _buf) = tmp.write_all_at(bytes, 0).await.into(); result?; tmp.sync_all().await?; @@ -536,7 +542,9 @@ async fn atomic_replace(dir: &Path, file_name: &str, bytes: Vec<u8>) -> io::Resu compio::fs::rename(&tmp_path, &final_path).await?; guard.defuse(); - let dir_file = compio::fs::File::open(dir).await?; + let dir_file = compio::fs::File::open(dir) + .await + .exit_on_descriptor_exhaustion(|| format!("opening directory {}", dir.display()))?; dir_file.sync_all().await?; Ok(()) } diff --git a/core/message_bus/src/accept.rs b/core/message_bus/src/accept.rs new file mode 100644 index 000000000..db3eb857e --- /dev/null +++ b/core/message_bus/src/accept.rs @@ -0,0 +1,42 @@ +// 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 does not stop the process on this, unlike the storage +/// paths: without a connection cap, any client that can open sockets could +/// then stop the node. +#[allow(clippy::future_not_send)] +pub async fn pause_after_accept_error(error: &io::Error) { + if is_descriptor_exhaustion(error) { + compio::time::sleep(DESCRIPTOR_EXHAUSTION_BACKOFF).await; + } +} diff --git a/core/message_bus/src/client_listener/tcp.rs b/core/message_bus/src/client_listener/tcp.rs index a14673774..d8bf4affc 100644 --- a/core/message_bus/src/client_listener/tcp.rs +++ b/core/message_bus/src/client_listener/tcp.rs @@ -27,6 +27,7 @@ //! writer + reader tasks via [`crate::installer`]. use crate::AcceptedClientFn; +use crate::accept::pause_after_accept_error; use crate::client_listener::bind_nodelay_listener; use crate::lifecycle::ShutdownToken; use compio::net::TcpListener; @@ -77,6 +78,7 @@ pub async fn run(listener: TcpListener, token: ShutdownToken, on_accepted: Accep } Err(e) => { error!("Client listener (TCP) accept failed: {e}"); + pause_after_accept_error(&e).await; } } } diff --git a/core/message_bus/src/client_listener/tcp_tls.rs b/core/message_bus/src/client_listener/tcp_tls.rs index cb3007633..4fe022f09 100644 --- a/core/message_bus/src/client_listener/tcp_tls.rs +++ b/core/message_bus/src/client_listener/tcp_tls.rs @@ -34,6 +34,7 @@ //! TCP does not carry over. use crate::AcceptedTlsClientFn; +use crate::accept::pause_after_accept_error; use crate::client_listener::bind_nodelay_listener; use crate::lifecycle::ShutdownToken; use crate::transports::tls::{TlsServerCredentials, install_default_crypto_provider}; @@ -122,6 +123,7 @@ pub async fn run( } Err(e) => { error!("Client listener (TCP-TLS) accept failed: {e}"); + pause_after_accept_error(&e).await; } } } diff --git a/core/message_bus/src/client_listener/ws.rs b/core/message_bus/src/client_listener/ws.rs index f61f8c7e4..8fe618ff6 100644 --- a/core/message_bus/src/client_listener/ws.rs +++ b/core/message_bus/src/client_listener/ws.rs @@ -35,6 +35,7 @@ //! without pulling `shard` in as a doc-only dep. use crate::AcceptedWsClientFn; +use crate::accept::pause_after_accept_error; use crate::client_listener::bind_nodelay_listener; use crate::lifecycle::ShutdownToken; use compio::net::TcpListener; @@ -91,6 +92,7 @@ pub async fn run(listener: TcpListener, token: ShutdownToken, on_accepted: Accep } Err(e) => { error!("Client listener (WS) accept failed: {e}"); + pause_after_accept_error(&e).await; } } } diff --git a/core/message_bus/src/client_listener/wss.rs b/core/message_bus/src/client_listener/wss.rs index 8308a4460..ac0933fe2 100644 --- a/core/message_bus/src/client_listener/wss.rs +++ b/core/message_bus/src/client_listener/wss.rs @@ -38,6 +38,7 @@ //! carry over. //! use crate::AcceptedWssClientFn; +use crate::accept::pause_after_accept_error; use crate::lifecycle::ShutdownToken; use crate::socket_opts::bind_reusable_tcp_listener; use crate::transports::tls::{TlsServerCredentials, install_default_crypto_provider}; @@ -116,6 +117,7 @@ pub async fn run( } Err(e) => { error!("Client listener (WSS) accept failed: {e}"); + pause_after_accept_error(&e).await; } } } diff --git a/core/message_bus/src/lib.rs b/core/message_bus/src/lib.rs index 8bd2a3fe2..61b649c69 100644 --- a/core/message_bus/src/lib.rs +++ b/core/message_bus/src/lib.rs @@ -75,6 +75,7 @@ //! transport (TCP, TCP-TLS, WS, WSS, QUIC) plugs in behind the same //! registry, fencing, and dispatch logic. +pub mod accept; pub mod cache; pub mod client_listener; pub mod config; diff --git a/core/message_bus/src/replica/listener.rs b/core/message_bus/src/replica/listener.rs index 71988e989..6d2843615 100644 --- a/core/message_bus/src/replica/listener.rs +++ b/core/message_bus/src/replica/listener.rs @@ -50,6 +50,7 @@ //! directional check itself runs on the owning shard, inside the //! handshake. +use crate::accept::pause_after_accept_error; use crate::lifecycle::ShutdownToken; use crate::socket_opts::bind_reusable_tcp_listener; use crate::{AcceptedReplicaFn, GenericHeader, Message}; @@ -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; } } } diff --git a/core/metadata/src/impls/metadata.rs b/core/metadata/src/impls/metadata.rs index 953f0e850..4bf11badb 100644 --- a/core/metadata/src/impls/metadata.rs +++ b/core/metadata/src/impls/metadata.rs @@ -62,6 +62,7 @@ use journal::superblock::{ use journal::{Journal, JournalHandle}; use message_bus::MessageBus; use server_common::Message; +use server_common::fatal::ExitOnDescriptorExhaustion; use server_common::iobuf::{Frozen, Owned}; use std::cell::{Cell, RefCell}; use std::mem::size_of; @@ -145,10 +146,12 @@ impl IggySnapshot { let tmp_path = path.with_extension("bin.tmp"); - let mut file = fs::File::create(&tmp_path).map_err(|e| SnapshotError::Persist { - stage: PersistStage::Write, - source: e, - })?; + let mut file = fs::File::create(&tmp_path) + .exit_on_descriptor_exhaustion(|| format!("creating {}", tmp_path.display())) + .map_err(|e| SnapshotError::Persist { + stage: PersistStage::Write, + source: e, + })?; file.write_all(encoded) .map_err(|e| SnapshotError::Persist { stage: PersistStage::Write, @@ -177,10 +180,12 @@ impl IggySnapshot { // Fsync the parent directory to ensure the rename is durable. if let Some(parent) = path.parent() { - let dir = fs::File::open(parent).map_err(|e| SnapshotError::Persist { - stage: PersistStage::DirSync, - source: e, - })?; + let dir = fs::File::open(parent) + .exit_on_descriptor_exhaustion(|| format!("opening directory {}", parent.display())) + .map_err(|e| SnapshotError::Persist { + stage: PersistStage::DirSync, + source: e, + })?; dir.sync_all().map_err(|e| SnapshotError::Persist { stage: PersistStage::DirSync, source: e, @@ -207,7 +212,8 @@ impl IggySnapshot { /// written in a format version this build does not read, or `SnapshotError` if the /// file cannot be read or deserialized. pub fn load(path: &Path) -> Result<(Self, u128), SnapshotError> { - let data = std::fs::read(path)?; + let data = std::fs::read(path) + .exit_on_descriptor_exhaustion(|| format!("reading {}", path.display()))?; let (payload, checksum) = split_trailer(&data, path)?; Ok((Self::decode(payload)?, checksum)) } @@ -1651,6 +1657,7 @@ where return Err(StateTransferUnavailable::NoSnapshot); } let sealed = std::fs::read(&path) + .exit_on_descriptor_exhaustion(|| format!("reading {}", path.display())) .map_err(|source| StateTransferUnavailable::SnapshotUnreadable(source.into()))?; // Verifies the trailer and hands back the payload alone. let (payload, _) = diff --git a/core/metadata/src/stm/stream.rs b/core/metadata/src/stm/stream.rs index afba796cb..caf735c4c 100644 --- a/core/metadata/src/stm/stream.rs +++ b/core/metadata/src/stm/stream.rs @@ -1660,6 +1660,20 @@ impl Streams { }) } + /// Total committed partition count across all topics (for the node-wide + /// `[metadata] partitions_max` admission check). + #[must_use] + pub fn partition_count(&self) -> usize { + self.inner.read(|inner| { + inner + .items + .iter() + .flat_map(|(_, stream)| stream.topics.iter()) + .map(|(_, topic)| topic.partitions.len()) + .sum() + }) + } + #[must_use] pub fn partition_count_context( &self, diff --git a/core/partitions/src/iggy_index_reader.rs b/core/partitions/src/iggy_index_reader.rs index bccdbe27f..a7719e8d5 100644 --- a/core/partitions/src/iggy_index_reader.rs +++ b/core/partitions/src/iggy_index_reader.rs @@ -20,6 +20,7 @@ use bytes::Buf; use compio::fs::{File, OpenOptions}; use compio::io::AsyncReadAtExt; use iggy_common::IggyError; +use server_common::fatal::ExitOnDescriptorExhaustion; use tracing::trace; /// Reader for the sparse index file written by [`crate::IggyIndexWriter`]. @@ -45,6 +46,7 @@ impl IggyIndexReader { .read(true) .open(file_path) .await + .exit_on_descriptor_exhaustion(|| format!("opening {file_path}")) .map_err(|_| IggyError::CannotReadFile)?; Ok(Self { file_path: file_path.to_owned(), diff --git a/core/partitions/src/iggy_index_writer.rs b/core/partitions/src/iggy_index_writer.rs index c91cea223..c30e1d764 100644 --- a/core/partitions/src/iggy_index_writer.rs +++ b/core/partitions/src/iggy_index_writer.rs @@ -18,6 +18,7 @@ use compio::fs::{File, OpenOptions}; use compio::io::AsyncWriteAtExt; use iggy_common::IggyError; +use server_common::fatal::ExitOnDescriptorExhaustion; use std::rc::Rc; use std::sync::atomic::{AtomicU64, Ordering}; use tracing::{error, trace}; @@ -51,6 +52,7 @@ impl IggyIndexWriter { let file = opts .open(file_path) .await + .exit_on_descriptor_exhaustion(|| format!("opening {file_path}")) .map_err(|_| IggyError::CannotReadFile)?; if file_exists { diff --git a/core/partitions/src/iggy_partition.rs b/core/partitions/src/iggy_partition.rs index c5e057881..6d5033248 100644 --- a/core/partitions/src/iggy_partition.rs +++ b/core/partitions/src/iggy_partition.rs @@ -1646,7 +1646,7 @@ where if !committed_restored && !append_restored { return; } - tracing::info!( + tracing::debug!( namespace_raw = self.consensus().group(), offset_frontier = frontier, offset_reserved = reserved, diff --git a/core/partitions/src/messages_writer.rs b/core/partitions/src/messages_writer.rs index 89dcb3cfb..c5d31e674 100644 --- a/core/partitions/src/messages_writer.rs +++ b/core/partitions/src/messages_writer.rs @@ -20,6 +20,7 @@ use compio::{ io::AsyncWriteAtExt, }; use iggy_common::{IggyByteSize, IggyError}; +use server_common::fatal::ExitOnDescriptorExhaustion; use server_common::fs_utils::preallocate_file; use server_common::iobuf::{Frozen, IOV_MAX}; use std::{ @@ -59,6 +60,7 @@ impl MessagesWriter { let file = opts .open(file_path) .await + .exit_on_descriptor_exhaustion(|| format!("opening {file_path}")) .map_err(|_| IggyError::CannotReadFile)?; if let Some(preallocate_size) = preallocate_size { diff --git a/core/partitions/src/offset_storage.rs b/core/partitions/src/offset_storage.rs index 266f5a280..4ad33658a 100644 --- a/core/partitions/src/offset_storage.rs +++ b/core/partitions/src/offset_storage.rs @@ -26,6 +26,7 @@ //! directory sync makes creation, replacement, or deletion durable. Offset callers //! own that directory sync, while purge marker writes include it before returning. +use server_common::fatal::ExitOnDescriptorExhaustion; use std::{io, path::Path}; use compio::{ @@ -184,6 +185,7 @@ pub async fn persist_offset_retained( .truncate(true) .open(path) .await + .exit_on_descriptor_exhaustion(|| format!("opening {path}")) .map_err(|_| IggyError::CannotOpenConsumerOffsetsFile(path.to_owned()))? }; let result = file.write_all_at(encode_offset_record(offset), 0).await.0; @@ -499,7 +501,12 @@ pub async fn read_purge_generation<S: DurableStorage>( async fn read_offset_record(path: &str) -> Result<Option<OffsetRecord>, IggyError> { // Absence answered by the open, not a `Path::exists()` probe: that is a BLOCKING // stat on the pump before every cold-key commit (see `persist_offset`). - let file = match OpenOptions::new().read(true).open(path).await { + let opened = OpenOptions::new() + .read(true) + .open(path) + .await + .exit_on_descriptor_exhaustion(|| format!("opening {path}")); + let file = match opened { Ok(file) => file, Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None), Err(_) => return Err(IggyError::CannotOpenConsumerOffsetsFile(path.to_owned())), diff --git a/core/partitions/src/poll_plan.rs b/core/partitions/src/poll_plan.rs index 8c3ca7131..68865d2db 100644 --- a/core/partitions/src/poll_plan.rs +++ b/core/partitions/src/poll_plan.rs @@ -30,6 +30,7 @@ use crate::{PollFragments, PollingConsumer}; use compio::io::AsyncReadAtExt; use iggy_binary_protocol::{WireError, batch}; use iggy_common::{ConsumerKind, IggyError}; +use server_common::fatal::exit_if_descriptors_exhausted; use server_common::iobuf::{Frozen, Owned}; use server_common::poll::PollHistoryId; use server_common::send_messages::{BatchIntegrity, COMMAND_HEADER_SIZE}; @@ -848,14 +849,18 @@ impl DiskReadPlan { ); } - /// Open a segment file for a disk poll, retrying transient IO failures (fd - /// pressure under heavy parallel load) so one failed syscall does not - /// silently collapse the poll into an empty result. + /// Open a segment file for a disk poll, retrying transient IO failures so + /// one failed syscall does not silently collapse the poll into an empty + /// result. Descriptor exhaustion is not retried: it stops the process. async fn open_segment_with_retry(&self, path: &str) -> Option<compio::fs::File> { for attempt in 0..3u8 { match compio::fs::File::open(path).await { Ok(file) => return Some(file), Err(error) => { + exit_if_descriptors_exhausted( + &error, + &format!("opening {path} for a disk poll"), + ); warn!( target: "iggy.partitions.diag", plane = "partitions", diff --git a/core/partitions/src/segment_anchor.rs b/core/partitions/src/segment_anchor.rs index 1891ce5fe..07f45b101 100644 --- a/core/partitions/src/segment_anchor.rs +++ b/core/partitions/src/segment_anchor.rs @@ -31,6 +31,7 @@ use crate::state_transfer::STAGING_SUFFIX; use compio::io::AsyncWriteAtExt; use consensus::state_artifact_checksum; +use server_common::fatal::ExitOnDescriptorExhaustion; use std::io; /// File extension for an anchor record, `{start_offset:020}.anchor` beside the @@ -155,7 +156,9 @@ pub async fn write_anchor(partition_dir: &str, anchor: SegmentAnchor) -> io::Res // unlinks it unconditionally, boot included, so a torn write leaves nothing // a later guard can read. let tmp_path = format!("{path}{STAGING_SUFFIX}"); - let mut file = compio::fs::File::create(&tmp_path).await?; + let mut file = compio::fs::File::create(&tmp_path) + .await + .exit_on_descriptor_exhaustion(|| format!("creating {tmp_path}"))?; let (result, _buf) = file .write_all_at(anchor.to_bytes().to_vec(), 0) .await @@ -194,7 +197,10 @@ pub async fn read_anchor( if length != ANCHOR_ENCODED_LEN as u64 { return Ok(None); } - match compio::fs::read(&path).await { + let read = compio::fs::read(&path) + .await + .exit_on_descriptor_exhaustion(|| format!("reading {path}")); + match read { Ok(bytes) => Ok(SegmentAnchor::from_bytes(&bytes)), Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(None), Err(error) => Err(error), diff --git a/core/partitions/src/segment_recovery.rs b/core/partitions/src/segment_recovery.rs index 14ff17313..f70d8b2ec 100644 --- a/core/partitions/src/segment_recovery.rs +++ b/core/partitions/src/segment_recovery.rs @@ -33,6 +33,7 @@ use crate::segment_anchor::ANCHOR_EXTENSION; use crate::state_transfer::STAGING_SUFFIX; use crate::{IggyIndex, IggyIndexReader, PartitionsConfig, Segment}; use iggy_common::{IggyByteSize, IggyError, MAX_MESSAGE_SIZE_UPPER_BYTES, PartitionStats}; +use server_common::fatal::ExitOnDescriptorExhaustion; use server_common::send_messages::{BatchHeader, COMMAND_HEADER_SIZE, decode_batch_slice}; use server_common::sharding::IggyNamespace; use server_common::{SegmentStorage, yield_to_reactor}; @@ -1111,7 +1112,9 @@ async fn ensure_contiguous_chain( fn sweep_scratch_files_and_collect_offsets( partition_path: &str, ) -> Result<Vec<u64>, PartitionRecoveryError> { - let entries = match fs::read_dir(partition_path) { + let entries = match fs::read_dir(partition_path) + .exit_on_descriptor_exhaustion(|| format!("listing directory {partition_path}")) + { Ok(entries) => entries, Err(source) if source.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()), Err(source) => { @@ -1244,6 +1247,7 @@ fn truncate_to(path: &str, target_size: u64) -> Result<(), PartitionRecoveryErro let file = fs::OpenOptions::new() .write(true) .open(path) + .exit_on_descriptor_exhaustion(|| format!("opening {path}")) .map_err(|source| { error!( path, @@ -1291,6 +1295,7 @@ fn stage_rebuilt_index(index_path: &str, entries: &[u8]) -> Result<String, Parti .create(true) .truncate(true) .open(&staging_path) + .exit_on_descriptor_exhaustion(|| format!("opening {staging_path}")) .map_err(|source| { error!( path = %staging_path, @@ -1343,6 +1348,7 @@ fn install_rebuilt_index( /// other mutation in this module (see [`FileScanner`]). fn fsync_dir(dir: &str) -> Result<(), PartitionRecoveryError> { fs::File::open(dir) + .exit_on_descriptor_exhaustion(|| format!("opening directory {dir}")) .and_then(|handle| handle.sync_all()) .map_err(|source| { error!( @@ -1471,6 +1477,7 @@ fn rename_into_fence(source_path: &str, target: &Path) -> Result<(), PartitionRe fn seed_empty_file(path: &str) -> Result<(), PartitionRecoveryError> { fs::File::create(path) + .exit_on_descriptor_exhaustion(|| format!("creating {path}")) .and_then(|file| file.sync_all()) .map_err(|source| { error!( @@ -2242,17 +2249,19 @@ async fn find_provable_index_anchor( searched_entries: 0, }); }; - let file = fs::File::open(index_path).map_err(|source| { - error!( - stream_id = identity.stream_id, - topic_id = identity.topic_id, - partition_id = identity.partition_id, - path = %index_path, - error = %source, - "failed to open sparse index for anchor search during recovery" - ); - PartitionRecoveryError::from(IggyError::CannotReadFile) - })?; + let file = fs::File::open(index_path) + .exit_on_descriptor_exhaustion(|| format!("opening {index_path}")) + .map_err(|source| { + error!( + stream_id = identity.stream_id, + topic_id = identity.topic_id, + partition_id = identity.partition_id, + path = %index_path, + error = %source, + "failed to open sparse index for anchor search during recovery" + ); + PartitionRecoveryError::from(IggyError::CannotReadFile) + })?; let mut raw = [0u8; IGGY_INDEX_SIZE]; // Lowest log byte a probed entry has already paid for; the next probed // entry is budgeted by the span from its own position up to here. @@ -2367,28 +2376,32 @@ async fn index_is_consistent( messages_size: u64, scratch: &mut ScanScratch, ) -> Result<IndexValidation, PartitionRecoveryError> { - let index_file = fs::File::open(index_path).map_err(|source| { - error!( - stream_id = identity.stream_id, - topic_id = identity.topic_id, - partition_id = identity.partition_id, - path = %index_path, - error = %source, - "failed to open sparse index for validation during recovery" - ); - PartitionRecoveryError::from(IggyError::CannotReadFile) - })?; - let messages_file = fs::File::open(messages_path).map_err(|source| { - error!( - stream_id = identity.stream_id, - topic_id = identity.topic_id, - partition_id = identity.partition_id, - path = %messages_path, - error = %source, - "failed to open segment log for sparse index validation" - ); - PartitionRecoveryError::from(IggyError::CannotReadFile) - })?; + let index_file = fs::File::open(index_path) + .exit_on_descriptor_exhaustion(|| format!("opening {index_path}")) + .map_err(|source| { + error!( + stream_id = identity.stream_id, + topic_id = identity.topic_id, + partition_id = identity.partition_id, + path = %index_path, + error = %source, + "failed to open sparse index for validation during recovery" + ); + PartitionRecoveryError::from(IggyError::CannotReadFile) + })?; + let messages_file = fs::File::open(messages_path) + .exit_on_descriptor_exhaustion(|| format!("opening {messages_path}")) + .map_err(|source| { + error!( + stream_id = identity.stream_id, + topic_id = identity.topic_id, + partition_id = identity.partition_id, + path = %messages_path, + error = %source, + "failed to open segment log for sparse index validation" + ); + PartitionRecoveryError::from(IggyError::CannotReadFile) + })?; let ScanScratch { window: index_window, spill: log_window, @@ -2503,17 +2516,19 @@ fn open_messages_file( identity: PartitionIdentity<'_>, messages_path: &str, ) -> Result<fs::File, PartitionRecoveryError> { - fs::File::open(messages_path).map_err(|source| { - error!( - stream_id = identity.stream_id, - topic_id = identity.topic_id, - partition_id = identity.partition_id, - path = %messages_path, - error = %source, - "failed to open a segment messages file during recovery" - ); - PartitionRecoveryError::from(IggyError::CannotReadFile) - }) + fs::File::open(messages_path) + .exit_on_descriptor_exhaustion(|| format!("opening {messages_path}")) + .map_err(|source| { + error!( + stream_id = identity.stream_id, + topic_id = identity.topic_id, + partition_id = identity.partition_id, + path = %messages_path, + error = %source, + "failed to open a segment messages file during recovery" + ); + PartitionRecoveryError::from(IggyError::CannotReadFile) + }) } /// The batch header at `position`, or `None` when the walk must stop there diff --git a/core/partitions/src/state_transfer.rs b/core/partitions/src/state_transfer.rs index e4ffae9e2..5832854b3 100644 --- a/core/partitions/src/state_transfer.rs +++ b/core/partitions/src/state_transfer.rs @@ -49,6 +49,7 @@ use journal::durable_storage::{DiskStorage, DurableStorage}; use journal::superblock::SuperblockStore; use message_bus::MessageBus; use server_common::Message; +use server_common::fatal::{ExitOnDescriptorExhaustion, exit_if_descriptors_exhausted}; use server_common::iobuf::Owned; use server_common::send_messages::{decode_batch_slice, decode_prepare_slice}; use server_common::{SegmentStorage, yield_to_reactor}; @@ -1548,7 +1549,8 @@ fn staging_paths(partition_dir: &str, start_offset: u64) -> (PathBuf, PathBuf) { /// error policies (propagate / silent skip / log-and-fail), which is what /// `sweep_staging_except`'s do-not-widen warning depends on. fn segment_dir_entries(partition_dir: &str) -> std::io::Result<Vec<PathBuf>> { - Ok(std::fs::read_dir(partition_dir)? + Ok(std::fs::read_dir(partition_dir) + .exit_on_descriptor_exhaustion(|| format!("listing directory {partition_dir}"))? .flatten() .map(|entry| entry.path()) .collect()) @@ -1593,7 +1595,8 @@ pub async fn mark_materialization_missing(directory: &str, revision: u64) -> std .truncate(true) .write(true) .open(&temporary) - .await?; + .await + .exit_on_descriptor_exhaustion(|| format!("opening {}", temporary.display()))?; file.write_all_at(revision.to_le_bytes().to_vec(), 0) .await .0?; @@ -1835,7 +1838,8 @@ async fn discard_offset_writes(planned: &[PlannedOffsetWrite]) { /// every other future on the pump keeps running through it. pub(crate) async fn fsync_dir(partition_dir: &str) -> std::io::Result<()> { compio::fs::File::open(partition_dir) - .await? + .await + .exit_on_descriptor_exhaustion(|| format!("opening directory {partition_dir}"))? .sync_all() .await } @@ -2997,6 +3001,7 @@ where // is an `open`+`close` per segment for no durability gain. let dir_handle = compio::fs::File::open(partition_dir) .await + .exit_on_descriptor_exhaustion(|| format!("opening directory {partition_dir}")) .map_err(|source| PartitionInstallError::SwapIo { path: partition_dir.to_owned(), source, @@ -3748,7 +3753,9 @@ async fn hash_segment_range( if from >= to { return Ok(()); } - let file = compio::fs::File::open(path).await?; + let file = compio::fs::File::open(path) + .await + .exit_on_descriptor_exhaustion(|| format!("opening {path}"))?; let mut position = from; // One buffer for the whole pass; `BufResult` hands it back per read // precisely so the alloc + memset are not paid per chunk. Re-allocated @@ -3838,6 +3845,7 @@ impl SegmentLoadError { // Everything unrecognised stays STALE: a short read past EOF is what a // racing GC unlink-and-recreate legitimately produces. const EIO: i32 = 5; + exit_if_descriptors_exhausted(&source, "loading a segment for state transfer"); if source.raw_os_error() == Some(EIO) { return Self::LocalFault(source); } @@ -3921,7 +3929,9 @@ enum OffsetDirEntry { /// alone rather than guessed at: every offset file is named by its id, so /// anything else is not ours. fn offset_dir_entries(dir: &str) -> Vec<OffsetDirEntry> { - let Ok(entries) = std::fs::read_dir(dir) else { + let Ok(entries) = + std::fs::read_dir(dir).exit_on_descriptor_exhaustion(|| format!("listing directory {dir}")) + else { return Vec::new(); }; entries @@ -3980,7 +3990,9 @@ fn segment_manifest_digest(manifest: &[consensus::StateArtifact]) -> u64 { } async fn write_staging_file(path: &Path, payload: Vec<u8>) -> std::io::Result<()> { - let mut file = compio::fs::File::create(path).await?; + let mut file = compio::fs::File::create(path) + .await + .exit_on_descriptor_exhaustion(|| format!("creating {}", path.display()))?; let (result, _) = file.write_all_at(payload, 0).await.into(); result?; file.sync_data().await?; diff --git a/core/server/config.toml b/core/server/config.toml index 33fcb3a65..e91b9f669 100644 --- a/core/server/config.toml +++ b/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, +# messages, open file descriptors) to the log from shard 0. "0" disables the +# line. +sysinfo_print_interval = "10 s" + # Encryption configuration [encryption] # Encrypt message payloads and user headers. Metadata and structural headers remain unencrypted. @@ -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 +# create a topic or partitions that would exceed it fails with +# PartitionsLimitReached. The check counts committed partitions only, so +# concurrent creates can go slightly over it. 0 means no cap. +partitions_max = 0 + # Per-partition consensus plane tunables. Unlike [metadata] (one shard-0 # plane), a pipeline exists per partition, so raising this multiplies pinned # request-buffer memory by the partition count. Keep it modest. diff --git a/core/server/src/boot/fd_limit.rs b/core/server/src/boot/fd_limit.rs new file mode 100644 index 000000000..18356a206 --- /dev/null +++ b/core/server/src/boot/fd_limit.rs @@ -0,0 +1,116 @@ +// 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. + +//! The process's own `RLIMIT_NOFILE`, raised once at startup. + +use nix::errno::Errno; +use nix::sys::resource::{Resource, getrlimit, setrlimit}; +use thiserror::Error; + +/// `OPEN_MAX` from `<sys/syslimits.h>`, which `libc` does not export. The +/// macOS hard limit is usually `RLIM_INFINITY`, and `setrlimit(2)` rejects +/// that as a soft `RLIMIT_NOFILE` with `EINVAL`, so the man page's recipe is +/// `min(OPEN_MAX, rlim_max)`. +#[cfg(target_vendor = "apple")] +const APPLE_OPEN_MAX: u64 = 10_240; + +/// `RLIMIT_NOFILE` around [`raise_open_file_limit`]. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct OpenFileLimit { + pub soft_before: u64, + pub soft: u64, + pub hard: u64, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Error)] +pub enum OpenFileLimitError { + #[error("cannot read RLIMIT_NOFILE: {0}")] + Read(Errno), + #[error( + "cannot raise the RLIMIT_NOFILE soft limit from {soft} to {target} (hard {hard}): {errno}" + )] + Raise { + errno: Errno, + soft: u64, + hard: u64, + target: u64, + }, +} + +/// Raise the soft `RLIMIT_NOFILE` to the hard limit (clamped on macOS). +/// A soft limit already at or above the target is left as it is. +/// +/// # Errors +/// +/// [`OpenFileLimitError::Read`] if the limit cannot be read, and +/// [`OpenFileLimitError::Raise`] if `setrlimit` rejects the new soft limit. +pub fn raise_open_file_limit() -> Result<OpenFileLimit, OpenFileLimitError> { + let (soft_before, hard) = + getrlimit(Resource::RLIMIT_NOFILE).map_err(OpenFileLimitError::Read)?; + let target = soft_target(hard); + if soft_before >= target { + return Ok(OpenFileLimit { + soft_before, + soft: soft_before, + hard, + }); + } + setrlimit(Resource::RLIMIT_NOFILE, target, hard).map_err(|errno| { + OpenFileLimitError::Raise { + errno, + soft: soft_before, + hard, + target, + } + })?; + Ok(OpenFileLimit { + soft_before, + soft: target, + hard, + }) +} + +#[cfg(target_vendor = "apple")] +const fn soft_target(hard: u64) -> u64 { + if hard < APPLE_OPEN_MAX { + hard + } else { + APPLE_OPEN_MAX + } +} + +#[cfg(not(target_vendor = "apple"))] +const fn soft_target(hard: u64) -> u64 { + hard +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn given_process_limit_when_raising_should_leave_reported_soft_limit_in_effect() { + let limit = raise_open_file_limit().expect("RLIMIT_NOFILE must be raisable in tests"); + + let (soft, hard) = getrlimit(Resource::RLIMIT_NOFILE).expect("RLIMIT_NOFILE readable"); + assert_eq!(soft, limit.soft); + assert_eq!(hard, limit.hard); + assert!(limit.soft >= limit.soft_before); + #[cfg(target_os = "linux")] + assert_eq!(soft, hard); + } +} diff --git a/core/server/src/boot/mod.rs b/core/server/src/boot/mod.rs index 0946af00e..7cb1c6a1c 100644 --- a/core/server/src/boot/mod.rs +++ b/core/server/src/boot/mod.rs @@ -24,6 +24,7 @@ //! the leaves hold the support it calls into. mod credentials; +mod fd_limit; mod handoff; mod listeners; mod recovery; @@ -34,6 +35,7 @@ mod topology; pub use crate::dispatch::host::ServerHost; pub use credentials::apply_default_root_credentials; +pub use fd_limit::{OpenFileLimit, OpenFileLimitError, raise_open_file_limit}; pub use threads::ShardHandles; use crate::boot::credentials::{ @@ -823,12 +825,29 @@ async fn shard_main( } else { None }; + + // Sysinfo printer: shard 0 only, since the line describes the whole + // process. A zero interval disables it. + let sysinfo_print_interval = config.logging.sysinfo_print_interval; + let sysinfo_printer_stop = if shard_id == 0 && !sysinfo_print_interval.is_zero() { + let (stop_tx, stop_rx) = channel(1); + let printer_shard = Rc::clone(&shard); + let interval = sysinfo_print_interval.get_duration(); + let printer_handle = compio::runtime::spawn(async move { + crate::sysinfo_printer::run_sysinfo_printer(printer_shard, stop_rx, interval).await; + }); + bus.track_background(printer_handle); + Some(stop_tx) + } else { + None + }; let mut stop_signals = StopSignals { pump: stop_tx, reconciler: reconcile_stop_tx, heartbeat: heartbeat_stop_tx, pat_cleaner: pat_cleaner_stop, segment_cleaner: segment_cleaner_stop, + sysinfo_printer: sysinfo_printer_stop, consumer_group_liveness: None, }; diff --git a/core/server/src/boot/threads.rs b/core/server/src/boot/threads.rs index 29d834284..12bfec460 100644 --- a/core/server/src/boot/threads.rs +++ b/core/server/src/boot/threads.rs @@ -845,6 +845,7 @@ pub(in crate::boot) struct StopSignals { pub(in crate::boot) heartbeat: Option<Sender<()>>, pub(in crate::boot) pat_cleaner: Option<Sender<()>>, pub(in crate::boot) segment_cleaner: Option<Sender<()>>, + pub(in crate::boot) sysinfo_printer: Option<Sender<()>>, pub(in crate::boot) consumer_group_liveness: Option<Sender<()>>, } @@ -857,6 +858,7 @@ impl StopSignals { &self.heartbeat, &self.pat_cleaner, &self.segment_cleaner, + &self.sysinfo_printer, &self.consumer_group_liveness, ] .into_iter() diff --git a/core/server/src/dispatch/mod.rs b/core/server/src/dispatch/mod.rs index 59ac32d8f..b091c5449 100644 --- a/core/server/src/dispatch/mod.rs +++ b/core/server/src/dispatch/mod.rs @@ -663,6 +663,7 @@ async fn handle_client_request<B, MJ, S, SB>( sessions, transport_client_id, max_tokens_per_user, + server_config.metadata.partitions_max, request, ) { Ok(rewritten) => rewritten, diff --git a/core/server/src/http/handlers.rs b/core/server/src/http/handlers.rs index 778692189..d0aba10e0 100644 --- a/core/server/src/http/handlers.rs +++ b/core/server/src/http/handlers.rs @@ -146,8 +146,9 @@ use crate::http::wire::{ use crate::reply_frame::{build_polled_messages_body, build_raw_pat_reply}; use crate::responses::connected_client_to_response; use crate::rewrite::{ - validate_option_keys, validate_topic_bounds, validate_topic_size_floor, - warn_unenforceable_topic_size, warn_unenforceable_topic_size_on_partition_add, + validate_option_keys, validate_partitions_limit, validate_topic_bounds, + validate_topic_size_floor, warn_unenforceable_topic_size, + warn_unenforceable_topic_size_on_partition_add, }; use crate::snapshot; @@ -998,6 +999,20 @@ pub(in crate::http) async fn create_topic( let max_topic_size = parsed.max_topic_size.unwrap_or(MaxTopicSize::ServerDefault); validate_topic_bounds(command.partitions_count, max_topic_size, segment_size) .map_err(WriteError::Rejected)?; + validate_partitions_limit( + state.server_config.metadata.partitions_max, + command.partitions_count, + || { + state + .shard + .plane + .metadata() + .mux_stm + .streams() + .partition_count() + }, + ) + .map_err(WriteError::Rejected)?; warn_unenforceable_topic_size( max_topic_size, segment_size, @@ -1158,6 +1173,12 @@ pub(in crate::http) async fn create_partitions( partitions_count: command.partitions_count, }; let metadata = state.shard.plane.metadata(); + validate_partitions_limit( + state.server_config.metadata.partitions_max, + request.partitions_count, + || metadata.mux_stm.streams().partition_count(), + ) + .map_err(WriteError::Rejected)?; warn_unenforceable_topic_size_on_partition_add( metadata.mux_stm.streams(), &request.stream_id, diff --git a/core/server/src/http/tls.rs b/core/server/src/http/tls.rs index 73f4108b8..e709de4ab 100644 --- a/core/server/src/http/tls.rs +++ b/core/server/src/http/tls.rs @@ -33,6 +33,7 @@ //! cross-transport invariant in `message_bus::client_listener`). Handshaken //! streams flow to the serve loop over a bounded channel. +use message_bus::accept::pause_after_accept_error; use std::future::Future; use std::net::SocketAddr; use std::path::Path; @@ -231,7 +232,10 @@ async fn accept_pump( Ok((stream, peer)) => { spawn_handshake(&acceptor, &connections, handshake_grace, stream, peer); } - Err(error) => error!(%error, "server HTTPS accept failed"), + Err(error) => { + error!(%error, "server HTTPS accept failed"); + pause_after_accept_error(&error).await; + } }, } } diff --git a/core/server/src/lib.rs b/core/server/src/lib.rs index f2273799b..998fcb2a0 100644 --- a/core/server/src/lib.rs +++ b/core/server/src/lib.rs @@ -62,6 +62,7 @@ pub(crate) mod partition_reconciler; pub(crate) mod personal_access_token_cleaner; pub(crate) mod segment_cleaner; pub(crate) mod snapshot; +pub(crate) mod sysinfo_printer; // support: shared plumbing. pub(crate) mod cluster_meta; diff --git a/core/server/src/main.rs b/core/server/src/main.rs index 8459a7100..d416cd182 100644 --- a/core/server/src/main.rs +++ b/core/server/src/main.rs @@ -23,11 +23,14 @@ mod banner; use args::Args; use clap::Parser; use configs::server::ServerConfig; -use server::boot::{apply_default_root_credentials, bootstrap, load_config, prepare_runtime_dirs}; +use server::boot::{ + apply_default_root_credentials, bootstrap, load_config, prepare_runtime_dirs, + raise_open_file_limit, +}; use server::server_error::ServerError; use server_common::log::Logging; use system_stats::capture_allowed_cpus; -use tracing::{error, info}; +use tracing::{error, info, warn}; fn main() -> Result<(), ServerError> { // This prelude must stay ahead of the first thread the process ever @@ -47,7 +50,7 @@ fn main() -> Result<(), ServerError> { #[cfg(all(feature = "mimalloc", not(feature = "disable-mimalloc")))] info!("Using mimalloc allocator"); #[cfg(not(all(feature = "mimalloc", not(feature = "disable-mimalloc"))))] - tracing::warn!("Using the default system allocator"); + warn!("Using the default system allocator"); if let Ok(env_path) = std::env::var("IGGY_ENV_PATH") { let _ = dotenvy::from_path(&env_path); } else { @@ -59,6 +62,19 @@ fn main() -> Result<(), ServerError> { // Before shard threads pin themselves: a pinned capture sees one core. capture_allowed_cpus(); + // Before bootstrap: partition persistence sizes its offset-file budget + // from the soft limit it reads first, and nothing else raises it except a + // side effect of sysinfo's first process refresh on Linux. + match raise_open_file_limit() { + Ok(limit) => info!( + soft_before = limit.soft_before, + soft = limit.soft, + hard = limit.hard, + "open-file limit (RLIMIT_NOFILE) set" + ), + Err(error) => warn!(error = %error, "open-file limit (RLIMIT_NOFILE) left unchanged"), + } + let bootstrap_runtime = match server_common::create_shard_executor() { Ok(rt) => rt, Err(e) => { diff --git a/core/server/src/responses.rs b/core/server/src/responses.rs index 05e1995be..14c075a56 100644 --- a/core/server/src/responses.rs +++ b/core/server/src/responses.rs @@ -494,10 +494,21 @@ fn aggregate_stats_totals( )) } -fn build_stats_response<B, MJ, S, SB>( +/// Node-wide object and message totals, shared by the `GetStats` reply and +/// the periodic sysinfo log line. +pub struct StatsTotals { + pub streams_count: u32, + pub topics_count: u32, + pub partitions_count: u32, + pub segments_count: u32, + pub messages_size_bytes: u64, + pub messages_count: u64, + pub consumer_groups_count: u32, +} + +pub fn stats_totals<B, MJ, S, SB>( shard: &Rc<ShellShard<B, MJ, S, SB>>, - clients_count: u32, -) -> Result<StatsResponse, IggyError> +) -> Result<StatsTotals, IggyError> where B: ShellBus, MJ: JournalHandle + 'static, @@ -526,7 +537,29 @@ where .streams() .consumer_group_count(), )?; + Ok(StatsTotals { + streams_count, + topics_count, + partitions_count, + segments_count, + messages_size_bytes, + messages_count, + consumer_groups_count, + }) +} +fn build_stats_response<B, MJ, S, SB>( + shard: &Rc<ShellShard<B, MJ, S, SB>>, + clients_count: u32, +) -> Result<StatsResponse, IggyError> +where + B: ShellBus, + MJ: JournalHandle + 'static, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, + S: 'static, + SB: SuperblockStore + 'static, +{ + let totals = stats_totals(shard)?; let system = probe_system_stats(); let (free_disk_space, total_disk_space) = stats_disk_space(); Ok(StatsResponse { @@ -540,14 +573,14 @@ where start_time: system.start_time, read_bytes: system.read_bytes, written_bytes: system.written_bytes, - messages_size_bytes, - streams_count, - topics_count, - partitions_count, - segments_count, - messages_count, + messages_size_bytes: totals.messages_size_bytes, + streams_count: totals.streams_count, + topics_count: totals.topics_count, + partitions_count: totals.partitions_count, + segments_count: totals.segments_count, + messages_count: totals.messages_count, clients_count, - consumer_groups_count, + consumer_groups_count: totals.consumer_groups_count, hostname: system.hostname, os_name: system.os_name, os_version: system.os_version, diff --git a/core/server/src/rewrite.rs b/core/server/src/rewrite.rs index 5f00c588c..7f4749bc4 100644 --- a/core/server/src/rewrite.rs +++ b/core/server/src/rewrite.rs @@ -113,6 +113,7 @@ pub fn tcp_chain<B, MJ, S, SB>( sessions: &Rc<RefCell<SessionManager>>, transport_client_id: u128, max_tokens_per_user: u32, + partitions_max: u32, request: Message<RoutedRequestHeader>, ) -> Result<(Message<RoutedRequestHeader>, Option<String>), RewriteDeny> where @@ -151,7 +152,7 @@ where stage: RewriteStage::UserPassword, error, })?; - static_bounds(shard, &request).map_err(|error| RewriteDeny { + static_bounds(shard, &request, partitions_max).map_err(|error| RewriteDeny { stage: RewriteStage::StaticBounds, error, })?; @@ -234,6 +235,31 @@ const fn validate_partitions_change_count(partitions_count: u32) -> Result<(), I validate_partitions_count(partitions_count) } +/// Node-wide `[metadata] partitions_max` cap, shared by the TCP and HTTP +/// ingresses for create-topic and create-partitions. Zero is no cap, and +/// `committed` is read only when a cap is set. +/// +/// A soft cap, like the PAT limit: it runs pre-consensus because the metadata +/// apply must not branch on node config, so it sees only this node's committed +/// count. Creates in flight together can overshoot it, and a lagging backup +/// counts fewer partitions than the primary. +pub fn validate_partitions_limit( + partitions_max: u32, + requested: u32, + committed: impl FnOnce() -> usize, +) -> Result<(), IggyError> { + if partitions_max == 0 { + return Ok(()); + } + let total = u64::try_from(committed()) + .unwrap_or(u64::MAX) + .saturating_add(u64::from(requested)); + if total > u64::from(partitions_max) { + return Err(IggyError::PartitionsLimitReached); + } + Ok(()) +} + /// Static create-topic bounds shared by the TCP and HTTP ingresses. Runs /// pre-consensus: a rejected request must not burn a replicated log entry, /// and `prepare_request` errors evict the session instead of denying typed. @@ -361,6 +387,7 @@ pub fn validate_option_keys(options: &WireOptions, known: &[&str]) -> Result<(), fn static_bounds<B, MJ, S, SB>( shard: &Rc<ShellShard<B, MJ, S, SB>>, request: &Message<RoutedRequestHeader>, + partitions_max: u32, ) -> Result<(), IggyError> where B: ShellBus, @@ -396,6 +423,9 @@ where .max_topic_size .unwrap_or(MaxTopicSize::ServerDefault); validate_topic_bounds(create_topic.partitions_count, max_topic_size, segment_size)?; + validate_partitions_limit(partitions_max, create_topic.partitions_count, || { + shard.plane.metadata().mux_stm.streams().partition_count() + })?; warn_unenforceable_topic_size( max_topic_size, segment_size, @@ -409,6 +439,11 @@ where .and_then(|create_partitions| { validate_partitions_change_count(create_partitions.partitions_count)?; let metadata = shard.plane.metadata(); + validate_partitions_limit( + partitions_max, + create_partitions.partitions_count, + || metadata.mux_stm.streams().partition_count(), + )?; warn_unenforceable_topic_size_on_partition_add( metadata.mux_stm.streams(), &create_partitions.stream_id, @@ -545,6 +580,32 @@ mod tests { assert!(validate_partitions_count(0).is_ok()); } + #[test] + fn given_no_partitions_cap_when_validating_should_admit_without_counting() { + let admitted = validate_partitions_limit(0, u32::MAX, || { + panic!("the committed count must not be read without a cap") + }); + + assert!(admitted.is_ok()); + } + + #[test] + fn given_create_reaching_partitions_cap_when_validating_should_admit() { + assert!(validate_partitions_limit(10, 4, || 6).is_ok()); + } + + #[test] + fn given_create_past_partitions_cap_when_validating_should_deny() { + assert!(matches!( + validate_partitions_limit(10, 5, || 6), + Err(IggyError::PartitionsLimitReached) + )); + assert!(matches!( + validate_partitions_limit(10, 1, || 10), + Err(IggyError::PartitionsLimitReached) + )); + } + #[test] fn zero_partitions_change_denies_pre_consensus() { // Adding or removing zero partitions is a no-op that would still burn diff --git a/core/server/src/sysinfo_printer.rs b/core/server/src/sysinfo_printer.rs new file mode 100644 index 000000000..e453c7591 --- /dev/null +++ b/core/server/src/sysinfo_printer.rs @@ -0,0 +1,197 @@ +// 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, probe_system_stats, stats_disk_space}; +use iggy_common::IggyByteSize; +use nix::sys::resource::{Resource, getrlimit}; +use shard::Receiver; +use std::fmt; +use std::rc::Rc; +use std::time::Duration; +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:?}"); + loop { + // `Ok(_)`: stop signalled -> exit. `Err(_)`: interval elapsed -> print. + match compio::time::timeout(interval, stop.recv()).await { + Ok(_) => break, + Err(_) => print_sysinfo(&shard).await, + } + } + trace!(shard = shard.id, "sysinfo printer exited"); +} + +async fn print_sysinfo(shard: &Rc<ServerShard>) { + let clients_count = shard.list_all_clients().await.len(); + let totals = match stats_totals(shard) { + Ok(totals) => totals, + Err(error) => { + error!(error = %error, "Failed to get system information"); + return; + } + }; + let (free_disk_space, total_disk_space) = stats_disk_space(); + let line = SysinfoLine { + system: probe_system_stats(), + totals, + clients_count, + free_disk_space, + total_disk_space, + open_files_limit: getrlimit(Resource::RLIMIT_NOFILE) + .ok() + .map(|(soft, _)| soft), + }; + info!("{line}"); +} + +/// One sample, rendered in the 0.8.2 server's layout plus open descriptors. +struct SysinfoLine { + system: SystemStats, + totals: StatsTotals, + clients_count: usize, + free_disk_space: u64, + total_disk_space: u64, + /// Soft `RLIMIT_NOFILE`: the ceiling `open()` actually fails at. + open_files_limit: Option<u64>, +} + +impl SysinfoLine { + // Precision loss starts above 2^53 bytes, far beyond any host's memory. + #[allow(clippy::cast_precision_loss)] + fn free_memory_percent(&self) -> f64 { + if self.system.total_memory == 0 { + return 0.0; + } + self.system.available_memory as f64 / self.system.total_memory as f64 * 100.0 + } +} + +impl fmt::Display for SysinfoLine { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + let system = &self.system; + write!( + f, + "CPU: {:.2}%/{:.2}% (IggyUsage/Total), Mem: {:.2}%/{}/{}/{} (Free/IggyUsage/TotalUsed/Total), Disk: {}/{} (Free/Total), IggyUsage: {}, Clients: {}, Messages: {}, Read: {}, Written: {}", + system.cpu_usage, + system.total_cpu_usage, + self.free_memory_percent(), + IggyByteSize::from(system.memory_usage), + IggyByteSize::from(system.total_memory.saturating_sub(system.available_memory)), + IggyByteSize::from(system.total_memory), + IggyByteSize::from(self.free_disk_space), + IggyByteSize::from(self.total_disk_space), + IggyByteSize::from(self.totals.messages_size_bytes), + self.clients_count, + self.totals.messages_count, + IggyByteSize::from(system.read_bytes), + IggyByteSize::from(system.written_bytes), + )?; + if system.threads_count > 0 { + write!(f, ", Threads: {}", system.threads_count)?; + } + if let Some(open_files) = system.open_files { + match self.open_files_limit { + Some(limit) => write!(f, ", OpenFDs: {open_files}/{limit} (Current/Max)")?, + None => write!(f, ", OpenFDs: {open_files}/unknown (Current/Max)")?, + } + } + Ok(()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn line(open_files: Option<u64>, open_files_limit: Option<u64>) -> SysinfoLine { + SysinfoLine { + system: SystemStats { + process_id: 1, + cpu_usage: 1.5, + total_cpu_usage: 20.0, + memory_usage: 1_000_000, + total_memory: 4_000_000, + available_memory: 1_000_000, + run_time: 0, + start_time: 0, + read_bytes: 0, + written_bytes: 0, + threads_count: 8, + open_files, + hostname: String::new(), + os_name: String::new(), + os_version: String::new(), + kernel_version: String::new(), + }, + totals: StatsTotals { + streams_count: 1, + topics_count: 1, + partitions_count: 1, + segments_count: 1, + messages_size_bytes: 0, + messages_count: 42, + consumer_groups_count: 0, + }, + clients_count: 3, + free_disk_space: 0, + total_disk_space: 0, + open_files_limit, + } + } + + #[test] + fn given_open_files_and_limit_when_rendering_should_print_current_and_max() { + let rendered = line(Some(12), Some(1024)).to_string(); + + assert!(rendered.starts_with("CPU: 1.50%/20.00% (IggyUsage/Total), Mem: 25.00%/")); + assert!(rendered.contains(", Clients: 3, Messages: 42, ")); + assert!(rendered.ends_with(", Threads: 8, OpenFDs: 12/1024 (Current/Max)")); + } + + #[test] + fn given_unreadable_limit_when_rendering_should_print_unknown_max() { + let rendered = line(Some(12), None).to_string(); + + assert!(rendered.ends_with(", OpenFDs: 12/unknown (Current/Max)")); + } + + #[test] + fn given_uncountable_open_files_when_rendering_should_omit_open_fds() { + let rendered = line(None, Some(1024)).to_string(); + + assert!(!rendered.contains("OpenFDs")); + } + + #[test] + fn given_zero_total_memory_when_rendering_should_report_zero_free_percent() { + let mut sample = line(None, None); + sample.system.total_memory = 0; + sample.system.available_memory = 0; + + assert!(sample.to_string().contains("Mem: 0.00%/")); + } +} diff --git a/core/server/src/sysinfo_probe.rs b/core/server/src/sysinfo_probe.rs index 206e59f35..0e08813f3 100644 --- a/core/server/src/sysinfo_probe.rs +++ b/core/server/src/sysinfo_probe.rs @@ -42,6 +42,9 @@ pub struct SystemStats { pub read_bytes: u64, pub written_bytes: u64, pub threads_count: u32, + /// Only the periodic sysinfo log line reads this. The `GetStats` wire + /// reply has no field for it. + pub open_files: Option<u64>, pub hostname: String, pub os_name: String, pub os_version: String, @@ -127,6 +130,7 @@ pub fn probe_system_stats() -> SystemStats { read_bytes: probe.read_bytes, written_bytes: probe.written_bytes, threads_count: probe.threads_count, + open_files: probe.open_files, hostname: host.hostname.clone(), os_name: host.os_name.clone(), os_version: host.os_version.clone(), diff --git a/core/server_common/src/fatal.rs b/core/server_common/src/fatal.rs new file mode 100644 index 000000000..93aa0564d --- /dev/null +++ b/core/server_common/src/fatal.rs @@ -0,0 +1,160 @@ +// 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 failed because the process (`EMFILE`) or the host + /// (`ENFILE`) has no free file descriptor. Every later open on every shard + /// fails the same way, so serving on only fails each 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.consensus.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. +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 { + exit_if_descriptors_exhausted(error, &operation()); + } + self + } +} + +/// Stop the process with [`FatalReason::DescriptorsExhausted`] if `error` is +/// `EMFILE` or `ENFILE`, and return otherwise. +/// +/// Storage paths only. An accept loop that stopped here would let any client +/// that can open sockets stop the node. +pub fn exit_if_descriptors_exhausted(error: &io::Error, operation: &str) { + if !is_descriptor_exhaustion(error) { + return; + } + 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})"), + ); +} + +/// `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) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn given_emfile_or_enfile_when_classifying_should_report_exhaustion() { + for errno in [Errno::EMFILE, Errno::ENFILE] { + assert!(is_descriptor_exhaustion(&io::Error::from_raw_os_error( + errno as i32 + ))); + } + } + + #[test] + fn given_other_error_when_classifying_should_not_report_exhaustion() { + assert!(!is_descriptor_exhaustion(&io::Error::from_raw_os_error( + Errno::ENOENT as i32 + ))); + assert!(!is_descriptor_exhaustion(&io::Error::other( + "Too many open files" + ))); + } + + #[test] + fn given_non_exhaustion_error_when_checking_result_should_pass_it_through() { + let result: io::Result<()> = Err(io::Error::from_raw_os_error(Errno::ENOENT as i32)); + + let checked = result.exit_on_descriptor_exhaustion(|| "opening a test file".to_owned()); + + assert_eq!( + checked.map_err(|error| error.kind()), + Err(io::ErrorKind::NotFound) + ); + } +} diff --git a/core/server_common/src/fs_utils.rs b/core/server_common/src/fs_utils.rs index c64090c56..3f6cf95bb 100644 --- a/core/server_common/src/fs_utils.rs +++ b/core/server_common/src/fs_utils.rs @@ -15,6 +15,7 @@ // specific language governing permissions and limitations // under the License. +use crate::fatal::ExitOnDescriptorExhaustion; use compio::fs; use std::io; use std::path::{Path, PathBuf}; @@ -115,7 +116,10 @@ pub async fn walk_dir(root: impl AsRef<Path>) -> io::Result<Vec<DirEntry>> { current_dir.to_str().map(|s| s.to_string()), )); - for entry in std::fs::read_dir(¤t_dir)? { + let entries = std::fs::read_dir(¤t_dir).exit_on_descriptor_exhaustion(|| { + format!("listing directory {}", current_dir.display()) + })?; + for entry in entries { let entry = entry?; let entry_path = entry.path(); let metadata = fs::symlink_metadata(&entry_path).await?; diff --git a/core/server_common/src/lib.rs b/core/server_common/src/lib.rs index e7b8c8151..8171035bc 100644 --- a/core/server_common/src/lib.rs +++ b/core/server_common/src/lib.rs @@ -22,6 +22,7 @@ mod consensus_message; pub mod crypto; pub mod diagnostics; pub mod executor; +pub mod fatal; pub mod fs_utils; pub mod iobuf; pub mod log; diff --git a/core/server_common/src/segment_storage/index_reader.rs b/core/server_common/src/segment_storage/index_reader.rs index b4fb65140..7f5cea46d 100644 --- a/core/server_common/src/segment_storage/index_reader.rs +++ b/core/server_common/src/segment_storage/index_reader.rs @@ -15,6 +15,7 @@ // specific language governing permissions and limitations // under the License. +use crate::fatal::ExitOnDescriptorExhaustion; use compio::fs::OpenOptions; use err_trail::ErrContext; use iggy_common::IggyError; @@ -36,6 +37,7 @@ impl IndexReader { .read(true) .open(file_path) .await + .exit_on_descriptor_exhaustion(|| format!("opening {file_path}")) .error(|e: &std::io::Error| format!("Failed to open index file: {file_path}. {e}")) .map_err(|_| IggyError::CannotReadFile)?; diff --git a/core/server_common/src/segment_storage/index_writer.rs b/core/server_common/src/segment_storage/index_writer.rs index 8740a8bb4..1a13b1c8c 100644 --- a/core/server_common/src/segment_storage/index_writer.rs +++ b/core/server_common/src/segment_storage/index_writer.rs @@ -15,6 +15,7 @@ // specific language governing permissions and limitations // under the License. +use crate::fatal::ExitOnDescriptorExhaustion; use compio::fs::File; use compio::fs::OpenOptions; use err_trail::ErrContext; @@ -50,6 +51,7 @@ impl IndexWriter { let file = opts .open(file_path) .await + .exit_on_descriptor_exhaustion(|| format!("opening {file_path}")) .error(|e: &std::io::Error| format!("Failed to open index file: {file_path}. {e}")) .map_err(|_| IggyError::CannotReadFile)?; diff --git a/core/server_common/src/segment_storage/messages_reader.rs b/core/server_common/src/segment_storage/messages_reader.rs index f7705713a..ba6f2859c 100644 --- a/core/server_common/src/segment_storage/messages_reader.rs +++ b/core/server_common/src/segment_storage/messages_reader.rs @@ -15,6 +15,7 @@ // specific language governing permissions and limitations // under the License. +use crate::fatal::ExitOnDescriptorExhaustion; use compio::fs::OpenOptions; use err_trail::ErrContext; use iggy_common::IggyError; @@ -37,6 +38,7 @@ impl MessagesReader { .read(true) .open(file_path) .await + .exit_on_descriptor_exhaustion(|| format!("opening {file_path}")) .error(|e: &std::io::Error| format!("Failed to open messages file: {file_path}. {e}")) .map_err(|_| IggyError::CannotReadFile)?; diff --git a/core/server_common/src/segment_storage/messages_writer.rs b/core/server_common/src/segment_storage/messages_writer.rs index 002402b8b..d1320f950 100644 --- a/core/server_common/src/segment_storage/messages_writer.rs +++ b/core/server_common/src/segment_storage/messages_writer.rs @@ -15,6 +15,7 @@ // specific language governing permissions and limitations // under the License. +use crate::fatal::ExitOnDescriptorExhaustion; use compio::fs::{File, OpenOptions}; use err_trail::ErrContext; use iggy_common::IggyError; @@ -56,6 +57,7 @@ impl MessagesWriter { let file = opts .open(file_path) .await + .exit_on_descriptor_exhaustion(|| format!("opening {file_path}")) .error(|err: &std::io::Error| { format!("Failed to open messages file: {file_path}, error: {err}") }) diff --git a/core/server_common/src/segment_storage/mod.rs b/core/server_common/src/segment_storage/mod.rs index d389f6c04..a39d7f94c 100644 --- a/core/server_common/src/segment_storage/mod.rs +++ b/core/server_common/src/segment_storage/mod.rs @@ -20,6 +20,7 @@ mod index_writer; mod messages_reader; mod messages_writer; +use crate::fatal::ExitOnDescriptorExhaustion; use iggy_common::IggyError; use std::path::Path; use std::rc::Rc; @@ -60,6 +61,7 @@ impl SegmentStorage { .truncate(false) .open(messages_path) .await + .exit_on_descriptor_exhaustion(|| format!("opening {messages_path}")) .map_err(|_| IggyError::CannotCreateSegmentLogFile(messages_path.to_owned()))?; let mut changed = !file_exists; if let Some(size) = preallocate_size diff --git a/core/system_stats/src/lib.rs b/core/system_stats/src/lib.rs index eb6fd4fa6..5f96ddff0 100644 --- a/core/system_stats/src/lib.rs +++ b/core/system_stats/src/lib.rs @@ -54,6 +54,10 @@ pub struct SystemProbe { pub read_bytes: u64, pub written_bytes: u64, pub threads_count: u32, + /// Descriptors the process holds open. `None` where sysinfo cannot + /// count them. On Linux the count includes the descriptor of the + /// `/proc/self/fd` scan itself. + pub open_files: Option<u64>, } impl SystemProbe { @@ -82,6 +86,7 @@ impl SystemProbe { read_bytes: 0, written_bytes: 0, threads_count: 0, + open_files: None, }; if let Some(process) = sys.process(pid) { @@ -95,6 +100,9 @@ impl SystemProbe { probe.threads_count = process .tasks() .map_or(0, |tasks| u32::try_from(tasks.len()).unwrap_or(u32::MAX)); + probe.open_files = process + .open_files() + .map(|count| u64::try_from(count).unwrap_or(u64::MAX)); if let Some(memory) = cgroup_scoped_memory(sys, process) { probe.total_memory = memory.total; @@ -196,6 +204,9 @@ mod tests { assert!(probe.memory_usage > 0); assert!(probe.total_memory > 0); assert!(probe.available_memory <= probe.total_memory); + // stdin, stdout and stderr alone make the count nonzero. + #[cfg(any(target_os = "linux", target_os = "macos"))] + assert!(probe.open_files.is_some_and(|count| count > 0)); } #[test] diff --git a/foreign/go/errors/errors.yaml b/foreign/go/errors/errors.yaml index 145cd569f..9d411d6c8 100644 --- a/foreign/go/errors/errors.yaml +++ b/foreign/go/errors/errors.yaml @@ -651,6 +651,10 @@ code: 2021 format: "too many topics" fields: [] +- name: PartitionsLimitReached + code: 2022 + format: "partitions limit reached, raise [metadata] partitions_max" + fields: [] - name: CannotCreatePartition code: 3000 format: "cannot create partition with id: %d for stream with id: %d and topic with id: %d" diff --git a/foreign/go/errors/errors_gen.go b/foreign/go/errors/errors_gen.go index 2f80a8a78..d97d3e420 100644 --- a/foreign/go/errors/errors_gen.go +++ b/foreign/go/errors/errors_gen.go @@ -1358,6 +1358,17 @@ func (e TooManyTopics) Is(target error) bool { return ok } +type PartitionsLimitReached struct{} + +func (e PartitionsLimitReached) Error() string { + return "partitions limit reached, raise [metadata] partitions_max" +} +func (e PartitionsLimitReached) Code() Code { return 2022 } +func (e PartitionsLimitReached) Is(target error) bool { + _, ok := target.(PartitionsLimitReached) + return ok +} + type CannotCreatePartition struct { PartitionId uint32 StreamId uint32 @@ -2778,6 +2789,7 @@ var ( ErrInvalidPartitionsCount = InvalidPartitionsCount{} ErrTopicDirectoryNotFound = TopicDirectoryNotFound{} ErrTooManyTopics = TooManyTopics{} + ErrPartitionsLimitReached = PartitionsLimitReached{} ErrCannotCreatePartition = CannotCreatePartition{} ErrCannotCreatePartitionsDirectory = CannotCreatePartitionsDirectory{} ErrCannotCreatePartitionDirectory = CannotCreatePartitionDirectory{} @@ -3022,6 +3034,7 @@ const ( InvalidPartitionsCountCode Code = 2019 TopicDirectoryNotFoundCode Code = 2020 TooManyTopicsCode Code = 2021 + PartitionsLimitReachedCode Code = 2022 CannotCreatePartitionCode Code = 3000 CannotCreatePartitionsDirectoryCode Code = 3001 CannotCreatePartitionDirectoryCode Code = 3002 @@ -3388,6 +3401,8 @@ func (c Code) String() string { return "TopicDirectoryNotFound" case TooManyTopicsCode: return "TooManyTopics" + case PartitionsLimitReachedCode: + return "PartitionsLimitReached" case CannotCreatePartitionCode: return "CannotCreatePartition" case CannotCreatePartitionsDirectoryCode: @@ -3873,6 +3888,8 @@ func FromCode(code Code) IggyError { return ErrTopicDirectoryNotFound case TooManyTopicsCode: return ErrTooManyTopics + case PartitionsLimitReachedCode: + return ErrPartitionsLimitReached case CannotCreatePartitionCode: return ErrCannotCreatePartition case CannotCreatePartitionsDirectoryCode: diff --git a/foreign/java/java-sdk/src/main/java/org/apache/iggy/exception/IggyErrorCode.java b/foreign/java/java-sdk/src/main/java/org/apache/iggy/exception/IggyErrorCode.java index 5d776f827..bbdb57ef3 100644 --- a/foreign/java/java-sdk/src/main/java/org/apache/iggy/exception/IggyErrorCode.java +++ b/foreign/java/java-sdk/src/main/java/org/apache/iggy/exception/IggyErrorCode.java @@ -89,6 +89,7 @@ public enum IggyErrorCode { INVALID_TOPIC_ID(2016), INVALID_REPLICATION_FACTOR(2018), TOO_MANY_TOPICS(2021), + PARTITIONS_LIMIT_REACHED(2022), // Partition errors PARTITION_NOT_FOUND(3007), diff --git a/foreign/java/java-sdk/src/test/java/org/apache/iggy/exception/IggyErrorCodeTest.java b/foreign/java/java-sdk/src/test/java/org/apache/iggy/exception/IggyErrorCodeTest.java index 45feef08b..c76e23f8c 100644 --- a/foreign/java/java-sdk/src/test/java/org/apache/iggy/exception/IggyErrorCodeTest.java +++ b/foreign/java/java-sdk/src/test/java/org/apache/iggy/exception/IggyErrorCodeTest.java @@ -211,6 +211,7 @@ class IggyErrorCodeTest { "2016, INVALID_TOPIC_ID", "2018, INVALID_REPLICATION_FACTOR", "2021, TOO_MANY_TOPICS", + "2022, PARTITIONS_LIMIT_REACHED", // Partition errors "3007, PARTITION_NOT_FOUND", diff --git a/foreign/node/src/wire/error.code.test.ts b/foreign/node/src/wire/error.code.test.ts index dc768e2d2..b495eff2f 100644 --- a/foreign/node/src/wire/error.code.test.ts +++ b/foreign/node/src/wire/error.code.test.ts @@ -19,6 +19,10 @@ import assert from 'node:assert/strict'; import { it } from 'node:test'; import { translateErrorCode } from './error.code.js'; +it('translates the partitions capacity error', () => { + assert.equal(translateErrorCode(2022), 'Partitions limit reached, raise [metadata] partitions_max'); +}); + it('translates the consumer-offset capacity error', () => { assert.equal(translateErrorCode(3024), 'Consumer offset limit reached for partition, raise [partition] consumer_offsets_max'); }); diff --git a/foreign/node/src/wire/error.code.ts b/foreign/node/src/wire/error.code.ts index 58773d957..f78a8e610 100644 --- a/foreign/node/src/wire/error.code.ts +++ b/foreign/node/src/wire/error.code.ts @@ -146,6 +146,7 @@ export const translateErrorCode = (code: number): string => { case '2019': return "Invalid partitions count"; case '2020': return "Topic directory: {0} not found"; case '2021': return "Too many topics"; + case '2022': return "Partitions limit reached, raise [metadata] partitions_max"; // TOPIC case '3000': return "Cannot create partition with ID: {0} for stream with ID: {1} and topic with ID: {2}"; diff --git a/foreign/swift/Sources/Iggy/Errors/IggyErrorCode.swift b/foreign/swift/Sources/Iggy/Errors/IggyErrorCode.swift index 1b6940d1f..a82413e19 100644 --- a/foreign/swift/Sources/Iggy/Errors/IggyErrorCode.swift +++ b/foreign/swift/Sources/Iggy/Errors/IggyErrorCode.swift @@ -147,6 +147,7 @@ public enum IggyErrorCode: UInt32, Sendable, Hashable, CaseIterable, Codable { case invalidPartitionsCount = 2019 case topicDirectoryNotFound = 2020 case tooManyTopics = 2021 + case partitionsLimitReached = 2022 case cannotCreatePartition = 3000 case cannotCreatePartitionsDirectory = 3001 case cannotCreatePartitionDirectory = 3002 @@ -393,6 +394,7 @@ extension IggyErrorCode { case .invalidPartitionsCount: "invalid_partitions_count" case .topicDirectoryNotFound: "topic_directory_not_found" case .tooManyTopics: "too_many_topics" + case .partitionsLimitReached: "partitions_limit_reached" case .cannotCreatePartition: "cannot_create_partition" case .cannotCreatePartitionsDirectory: "cannot_create_partitions_directory" case .cannotCreatePartitionDirectory: "cannot_create_partition_directory"
