github-actions[bot] commented on code in PR #4311: URL: https://github.com/apache/iggy/pull/4311#discussion_r4127107200
########## core/partitions/src/io.rs: ########## @@ -0,0 +1,1382 @@ +// 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. + +use std::cell::{Cell, RefCell}; +use std::collections::HashSet; +use std::io::IoSlice; +use std::path::Path; +use std::rc::Rc; +use std::task::{Poll, Waker}; +use std::time::Duration; + +use compio::fs::File; +use consensus::VsrState; +use futures::{StreamExt, TryStreamExt}; +use iggy_binary_protocol::PrepareHeader; +use iggy_common::MAX_MESSAGE_SIZE_UPPER_BYTES; +use iggy_common::{IggyByteSize, IggyError}; +use journal::durable_storage::{DiskStorage, DurableStorage}; +use journal::local_gate::OwnedLocalGateGuard; +use journal::superblock::{PingPongSuperblock, SuperblockStore}; +use server_common::SegmentStorage; +use server_common::iobuf::{Frozen, IOV_MAX}; +use server_common::poll::PollHistoryId; +use server_common::send_messages::COMMAND_HEADER_SIZE; +use server_common::sharding::IggyNamespace; +use tracing::warn; + +use crate::iggy_index::IGGY_INDEX_SIZE; +use crate::offset_storage::{ + OffsetFilePermit, PersistedOffset, delete_persisted_offset, + delete_persisted_offset_with_storage, persist_offset, persist_offset_retained, + persist_purge_generation_with_storage, read_offset_max, +}; +use crate::{IggyIndexWriter, MessagesWriter, Segment}; + +// Slot cells, Rc headers, the executor task header and completion token storage. +const IO_CONTROL_ALLOCATION_RESERVE: usize = 4096; +// Source/destination siblings, File/OpenOptions paths and driver C strings can coexist. +const FILE_PATH_COPIES_MAX: usize = 8; +// glibc __alloc_dir bounds its filesystem-sized readdir buffer at 1 MiB. +// The control reserve separately covers DIR metadata and Rust iterator ownership. +const DIRECTORY_ITERATION_SCRATCH_MAX: usize = 1024 * 1024; +const PARTITION_IO_DRAIN_TIMEOUT: Duration = Duration::from_secs(30); + +/// Process-unique identity of one partition owner, preserved across its views. +/// A replacement in the same namespace always has a different identity. +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, PartialOrd, Ord)] +pub struct PartitionIncarnation(PollHistoryId); + +/// Captured file work. Execution never accesses installed partition state or +/// advances writer cursors; only the owner can accept the returned result. +pub struct MaterializationIoJob { + pub(crate) allocation_charge: usize, + pub(crate) target: MaterializationTarget, + pub(crate) batches: Vec<Frozen<4096>>, + pub(crate) indexes: Vec<u8>, + pub(crate) bodies_written: Option<u64>, + pub(crate) written: u64, + pub(crate) completes: bool, +} + +pub struct MaterializationIoResult { + pub(crate) target: MaterializationTarget, + pub(crate) outcome: Result<(u64, u64), IggyError>, +} + +pub enum PartitionIoJob<SB = PingPongSuperblock> { + Materialize(MaterializationIoJob), + OffsetWrite(OffsetIoJob), + OffsetDelete(OffsetDeleteIoJob), + OffsetDirectories(OffsetDirectoriesIoJob), + SegmentDirectory(String), + IndexSync(Rc<IggyIndexWriter>), + RemoveSegment { + namespace: IggyNamespace, + paths: [Option<String>; 3], + strict: bool, + }, + EmptySegment(SegmentIoJob), + PurgeCleanup { + directory: String, + bodies: bool, + }, + PurgeOffsets { + directory: String, + known: HashSet<u32>, + /// Sorted IDs preserved by a state-transfer install. + retained: Vec<u32>, + }, + PurgeGeneration { + path: String, + generation: u64, + revision: u64, + }, + Quarantine { + directory: String, + revision: u64, + replicated: bool, + }, + Superblock(SuperblockIoJob<SB>), + Transfer(TransferFileJob), + Rotate { + target: RotationTarget, + job: SegmentIoJob, + }, +} + +pub enum PartitionIoResult { + Materialize(MaterializationIoResult), + OffsetWrite(OffsetIoResult), + OffsetDelete(Result<bool, IggyError>), + OffsetDirectories(OffsetDirectoriesIoResult), + SegmentDirectory(std::io::Result<()>), + IndexSync { + writer: Rc<IggyIndexWriter>, + outcome: Result<(), IggyError>, + }, + SegmentRemoved(Result<(), IggyError>), + EmptySegment(Result<InstalledSegment, IggyError>), + PurgeCleanup(std::io::Result<()>), + PurgeOffsets(PurgeOffsetsIoResult), + PurgeGeneration(Result<(), IggyError>), + Quarantine(std::io::Result<String>), + Superblock(SuperblockIoResult), + Transfer(TransferFileResult), + Rotate { + target: RotationTarget, + outcome: Result<InstalledSegment, IggyError>, + }, +} + +/// Identifies the owner continuation that may accept a completed phase. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum PartitionIoContinuation { + Install, + Commit, + NoAck, + Materialization, + Superblock, + Checkpoint, + Retention, + Purge, + Quarantine, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct PartitionIoIdentity { + pub namespace: IggyNamespace, + pub incarnation: PartitionIncarnation, + pub history: PollHistoryId, + pub sequence: u64, + pub continuation: PartitionIoContinuation, + pub local_order: Option<consensus::LocalRequestOrder>, +} + +pub struct CapturedPartitionIo<SB> { + pub identity: PartitionIoIdentity, + pub job: PartitionIoJob<SB>, + pub gate: Option<OwnedLocalGateGuard>, + pub quiescence: Rc<PartitionIoQuiescence>, +} + +/// Captured before tombstoning, so teardown never depends on mounted lookup. +#[derive(Default)] +pub struct PartitionIoQuiescence { + active: Cell<Option<PartitionIoIdentity>>, + retiring: Cell<bool>, + deleting: Cell<bool>, + interrupted: Cell<bool>, + waiter: RefCell<Option<Waker>>, +} + +pub struct PartitionTeardown { + pub(crate) io: Rc<PartitionIoQuiescence>, + pub(crate) persistence: Option<Rc<crate::PartitionPersistence>>, +} + +impl PartitionIoQuiescence { + /// Called only after the matching file future returned and its result was consumed. + /// A removed owner still needs this settlement before its files can be deleted. + pub fn settle(&self, identity: PartitionIoIdentity) { + if !self.interrupted.get() && self.active.get() == Some(identity) { + self.set(None); + } + } + + pub(crate) const fn get(&self) -> Option<PartitionIoIdentity> { + self.active.get() + } + + pub(crate) fn set(&self, active: Option<PartitionIoIdentity>) { + self.active.set(active); + if active.is_none() + && let Some(waiter) = self.waiter.borrow_mut().take() + { + waiter.wake(); + } + } + + pub(crate) fn retire(&self) { + self.retiring.set(true); + } + + pub(crate) const fn is_retiring(&self) -> bool { + self.retiring.get() + } + + pub(crate) fn delete(&self) { + self.retiring.set(true); + self.deleting.set(true); + } + + pub(crate) const fn is_deleting(&self) -> bool { + self.deleting.get() + } + + pub(crate) const fn is_interrupted(&self) -> bool { + self.interrupted.get() + } + + /// Safe in a dropped task: no allocation, callback or resource release. + pub fn interrupt(&self) { Review Comment: warning: `interrupt()` sets the flag but wakes no waiter, and `settle()` refuses to clear `active` once interrupted, so a `drain()` parked earlier waits out the 30 s timeout instead of returning the interrupted-writer error. Wake the parked waiter when the interrupt fires, for example from the lane's interrupted scan, which already holds the slot's quiescence handle. ########## core/shard/src/partition_io.rs: ########## @@ -0,0 +1,1937 @@ +// 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. + +use std::cell::{Cell, RefCell}; +use std::collections::{BTreeSet, VecDeque}; +use std::rc::Rc; +use std::task::Waker; + +use consensus::PartitionsHandle; +use crossfire::{RecvError, TryRecvError}; +use journal::local_gate::OwnedLocalGateGuard; +use journal::superblock::SuperblockStore; +use message_bus::MessageBus; +use partitions::{ + CapturedPartitionIo, PartitionIncarnation, PartitionIoIdentity, PartitionIoResources, + PartitionIoResult, +}; +use server_common::sharding::IggyNamespace; +use thiserror::Error; + +use crate::{IggyShard, Receiver, Sender, channel}; + +pub const DEFAULT_PARTITION_IO_CAPACITY: usize = 16; +pub const DEFAULT_PARTITION_IO_BYTES: usize = 256 * 1024 * 1024; +pub const PARTITION_IO_CAPACITY_MAX: usize = 1 << 20; + +#[derive(Clone, Copy, Debug)] +pub struct PartitionIoLimits { + capacity: usize, + bytes_max: usize, +} + +#[derive(Debug, Error)] +pub enum PartitionIoLimitsError { + #[error("sharding.partition_io_capacity must be in 1..={PARTITION_IO_CAPACITY_MAX}; got {0}")] + Capacity(usize), + #[error("partition I/O allocation charge exceeds addressable memory")] + Overflow, + #[error( + "sharding.partition_io_bytes_max must be at least {minimum} and fit addressable memory; got {value}" + )] + Bytes { value: usize, minimum: usize }, +} + +impl PartitionIoLimits { + /// Resolve omitted bytes using the same allocation calculation as dispatch. + /// + /// # Errors + /// Rejects invalid slot counts, arithmetic overflow and undersized byte limits. + pub fn new(capacity: usize, bytes_max: Option<usize>) -> Result<Self, PartitionIoLimitsError> { + if capacity == 0 || capacity > PARTITION_IO_CAPACITY_MAX { + return Err(PartitionIoLimitsError::Capacity(capacity)); + } + let minimum = partitions::largest_legal_job_charge() + .filter(|charge| isize::try_from(*charge).is_ok()) + .ok_or(PartitionIoLimitsError::Overflow)?; + let bytes_max = bytes_max.unwrap_or_else(|| DEFAULT_PARTITION_IO_BYTES.max(minimum)); Review Comment: warning: `DEFAULT_PARTITION_IO_BYTES` never wins this `max`, because one job's charge starts at `u32::MAX`, so the resolved budget is just over 4 GiB while `config.toml` and the `ShardingConfig` field doc advertise 256 MiB. Delete the constant and state the real resolved minimum where operators set `partition_io_bytes_max`. ########## core/shard/src/partition_io.rs: ########## @@ -0,0 +1,1937 @@ +// 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. + +use std::cell::{Cell, RefCell}; +use std::collections::{BTreeSet, VecDeque}; +use std::rc::Rc; +use std::task::Waker; + +use consensus::PartitionsHandle; +use crossfire::{RecvError, TryRecvError}; +use journal::local_gate::OwnedLocalGateGuard; +use journal::superblock::SuperblockStore; +use message_bus::MessageBus; +use partitions::{ + CapturedPartitionIo, PartitionIncarnation, PartitionIoIdentity, PartitionIoResources, + PartitionIoResult, +}; +use server_common::sharding::IggyNamespace; +use thiserror::Error; + +use crate::{IggyShard, Receiver, Sender, channel}; + +pub const DEFAULT_PARTITION_IO_CAPACITY: usize = 16; +pub const DEFAULT_PARTITION_IO_BYTES: usize = 256 * 1024 * 1024; +pub const PARTITION_IO_CAPACITY_MAX: usize = 1 << 20; + +#[derive(Clone, Copy, Debug)] +pub struct PartitionIoLimits { + capacity: usize, + bytes_max: usize, +} + +#[derive(Debug, Error)] +pub enum PartitionIoLimitsError { + #[error("sharding.partition_io_capacity must be in 1..={PARTITION_IO_CAPACITY_MAX}; got {0}")] + Capacity(usize), + #[error("partition I/O allocation charge exceeds addressable memory")] + Overflow, + #[error( + "sharding.partition_io_bytes_max must be at least {minimum} and fit addressable memory; got {value}" + )] + Bytes { value: usize, minimum: usize }, +} + +impl PartitionIoLimits { + /// Resolve omitted bytes using the same allocation calculation as dispatch. + /// + /// # Errors + /// Rejects invalid slot counts, arithmetic overflow and undersized byte limits. + pub fn new(capacity: usize, bytes_max: Option<usize>) -> Result<Self, PartitionIoLimitsError> { + if capacity == 0 || capacity > PARTITION_IO_CAPACITY_MAX { + return Err(PartitionIoLimitsError::Capacity(capacity)); + } + let minimum = partitions::largest_legal_job_charge() + .filter(|charge| isize::try_from(*charge).is_ok()) + .ok_or(PartitionIoLimitsError::Overflow)?; + let bytes_max = bytes_max.unwrap_or_else(|| DEFAULT_PARTITION_IO_BYTES.max(minimum)); + if bytes_max < minimum || bytes_max > isize::MAX as usize { + return Err(PartitionIoLimitsError::Bytes { + value: bytes_max, + minimum, + }); + } + Ok(Self { + capacity, + bytes_max, + }) + } + + #[must_use] + pub const fn capacity(self) -> usize { + self.capacity + } + + #[must_use] + pub const fn bytes_max(self) -> usize { + self.bytes_max + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct PartitionIoToken { + slot: usize, + identity: PartitionIoIdentity, +} + +#[derive(Clone, Copy, PartialEq, Eq)] +enum SlotState { + Reserved, + Running, + Queued, + Settled, + Interrupted, +} + +struct PartitionIoSlot<SB> { + namespace: IggyNamespace, + incarnation: PartitionIncarnation, + identity: Cell<Option<PartitionIoIdentity>>, + state: Cell<SlotState>, + charge: usize, + result: RefCell<Option<PartitionIoResult>>, + resources: RefCell<Option<PartitionIoResources<SB>>>, + gate: RefCell<Option<OwnedLocalGateGuard>>, + quiescence: RefCell<Option<Rc<partitions::PartitionIoQuiescence>>>, +} + +#[derive(Default)] +struct ReadyPartitions { + queue: RefCell<VecDeque<(IggyNamespace, PartitionIncarnation)>>, + present: RefCell<BTreeSet<(IggyNamespace, PartitionIncarnation)>>, + waker: RefCell<Option<Waker>>, +} + +#[derive(Default)] +struct IoCounters { + active: Cell<usize>, + queued: Cell<usize>, + quarantined: Cell<usize>, +} + +impl ReadyPartitions { + fn notify(&self, namespace: IggyNamespace, incarnation: PartitionIncarnation) { + if !self.present.borrow_mut().insert((namespace, incarnation)) { + return; + } + self.queue.borrow_mut().push_back((namespace, incarnation)); + let waker = self.waker.borrow().clone(); + if let Some(waker) = waker { + waker.wake(); + } + } +} + +/// Slots outlive mounted lookup. Tokens never carry local file handles or results. +pub struct PartitionIoLane<SB> { + pub(crate) limits: PartitionIoLimits, + slots: RefCell<Vec<Option<Rc<PartitionIoSlot<SB>>>>>, + charged: Cell<usize>, + sender: Sender<PartitionIoToken>, + receiver: Receiver<PartitionIoToken>, + ready: Rc<ReadyPartitions>, + interrupted: Rc<Cell<bool>>, + undelivered: Rc<Cell<bool>>, + closed: Cell<bool>, + capacity_blocked: Cell<bool>, + counters: Rc<IoCounters>, + #[cfg(test)] + execution_gate: RefCell<Option<futures::channel::oneshot::Receiver<()>>>, +} + +impl<SB: SuperblockStore> PartitionIoLane<SB> { + pub(crate) fn new(limits: PartitionIoLimits) -> Self { + let (sender, receiver) = channel(limits.capacity); + Self { + limits, + slots: RefCell::new((0..limits.capacity).map(|_| None).collect()), + charged: Cell::new(0), + sender, + receiver, + ready: Rc::default(), + interrupted: Rc::new(Cell::new(false)), + undelivered: Rc::new(Cell::new(false)), + closed: Cell::new(false), + capacity_blocked: Cell::new(false), + counters: Rc::default(), + #[cfg(test)] + execution_gate: RefCell::new(None), + } + } + + pub(crate) fn notifier(&self) -> partitions::PartitionIoNotifier { + let ready = Rc::clone(&self.ready); + Rc::new(move |namespace, incarnation| ready.notify(namespace, incarnation)) + } + + pub(crate) fn register_waker(&self, waker: &Waker) { + let mut current = self.ready.waker.borrow_mut(); + if current + .as_ref() + .is_none_or(|current| !current.will_wake(waker)) + { + *current = Some(waker.clone()); + } + } + + pub(crate) fn has_ready(&self) -> bool { + self.head().is_some() || self.undelivered.get() || self.interrupted.get() + } + + pub(crate) fn head(&self) -> Option<(IggyNamespace, PartitionIncarnation)> { + if !self.capacity_blocked.get() { + return self.ready.queue.borrow().front().copied(); + } + let present = self.ready.present.borrow(); + self.slots.borrow().iter().flatten().find_map(|slot| { + (slot.state.get() == SlotState::Settled + && present.contains(&(slot.namespace, slot.incarnation))) + .then_some((slot.namespace, slot.incarnation)) + }) + } + + pub(crate) fn pop_ready(&self, namespace: IggyNamespace, incarnation: PartitionIncarnation) { Review Comment: warning: `pop_ready` rebuilds the ready queue with `VecDeque::retain`, so every pop scans the whole queue even though the pair just came from `head()`. Pop the front entry when it matches and keep the scan for the capacity-blocked case only. ########## core/configs/src/server_config/sharding.rs: ########## @@ -101,6 +101,11 @@ pub struct ShardingConfig { /// rejects new disk polls before I/O. Main and reply inbox traffic uses /// separate capacity. This counts operations, not retained message bytes. pub poll_completion_capacity: usize, + /// Active partition file jobs and results awaiting owner acceptance. + pub partition_io_capacity: usize, + /// Retained job allocations, resolved at boot against the largest legal record. + /// Omission selects max(256 MiB, the single-job minimum). + pub partition_io_bytes_max: Option<usize>, Review Comment: nit: `partition_io_bytes_max` is the only byte-valued server key typed as a raw integer; `wal_bytes_max`, `evicted_ring_bytes_max` and `transfer_artifact_bytes_max` all use `IggyByteSize`. Type this key as `Option<​IggyByteSize>` and convert at boot. ########## core/shard/src/partition_io.rs: ########## @@ -0,0 +1,1937 @@ +// 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. + +use std::cell::{Cell, RefCell}; +use std::collections::{BTreeSet, VecDeque}; +use std::rc::Rc; +use std::task::Waker; + +use consensus::PartitionsHandle; +use crossfire::{RecvError, TryRecvError}; +use journal::local_gate::OwnedLocalGateGuard; +use journal::superblock::SuperblockStore; +use message_bus::MessageBus; +use partitions::{ + CapturedPartitionIo, PartitionIncarnation, PartitionIoIdentity, PartitionIoResources, + PartitionIoResult, +}; +use server_common::sharding::IggyNamespace; +use thiserror::Error; + +use crate::{IggyShard, Receiver, Sender, channel}; + +pub const DEFAULT_PARTITION_IO_CAPACITY: usize = 16; +pub const DEFAULT_PARTITION_IO_BYTES: usize = 256 * 1024 * 1024; +pub const PARTITION_IO_CAPACITY_MAX: usize = 1 << 20; + +#[derive(Clone, Copy, Debug)] +pub struct PartitionIoLimits { + capacity: usize, + bytes_max: usize, +} + +#[derive(Debug, Error)] +pub enum PartitionIoLimitsError { + #[error("sharding.partition_io_capacity must be in 1..={PARTITION_IO_CAPACITY_MAX}; got {0}")] + Capacity(usize), + #[error("partition I/O allocation charge exceeds addressable memory")] + Overflow, + #[error( + "sharding.partition_io_bytes_max must be at least {minimum} and fit addressable memory; got {value}" + )] + Bytes { value: usize, minimum: usize }, +} + +impl PartitionIoLimits { + /// Resolve omitted bytes using the same allocation calculation as dispatch. + /// + /// # Errors + /// Rejects invalid slot counts, arithmetic overflow and undersized byte limits. + pub fn new(capacity: usize, bytes_max: Option<usize>) -> Result<Self, PartitionIoLimitsError> { + if capacity == 0 || capacity > PARTITION_IO_CAPACITY_MAX { + return Err(PartitionIoLimitsError::Capacity(capacity)); + } + let minimum = partitions::largest_legal_job_charge() + .filter(|charge| isize::try_from(*charge).is_ok()) + .ok_or(PartitionIoLimitsError::Overflow)?; + let bytes_max = bytes_max.unwrap_or_else(|| DEFAULT_PARTITION_IO_BYTES.max(minimum)); + if bytes_max < minimum || bytes_max > isize::MAX as usize { + return Err(PartitionIoLimitsError::Bytes { + value: bytes_max, + minimum, + }); + } + Ok(Self { + capacity, + bytes_max, + }) + } + + #[must_use] + pub const fn capacity(self) -> usize { Review Comment: simplification: `capacity()` has no caller in the repository; `PartitionIoLane::new` reads the private field directly. Delete the accessor. ########## core/shard/src/partition_io.rs: ########## @@ -0,0 +1,1937 @@ +// 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. + +use std::cell::{Cell, RefCell}; +use std::collections::{BTreeSet, VecDeque}; +use std::rc::Rc; +use std::task::Waker; + +use consensus::PartitionsHandle; +use crossfire::{RecvError, TryRecvError}; +use journal::local_gate::OwnedLocalGateGuard; +use journal::superblock::SuperblockStore; +use message_bus::MessageBus; +use partitions::{ + CapturedPartitionIo, PartitionIncarnation, PartitionIoIdentity, PartitionIoResources, + PartitionIoResult, +}; +use server_common::sharding::IggyNamespace; +use thiserror::Error; + +use crate::{IggyShard, Receiver, Sender, channel}; + +pub const DEFAULT_PARTITION_IO_CAPACITY: usize = 16; +pub const DEFAULT_PARTITION_IO_BYTES: usize = 256 * 1024 * 1024; +pub const PARTITION_IO_CAPACITY_MAX: usize = 1 << 20; + +#[derive(Clone, Copy, Debug)] +pub struct PartitionIoLimits { + capacity: usize, + bytes_max: usize, +} + +#[derive(Debug, Error)] +pub enum PartitionIoLimitsError { + #[error("sharding.partition_io_capacity must be in 1..={PARTITION_IO_CAPACITY_MAX}; got {0}")] + Capacity(usize), + #[error("partition I/O allocation charge exceeds addressable memory")] + Overflow, + #[error( + "sharding.partition_io_bytes_max must be at least {minimum} and fit addressable memory; got {value}" + )] + Bytes { value: usize, minimum: usize }, +} + +impl PartitionIoLimits { + /// Resolve omitted bytes using the same allocation calculation as dispatch. + /// + /// # Errors + /// Rejects invalid slot counts, arithmetic overflow and undersized byte limits. + pub fn new(capacity: usize, bytes_max: Option<usize>) -> Result<Self, PartitionIoLimitsError> { + if capacity == 0 || capacity > PARTITION_IO_CAPACITY_MAX { + return Err(PartitionIoLimitsError::Capacity(capacity)); + } + let minimum = partitions::largest_legal_job_charge() + .filter(|charge| isize::try_from(*charge).is_ok()) + .ok_or(PartitionIoLimitsError::Overflow)?; + let bytes_max = bytes_max.unwrap_or_else(|| DEFAULT_PARTITION_IO_BYTES.max(minimum)); + if bytes_max < minimum || bytes_max > isize::MAX as usize { + return Err(PartitionIoLimitsError::Bytes { + value: bytes_max, + minimum, + }); + } + Ok(Self { + capacity, + bytes_max, + }) + } + + #[must_use] + pub const fn capacity(self) -> usize { + self.capacity + } + + #[must_use] + pub const fn bytes_max(self) -> usize { + self.bytes_max + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct PartitionIoToken { + slot: usize, + identity: PartitionIoIdentity, +} + +#[derive(Clone, Copy, PartialEq, Eq)] +enum SlotState { + Reserved, + Running, + Queued, + Settled, + Interrupted, +} + +struct PartitionIoSlot<SB> { + namespace: IggyNamespace, + incarnation: PartitionIncarnation, + identity: Cell<Option<PartitionIoIdentity>>, + state: Cell<SlotState>, + charge: usize, + result: RefCell<Option<PartitionIoResult>>, + resources: RefCell<Option<PartitionIoResources<SB>>>, + gate: RefCell<Option<OwnedLocalGateGuard>>, + quiescence: RefCell<Option<Rc<partitions::PartitionIoQuiescence>>>, +} + +#[derive(Default)] +struct ReadyPartitions { + queue: RefCell<VecDeque<(IggyNamespace, PartitionIncarnation)>>, + present: RefCell<BTreeSet<(IggyNamespace, PartitionIncarnation)>>, + waker: RefCell<Option<Waker>>, +} + +#[derive(Default)] +struct IoCounters { + active: Cell<usize>, + queued: Cell<usize>, + quarantined: Cell<usize>, +} + +impl ReadyPartitions { + fn notify(&self, namespace: IggyNamespace, incarnation: PartitionIncarnation) { + if !self.present.borrow_mut().insert((namespace, incarnation)) { + return; + } + self.queue.borrow_mut().push_back((namespace, incarnation)); + let waker = self.waker.borrow().clone(); + if let Some(waker) = waker { + waker.wake(); + } + } +} + +/// Slots outlive mounted lookup. Tokens never carry local file handles or results. +pub struct PartitionIoLane<SB> { + pub(crate) limits: PartitionIoLimits, + slots: RefCell<Vec<Option<Rc<PartitionIoSlot<SB>>>>>, + charged: Cell<usize>, + sender: Sender<PartitionIoToken>, + receiver: Receiver<PartitionIoToken>, + ready: Rc<ReadyPartitions>, + interrupted: Rc<Cell<bool>>, + undelivered: Rc<Cell<bool>>, + closed: Cell<bool>, + capacity_blocked: Cell<bool>, + counters: Rc<IoCounters>, + #[cfg(test)] + execution_gate: RefCell<Option<futures::channel::oneshot::Receiver<()>>>, +} + +impl<SB: SuperblockStore> PartitionIoLane<SB> { + pub(crate) fn new(limits: PartitionIoLimits) -> Self { + let (sender, receiver) = channel(limits.capacity); + Self { + limits, + slots: RefCell::new((0..limits.capacity).map(|_| None).collect()), + charged: Cell::new(0), + sender, + receiver, + ready: Rc::default(), + interrupted: Rc::new(Cell::new(false)), + undelivered: Rc::new(Cell::new(false)), + closed: Cell::new(false), + capacity_blocked: Cell::new(false), + counters: Rc::default(), + #[cfg(test)] + execution_gate: RefCell::new(None), + } + } + + pub(crate) fn notifier(&self) -> partitions::PartitionIoNotifier { + let ready = Rc::clone(&self.ready); + Rc::new(move |namespace, incarnation| ready.notify(namespace, incarnation)) + } + + pub(crate) fn register_waker(&self, waker: &Waker) { + let mut current = self.ready.waker.borrow_mut(); + if current + .as_ref() + .is_none_or(|current| !current.will_wake(waker)) + { + *current = Some(waker.clone()); + } + } + + pub(crate) fn has_ready(&self) -> bool { + self.head().is_some() || self.undelivered.get() || self.interrupted.get() + } + + pub(crate) fn head(&self) -> Option<(IggyNamespace, PartitionIncarnation)> { + if !self.capacity_blocked.get() { + return self.ready.queue.borrow().front().copied(); + } + let present = self.ready.present.borrow(); + self.slots.borrow().iter().flatten().find_map(|slot| { + (slot.state.get() == SlotState::Settled + && present.contains(&(slot.namespace, slot.incarnation))) + .then_some((slot.namespace, slot.incarnation)) + }) + } + + pub(crate) fn pop_ready(&self, namespace: IggyNamespace, incarnation: PartitionIncarnation) { + self.ready + .queue + .borrow_mut() + .retain(|queued| *queued != (namespace, incarnation)); + self.ready + .present + .borrow_mut() + .remove(&(namespace, incarnation)); + } + + pub(crate) fn reschedule(&self, namespace: IggyNamespace, incarnation: PartitionIncarnation) { + self.ready.notify(namespace, incarnation); + } + + pub(crate) fn try_reserve( + &self, + namespace: IggyNamespace, + incarnation: PartitionIncarnation, + charge: usize, + ) -> Option<usize> { + if self.closed.get() { + return None; + } + let mut slots = self.slots.borrow_mut(); + if let Some((index, existing)) = slots.iter().enumerate().find_map(|(index, slot)| { Review Comment: simplification: `try_reserve` and `retained` each scan `slots` for the same `(namespace, incarnation)` pair. Extract one `slot_of()` lookup and call it from both. ########## core/shard/src/partition_io.rs: ########## @@ -0,0 +1,1937 @@ +// 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. + +use std::cell::{Cell, RefCell}; +use std::collections::{BTreeSet, VecDeque}; +use std::rc::Rc; +use std::task::Waker; + +use consensus::PartitionsHandle; +use crossfire::{RecvError, TryRecvError}; +use journal::local_gate::OwnedLocalGateGuard; +use journal::superblock::SuperblockStore; +use message_bus::MessageBus; +use partitions::{ + CapturedPartitionIo, PartitionIncarnation, PartitionIoIdentity, PartitionIoResources, + PartitionIoResult, +}; +use server_common::sharding::IggyNamespace; +use thiserror::Error; + +use crate::{IggyShard, Receiver, Sender, channel}; + +pub const DEFAULT_PARTITION_IO_CAPACITY: usize = 16; +pub const DEFAULT_PARTITION_IO_BYTES: usize = 256 * 1024 * 1024; +pub const PARTITION_IO_CAPACITY_MAX: usize = 1 << 20; + +#[derive(Clone, Copy, Debug)] +pub struct PartitionIoLimits { + capacity: usize, + bytes_max: usize, +} + +#[derive(Debug, Error)] +pub enum PartitionIoLimitsError { + #[error("sharding.partition_io_capacity must be in 1..={PARTITION_IO_CAPACITY_MAX}; got {0}")] + Capacity(usize), + #[error("partition I/O allocation charge exceeds addressable memory")] + Overflow, + #[error( + "sharding.partition_io_bytes_max must be at least {minimum} and fit addressable memory; got {value}" + )] + Bytes { value: usize, minimum: usize }, +} + +impl PartitionIoLimits { + /// Resolve omitted bytes using the same allocation calculation as dispatch. + /// + /// # Errors + /// Rejects invalid slot counts, arithmetic overflow and undersized byte limits. + pub fn new(capacity: usize, bytes_max: Option<usize>) -> Result<Self, PartitionIoLimitsError> { + if capacity == 0 || capacity > PARTITION_IO_CAPACITY_MAX { + return Err(PartitionIoLimitsError::Capacity(capacity)); + } + let minimum = partitions::largest_legal_job_charge() + .filter(|charge| isize::try_from(*charge).is_ok()) + .ok_or(PartitionIoLimitsError::Overflow)?; + let bytes_max = bytes_max.unwrap_or_else(|| DEFAULT_PARTITION_IO_BYTES.max(minimum)); + if bytes_max < minimum || bytes_max > isize::MAX as usize { + return Err(PartitionIoLimitsError::Bytes { + value: bytes_max, + minimum, + }); + } + Ok(Self { + capacity, + bytes_max, + }) + } + + #[must_use] + pub const fn capacity(self) -> usize { + self.capacity + } + + #[must_use] + pub const fn bytes_max(self) -> usize { + self.bytes_max + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct PartitionIoToken { + slot: usize, + identity: PartitionIoIdentity, +} + +#[derive(Clone, Copy, PartialEq, Eq)] +enum SlotState { + Reserved, + Running, + Queued, + Settled, + Interrupted, +} + +struct PartitionIoSlot<SB> { + namespace: IggyNamespace, + incarnation: PartitionIncarnation, + identity: Cell<Option<PartitionIoIdentity>>, + state: Cell<SlotState>, + charge: usize, + result: RefCell<Option<PartitionIoResult>>, + resources: RefCell<Option<PartitionIoResources<SB>>>, + gate: RefCell<Option<OwnedLocalGateGuard>>, + quiescence: RefCell<Option<Rc<partitions::PartitionIoQuiescence>>>, +} + +#[derive(Default)] +struct ReadyPartitions { + queue: RefCell<VecDeque<(IggyNamespace, PartitionIncarnation)>>, + present: RefCell<BTreeSet<(IggyNamespace, PartitionIncarnation)>>, + waker: RefCell<Option<Waker>>, +} + +#[derive(Default)] +struct IoCounters { + active: Cell<usize>, + queued: Cell<usize>, + quarantined: Cell<usize>, +} + +impl ReadyPartitions { + fn notify(&self, namespace: IggyNamespace, incarnation: PartitionIncarnation) { + if !self.present.borrow_mut().insert((namespace, incarnation)) { + return; + } + self.queue.borrow_mut().push_back((namespace, incarnation)); + let waker = self.waker.borrow().clone(); + if let Some(waker) = waker { + waker.wake(); + } + } +} + +/// Slots outlive mounted lookup. Tokens never carry local file handles or results. +pub struct PartitionIoLane<SB> { + pub(crate) limits: PartitionIoLimits, + slots: RefCell<Vec<Option<Rc<PartitionIoSlot<SB>>>>>, + charged: Cell<usize>, + sender: Sender<PartitionIoToken>, + receiver: Receiver<PartitionIoToken>, + ready: Rc<ReadyPartitions>, + interrupted: Rc<Cell<bool>>, + undelivered: Rc<Cell<bool>>, + closed: Cell<bool>, + capacity_blocked: Cell<bool>, + counters: Rc<IoCounters>, + #[cfg(test)] + execution_gate: RefCell<Option<futures::channel::oneshot::Receiver<()>>>, +} + +impl<SB: SuperblockStore> PartitionIoLane<SB> { + pub(crate) fn new(limits: PartitionIoLimits) -> Self { + let (sender, receiver) = channel(limits.capacity); + Self { + limits, + slots: RefCell::new((0..limits.capacity).map(|_| None).collect()), + charged: Cell::new(0), + sender, + receiver, + ready: Rc::default(), + interrupted: Rc::new(Cell::new(false)), + undelivered: Rc::new(Cell::new(false)), + closed: Cell::new(false), + capacity_blocked: Cell::new(false), + counters: Rc::default(), + #[cfg(test)] + execution_gate: RefCell::new(None), + } + } + + pub(crate) fn notifier(&self) -> partitions::PartitionIoNotifier { + let ready = Rc::clone(&self.ready); + Rc::new(move |namespace, incarnation| ready.notify(namespace, incarnation)) + } + + pub(crate) fn register_waker(&self, waker: &Waker) { + let mut current = self.ready.waker.borrow_mut(); + if current + .as_ref() + .is_none_or(|current| !current.will_wake(waker)) + { + *current = Some(waker.clone()); + } + } + + pub(crate) fn has_ready(&self) -> bool { + self.head().is_some() || self.undelivered.get() || self.interrupted.get() + } + + pub(crate) fn head(&self) -> Option<(IggyNamespace, PartitionIncarnation)> { + if !self.capacity_blocked.get() { + return self.ready.queue.borrow().front().copied(); + } + let present = self.ready.present.borrow(); + self.slots.borrow().iter().flatten().find_map(|slot| { + (slot.state.get() == SlotState::Settled + && present.contains(&(slot.namespace, slot.incarnation))) + .then_some((slot.namespace, slot.incarnation)) + }) + } + + pub(crate) fn pop_ready(&self, namespace: IggyNamespace, incarnation: PartitionIncarnation) { + self.ready + .queue + .borrow_mut() + .retain(|queued| *queued != (namespace, incarnation)); + self.ready + .present + .borrow_mut() + .remove(&(namespace, incarnation)); + } + + pub(crate) fn reschedule(&self, namespace: IggyNamespace, incarnation: PartitionIncarnation) { + self.ready.notify(namespace, incarnation); + } + + pub(crate) fn try_reserve( + &self, + namespace: IggyNamespace, + incarnation: PartitionIncarnation, + charge: usize, + ) -> Option<usize> { + if self.closed.get() { + return None; + } + let mut slots = self.slots.borrow_mut(); + if let Some((index, existing)) = slots.iter().enumerate().find_map(|(index, slot)| { + slot.as_ref() + .filter(|slot| slot.namespace == namespace && slot.incarnation == incarnation) + .map(|slot| (index, slot)) + }) { + return (existing.state.get() == SlotState::Settled && charge <= existing.charge) + .then_some(index); + } + let charged = self.charged.get().checked_add(charge)?; + if charged > self.limits.bytes_max { + return None; + } + let index = slots.iter().position(Option::is_none)?; + slots[index] = Some(Rc::new(PartitionIoSlot { + namespace, + incarnation, + identity: Cell::new(None), + state: Cell::new(SlotState::Reserved), + charge, + result: RefCell::new(None), + resources: RefCell::new(None), + gate: RefCell::new(None), + quiescence: RefCell::new(None), + })); + self.charged.set(charged); + Some(index) + } + + pub(crate) fn dispatch( + &self, + index: usize, + captured: CapturedPartitionIo<SB>, + bus: &impl MessageBus, + ) where + SB: 'static, + { + let slot = Rc::clone( + self.slots.borrow()[index] + .as_ref() + .expect("reserved partition I/O slot"), + ); + let CapturedPartitionIo { + identity, + job, + gate, + quiescence, + } = captured; + slot.identity.set(Some(identity)); + *slot.resources.borrow_mut() = Some(job.retain_resources()); + *slot.gate.borrow_mut() = gate; + *slot.quiescence.borrow_mut() = Some(Rc::clone(&quiescence)); + slot.state.set(SlotState::Running); + self.counters.active.set(self.counters.active.get() + 1); + let marker = InterruptionMarker { + slot: Rc::clone(&slot), + interrupted: Rc::clone(&self.interrupted), + quiescence, + counters: Rc::clone(&self.counters), + }; + let sender = self.sender.clone(); + let undelivered = Rc::clone(&self.undelivered); + let ready = Rc::clone(&self.ready); + let counters = Rc::clone(&self.counters); + #[cfg(test)] + let execution_gate = self.execution_gate.borrow_mut().take(); + bus.spawn(async move { + #[cfg(test)] + if let Some(execution_gate) = execution_gate { + execution_gate + .await + .expect("test releases captured file job"); + } + let result = job.execute().await; + *slot.result.borrow_mut() = Some(result); + slot.state.set(SlotState::Queued); + counters.active.set(counters.active.get() - 1); + counters.queued.set(counters.queued.get() + 1); + if sender + .try_send(PartitionIoToken { + slot: index, + identity, + }) + .is_err() + { + undelivered.set(true); + if let Some(waker) = ready.waker.borrow().as_ref() { + waker.wake_by_ref(); + } + } + drop(marker); + }); + } + + #[allow(clippy::future_not_send)] + pub(crate) async fn recv(&self) -> Result<PartitionIoToken, RecvError> { + self.receiver.recv().await + } + + pub(crate) fn try_recv(&self) -> Result<PartitionIoToken, TryRecvError> { + self.receiver.try_recv().or_else(|error| { + if !self.undelivered.get() { + return Err(error); + } + self.slots + .borrow() + .iter() + .enumerate() + .find_map(|(index, slot)| { + let slot = slot.as_ref()?; + (slot.state.get() == SlotState::Queued).then(|| PartitionIoToken { + slot: index, + identity: slot.identity.get().expect("queued slot has an identity"), + }) + }) + .ok_or_else(|| { + self.undelivered.set(false); + error + }) + }) + } + + pub(crate) fn take_result(&self, token: PartitionIoToken) -> Option<PartitionIoResult> { + let slots = self.slots.borrow(); + let slot = slots.get(token.slot)?.as_ref()?; + if slot.identity.get() != Some(token.identity) || slot.state.get() != SlotState::Queued { + return None; + } + let result = slot.result.borrow_mut().take()?; + slot.state.set(SlotState::Settled); + self.counters.queued.set(self.counters.queued.get() - 1); + Some(result) + } + + pub(crate) fn settle(&self, token: PartitionIoToken, retain: bool) { + let slot = self.slots.borrow()[token.slot].clone(); + let Some(slot) = slot.filter(|slot| { + slot.identity.get() == Some(token.identity) && slot.state.get() == SlotState::Settled + }) else { + return; + }; + if let Some(guard) = slot.gate.borrow_mut().take() { + guard.release(); + } + slot.resources.borrow_mut().take(); + if let Some(quiescence) = slot.quiescence.borrow_mut().take() { + quiescence.settle(token.identity); + } + if !retain { + self.release(token.slot); + } + } + + pub(crate) fn release(&self, index: usize) { + let mut slots = self.slots.borrow_mut(); + if slots[index].as_ref().is_some_and(|slot| { + matches!(slot.state.get(), SlotState::Reserved | SlotState::Settled) + }) { + let slot = slots[index].take().expect("settled slot exists"); + self.charged.set(self.charged.get() - slot.charge); + self.capacity_blocked.set(false); + } + } + + pub(crate) fn retained( + &self, + namespace: IggyNamespace, + incarnation: PartitionIncarnation, + ) -> Option<PartitionIoToken> { + self.slots + .borrow() + .iter() + .enumerate() + .find_map(|(index, slot)| { + let slot = slot.as_ref()?; + (slot.namespace == namespace + && slot.incarnation == incarnation + && slot.state.get() == SlotState::Settled) + .then(|| { + slot.identity.get().map(|identity| PartitionIoToken { + slot: index, + identity, + }) + }) + .flatten() + }) + } + + pub(crate) fn interrupted(&self) -> Vec<PartitionIoIdentity> { + if !self.interrupted.replace(false) { + return Vec::new(); + } + self.slots + .borrow() + .iter() + .flatten() + .filter(|slot| slot.state.get() == SlotState::Interrupted) + .filter_map(|slot| slot.identity.get()) + .collect() + } + + pub(crate) fn outstanding(&self) -> usize { + self.slots + .borrow() + .iter() + .flatten() + .filter(|slot| slot.state.get() != SlotState::Interrupted) + .count() + } + + pub(crate) fn close(&self) { + self.closed.set(true); + self.ready.waker.borrow_mut().take(); + } + + fn record_metrics(&self, metrics: &crate::metrics::ShardMetrics) { + metrics.set_partition_io( + self.counters.active.get(), + self.counters.queued.get(), + self.charged.get(), + self.ready.present.borrow().len(), + self.counters.quarantined.get(), + ); + } +} + +impl<B: MessageBus + 'static, MJ, S, M, T, SB: SuperblockStore + 'static> + IggyShard<B, MJ, S, M, T, SB> +where + MJ: crate::JournalHandle, + MJ::Target: journal::Journal< + Entry = server_common::Message<iggy_binary_protocol::PrepareHeader>, + Header = iggy_binary_protocol::PrepareHeader, + >, + M: crate::RestorableMetadataStm, + T: crate::ShardsTable, +{ + pub(crate) fn accept_partition_io_completion(&self, token: PartitionIoToken) { + let Some(result) = self.partition_io.take_result(token) else { + return; + }; + let retained = if let Some(partition) = self + .plane + .partitions() + .get_io_owner(&token.identity.namespace) + .filter(|partition| partition.incarnation() == token.identity.incarnation) + { + if let Err(error) = partition.accept_io(token.identity, result) { + tracing::error!(namespace_raw = token.identity.namespace.inner(), %error, "partition I/O acceptance failed"); + } + if token.identity.continuation == partitions::PartitionIoContinuation::Retention { + self.drop_partition_transfer_state(token.identity.namespace, partition); + } + partition.retains_io_reservation(token.identity) + } else { + drop(result); + false + }; + self.partition_io.settle(token, retained); + } + + /// Bounded completion and continuation service after each ordinary pump event. + #[allow(clippy::future_not_send, clippy::too_many_lines)] + pub(crate) async fn service_partition_io(&self) -> bool { Review Comment: simplification: `service_partition_io` returns a `bool` that every call site discards. Return `()` and delete `progressed` and the two early returns that carry it. ########## core/shard/src/partition_io.rs: ########## @@ -0,0 +1,1937 @@ +// 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. + +use std::cell::{Cell, RefCell}; +use std::collections::{BTreeSet, VecDeque}; +use std::rc::Rc; +use std::task::Waker; + +use consensus::PartitionsHandle; +use crossfire::{RecvError, TryRecvError}; +use journal::local_gate::OwnedLocalGateGuard; +use journal::superblock::SuperblockStore; +use message_bus::MessageBus; +use partitions::{ + CapturedPartitionIo, PartitionIncarnation, PartitionIoIdentity, PartitionIoResources, + PartitionIoResult, +}; +use server_common::sharding::IggyNamespace; +use thiserror::Error; + +use crate::{IggyShard, Receiver, Sender, channel}; + +pub const DEFAULT_PARTITION_IO_CAPACITY: usize = 16; +pub const DEFAULT_PARTITION_IO_BYTES: usize = 256 * 1024 * 1024; +pub const PARTITION_IO_CAPACITY_MAX: usize = 1 << 20; + +#[derive(Clone, Copy, Debug)] +pub struct PartitionIoLimits { + capacity: usize, + bytes_max: usize, +} + +#[derive(Debug, Error)] +pub enum PartitionIoLimitsError { + #[error("sharding.partition_io_capacity must be in 1..={PARTITION_IO_CAPACITY_MAX}; got {0}")] + Capacity(usize), + #[error("partition I/O allocation charge exceeds addressable memory")] + Overflow, + #[error( + "sharding.partition_io_bytes_max must be at least {minimum} and fit addressable memory; got {value}" + )] + Bytes { value: usize, minimum: usize }, +} + +impl PartitionIoLimits { + /// Resolve omitted bytes using the same allocation calculation as dispatch. + /// + /// # Errors + /// Rejects invalid slot counts, arithmetic overflow and undersized byte limits. + pub fn new(capacity: usize, bytes_max: Option<usize>) -> Result<Self, PartitionIoLimitsError> { + if capacity == 0 || capacity > PARTITION_IO_CAPACITY_MAX { + return Err(PartitionIoLimitsError::Capacity(capacity)); + } + let minimum = partitions::largest_legal_job_charge() + .filter(|charge| isize::try_from(*charge).is_ok()) + .ok_or(PartitionIoLimitsError::Overflow)?; + let bytes_max = bytes_max.unwrap_or_else(|| DEFAULT_PARTITION_IO_BYTES.max(minimum)); + if bytes_max < minimum || bytes_max > isize::MAX as usize { + return Err(PartitionIoLimitsError::Bytes { + value: bytes_max, + minimum, + }); + } + Ok(Self { + capacity, + bytes_max, + }) + } + + #[must_use] + pub const fn capacity(self) -> usize { + self.capacity + } + + #[must_use] + pub const fn bytes_max(self) -> usize { + self.bytes_max + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct PartitionIoToken { + slot: usize, + identity: PartitionIoIdentity, +} + +#[derive(Clone, Copy, PartialEq, Eq)] +enum SlotState { + Reserved, + Running, + Queued, + Settled, + Interrupted, +} + +struct PartitionIoSlot<SB> { + namespace: IggyNamespace, + incarnation: PartitionIncarnation, + identity: Cell<Option<PartitionIoIdentity>>, + state: Cell<SlotState>, + charge: usize, + result: RefCell<Option<PartitionIoResult>>, + resources: RefCell<Option<PartitionIoResources<SB>>>, + gate: RefCell<Option<OwnedLocalGateGuard>>, + quiescence: RefCell<Option<Rc<partitions::PartitionIoQuiescence>>>, +} + +#[derive(Default)] +struct ReadyPartitions { + queue: RefCell<VecDeque<(IggyNamespace, PartitionIncarnation)>>, + present: RefCell<BTreeSet<(IggyNamespace, PartitionIncarnation)>>, + waker: RefCell<Option<Waker>>, +} + +#[derive(Default)] +struct IoCounters { + active: Cell<usize>, + queued: Cell<usize>, + quarantined: Cell<usize>, +} + +impl ReadyPartitions { + fn notify(&self, namespace: IggyNamespace, incarnation: PartitionIncarnation) { + if !self.present.borrow_mut().insert((namespace, incarnation)) { + return; + } + self.queue.borrow_mut().push_back((namespace, incarnation)); + let waker = self.waker.borrow().clone(); + if let Some(waker) = waker { + waker.wake(); + } + } +} + +/// Slots outlive mounted lookup. Tokens never carry local file handles or results. +pub struct PartitionIoLane<SB> { + pub(crate) limits: PartitionIoLimits, + slots: RefCell<Vec<Option<Rc<PartitionIoSlot<SB>>>>>, + charged: Cell<usize>, + sender: Sender<PartitionIoToken>, + receiver: Receiver<PartitionIoToken>, + ready: Rc<ReadyPartitions>, + interrupted: Rc<Cell<bool>>, + undelivered: Rc<Cell<bool>>, + closed: Cell<bool>, + capacity_blocked: Cell<bool>, + counters: Rc<IoCounters>, + #[cfg(test)] + execution_gate: RefCell<Option<futures::channel::oneshot::Receiver<()>>>, +} + +impl<SB: SuperblockStore> PartitionIoLane<SB> { + pub(crate) fn new(limits: PartitionIoLimits) -> Self { + let (sender, receiver) = channel(limits.capacity); + Self { + limits, + slots: RefCell::new((0..limits.capacity).map(|_| None).collect()), + charged: Cell::new(0), + sender, + receiver, + ready: Rc::default(), + interrupted: Rc::new(Cell::new(false)), + undelivered: Rc::new(Cell::new(false)), + closed: Cell::new(false), + capacity_blocked: Cell::new(false), + counters: Rc::default(), + #[cfg(test)] + execution_gate: RefCell::new(None), + } + } + + pub(crate) fn notifier(&self) -> partitions::PartitionIoNotifier { + let ready = Rc::clone(&self.ready); + Rc::new(move |namespace, incarnation| ready.notify(namespace, incarnation)) + } + + pub(crate) fn register_waker(&self, waker: &Waker) { + let mut current = self.ready.waker.borrow_mut(); + if current + .as_ref() + .is_none_or(|current| !current.will_wake(waker)) + { + *current = Some(waker.clone()); + } + } + + pub(crate) fn has_ready(&self) -> bool { + self.head().is_some() || self.undelivered.get() || self.interrupted.get() + } + + pub(crate) fn head(&self) -> Option<(IggyNamespace, PartitionIncarnation)> { + if !self.capacity_blocked.get() { + return self.ready.queue.borrow().front().copied(); + } + let present = self.ready.present.borrow(); + self.slots.borrow().iter().flatten().find_map(|slot| { + (slot.state.get() == SlotState::Settled + && present.contains(&(slot.namespace, slot.incarnation))) + .then_some((slot.namespace, slot.incarnation)) + }) + } + + pub(crate) fn pop_ready(&self, namespace: IggyNamespace, incarnation: PartitionIncarnation) { + self.ready + .queue + .borrow_mut() + .retain(|queued| *queued != (namespace, incarnation)); + self.ready + .present + .borrow_mut() + .remove(&(namespace, incarnation)); + } + + pub(crate) fn reschedule(&self, namespace: IggyNamespace, incarnation: PartitionIncarnation) { + self.ready.notify(namespace, incarnation); + } + + pub(crate) fn try_reserve( + &self, + namespace: IggyNamespace, + incarnation: PartitionIncarnation, + charge: usize, + ) -> Option<usize> { + if self.closed.get() { + return None; + } + let mut slots = self.slots.borrow_mut(); + if let Some((index, existing)) = slots.iter().enumerate().find_map(|(index, slot)| { + slot.as_ref() + .filter(|slot| slot.namespace == namespace && slot.incarnation == incarnation) + .map(|slot| (index, slot)) + }) { + return (existing.state.get() == SlotState::Settled && charge <= existing.charge) + .then_some(index); + } + let charged = self.charged.get().checked_add(charge)?; + if charged > self.limits.bytes_max { + return None; + } + let index = slots.iter().position(Option::is_none)?; + slots[index] = Some(Rc::new(PartitionIoSlot { + namespace, + incarnation, + identity: Cell::new(None), + state: Cell::new(SlotState::Reserved), + charge, + result: RefCell::new(None), + resources: RefCell::new(None), + gate: RefCell::new(None), + quiescence: RefCell::new(None), + })); + self.charged.set(charged); + Some(index) + } + + pub(crate) fn dispatch( + &self, + index: usize, + captured: CapturedPartitionIo<SB>, + bus: &impl MessageBus, + ) where + SB: 'static, + { + let slot = Rc::clone( + self.slots.borrow()[index] + .as_ref() + .expect("reserved partition I/O slot"), + ); + let CapturedPartitionIo { + identity, + job, + gate, + quiescence, + } = captured; + slot.identity.set(Some(identity)); + *slot.resources.borrow_mut() = Some(job.retain_resources()); + *slot.gate.borrow_mut() = gate; + *slot.quiescence.borrow_mut() = Some(Rc::clone(&quiescence)); + slot.state.set(SlotState::Running); + self.counters.active.set(self.counters.active.get() + 1); + let marker = InterruptionMarker { + slot: Rc::clone(&slot), + interrupted: Rc::clone(&self.interrupted), + quiescence, + counters: Rc::clone(&self.counters), + }; + let sender = self.sender.clone(); + let undelivered = Rc::clone(&self.undelivered); + let ready = Rc::clone(&self.ready); + let counters = Rc::clone(&self.counters); + #[cfg(test)] + let execution_gate = self.execution_gate.borrow_mut().take(); + bus.spawn(async move { + #[cfg(test)] + if let Some(execution_gate) = execution_gate { + execution_gate + .await + .expect("test releases captured file job"); + } + let result = job.execute().await; + *slot.result.borrow_mut() = Some(result); + slot.state.set(SlotState::Queued); + counters.active.set(counters.active.get() - 1); + counters.queued.set(counters.queued.get() + 1); + if sender Review Comment: simplification: Each slot holds at most one token and the channel capacity equals the slot count, so `try_send` never fails and the `undelivered` flag and its slot-scan fallback in `try_recv` are dead. Delete the flag and the fallback. ########## core/configs/src/server_config/sharding.rs: ########## @@ -175,6 +182,17 @@ impl Default for ShardingConfig { impl Validatable<ConfigurationError> for ShardingConfig { fn validate(&self) -> Result<(), ConfigurationError> { + if self.partition_io_capacity == 0 + || self.partition_io_capacity > INBOX_CAPACITY_MAX Review Comment: nit: The partition I/O capacity cap reuses `INBOX_CAPACITY_MAX`, whose doc covers only inbox and completion lanes, and it duplicates `PARTITION_IO_CAPACITY_MAX` in the shard crate. Use one shared constant and extend its doc. ########## core/shard/src/lib.rs: ########## @@ -598,6 +600,8 @@ impl std::ops::Deref for TaggedSender { #[derive(Debug, thiserror::Error)] #[non_exhaustive] pub enum ShardCtorError { + #[error(transparent)] + PartitionIoLimits(#[from] PartitionIoLimitsError), Review Comment: nit: The enum doc says both constructor failures are bootstrap programming errors, but this new variant carries a limit-validation failure, and the `# Errors` lists on `IggyShard::new` and `build` omit it. Document the variant in all three places. ########## core/configs/src/server_config/sharding.rs: ########## @@ -175,6 +182,17 @@ impl Default for ShardingConfig { impl Validatable<ConfigurationError> for ShardingConfig { fn validate(&self) -> Result<(), ConfigurationError> { + if self.partition_io_capacity == 0 + || self.partition_io_capacity > INBOX_CAPACITY_MAX + || self + .partition_io_bytes_max + .is_some_and(|bytes| bytes == 0 || bytes > isize::MAX as usize) + { + eprintln!( + "Invalid sharding configuration: partition I/O limits must be positive and fit addressable capacity" Review Comment: nit: This message covers both new keys and names neither the key, the value, nor the bound, while every sibling block prints all three. Validate each key on its own and name it in the message. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
