mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-01 02:52:15 +00:00
Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 6befa4d042 | |||
| 06c70fcd14 |
@@ -12,18 +12,24 @@
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
# Companion to ci.yml for the required "Test and Lint" status check.
|
||||
# Companion to ci.yml for required status checks.
|
||||
#
|
||||
# ci.yml skips docs-only pull requests via paths-ignore, but the branch
|
||||
# ruleset requires a check named "Test and Lint" — without this workflow a
|
||||
# docs-only PR would wait on that check forever. This workflow triggers on
|
||||
# exactly the paths ci.yml ignores and reports an instant success under the
|
||||
# same job name. Mixed PRs trigger both workflows and the real check still
|
||||
# gates: a required check with any failing run blocks the merge.
|
||||
# ci.yml skips docs-only pull requests via paths-ignore, but the branch ruleset
|
||||
# requires a check named "Test and Lint" — without this workflow a docs-only PR
|
||||
# would wait on it forever. This workflow triggers on exactly the paths ci.yml
|
||||
# ignores and reports success under the same job name. Mixed PRs trigger both
|
||||
# workflows and the real check still gates: a required check with any failing
|
||||
# run blocks the merge.
|
||||
# https://docs.github.com/en/repositories/configuring-branches-and-merges-in-your-repository/defining-the-mergeability-of-pull-requests/troubleshooting-required-status-checks#handling-skipped-but-required-checks
|
||||
#
|
||||
# "Quick Checks" is mirrored here ahead of the ruleset change that will make it
|
||||
# required too (rustfs/backlog#1599). It is not required yet, so today this job
|
||||
# is inert; adding it first is what lets that ruleset change land without
|
||||
# stranding docs-only PRs on a check nobody reports.
|
||||
#
|
||||
# Keep the paths list below in sync with the pull_request paths-ignore list
|
||||
# in ci.yml.
|
||||
# in ci.yml, and keep the quick-checks steps below byte-identical to the
|
||||
# quick-checks job in ci.yml.
|
||||
|
||||
name: Continuous Integration (docs only)
|
||||
|
||||
@@ -52,9 +58,63 @@ permissions:
|
||||
contents: read
|
||||
|
||||
jobs:
|
||||
# Deliberately NOT a bare `echo`. Once "Quick Checks" becomes a required
|
||||
# check, ci.yml gates every expensive job behind it, so a mixed PR reports
|
||||
# two check runs with this name: the real one (45-51s) and this companion.
|
||||
# GitHub has no written contract for how it picks between same-named
|
||||
# required check runs ("latest wins" vs "any failure blocks"), so instead of
|
||||
# relying on ordering we make both runs execute the same commands against
|
||||
# the same merge ref — their conclusions are then necessarily identical and
|
||||
# the choice does not matter. Keep these steps byte-identical to the
|
||||
# quick-checks job in ci.yml (a guard script that asserts this, and the paths
|
||||
# sync below, is tracked in rustfs/backlog#1603).
|
||||
#
|
||||
# For a genuinely docs-only PR this adds no strictness (no code changed, so
|
||||
# fmt and the guards always pass) and costs ~50s of ubuntu-latest.
|
||||
quick-checks:
|
||||
name: Quick Checks
|
||||
runs-on: ubuntu-latest
|
||||
timeout-minutes: 10
|
||||
steps:
|
||||
- name: Checkout repository
|
||||
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
|
||||
|
||||
- name: Install ripgrep
|
||||
run: sudo apt-get update && sudo apt-get install -y ripgrep
|
||||
|
||||
- name: Install Rust toolchain
|
||||
uses: dtolnay/rust-toolchain@29eef336d9b2848a0b548edc03f92a220660cdb8 # stable
|
||||
with:
|
||||
components: rustfmt
|
||||
|
||||
- name: Check code formatting
|
||||
run: cargo fmt --all --check
|
||||
|
||||
- name: Check unsafe code allowances
|
||||
run: ./scripts/check_unsafe_code_allowances.sh
|
||||
|
||||
- name: Check layered dependencies
|
||||
run: ./scripts/check_layer_dependencies.sh
|
||||
|
||||
- name: Check architecture migration rules
|
||||
run: ./scripts/check_architecture_migration_rules.sh
|
||||
|
||||
- name: Check tokio io-uring feature guard
|
||||
run: ./scripts/check_no_tokio_io_uring.sh
|
||||
|
||||
- name: Check extension schema boundaries
|
||||
run: ./scripts/check_extension_schema_boundaries.sh
|
||||
|
||||
- name: Check body-cache whitelist guard
|
||||
run: ./scripts/check_body_cache_whitelist.sh
|
||||
|
||||
- name: Check no planning docs committed
|
||||
run: ./scripts/check_no_planning_docs.sh
|
||||
|
||||
test-and-lint:
|
||||
name: Test and Lint
|
||||
runs-on: ubuntu-latest
|
||||
timeout-minutes: 10
|
||||
steps:
|
||||
- name: Checkout repository
|
||||
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
|
||||
|
||||
@@ -96,6 +96,10 @@ jobs:
|
||||
|
||||
# Fast, compile-free checks that fail early so contributors get feedback in
|
||||
# ~1 minute instead of waiting for the full test job.
|
||||
#
|
||||
# These steps are mirrored byte-for-byte in ci-docs-only.yml so that a mixed
|
||||
# PR, which reports two check runs named "Quick Checks", cannot get one red
|
||||
# and one green. Edit both jobs together.
|
||||
quick-checks:
|
||||
name: Quick Checks
|
||||
if: github.event_name != 'pull_request' || github.event.action != 'closed'
|
||||
@@ -448,6 +452,13 @@ jobs:
|
||||
|
||||
uring-integration:
|
||||
name: io_uring Integration (real)
|
||||
# The pull_request trigger includes `closed` purely so the concurrency
|
||||
# group cancels in-flight runs of a closed PR; every other job opts out of
|
||||
# that run with this guard (or is skipped through its `needs` chain). This
|
||||
# job had neither, so each closed/merged PR really ran the whole io_uring
|
||||
# suite (measured 4m17s / 7m19s / 7m31s on runs 30678272341 / 30678117601 /
|
||||
# 30662728539) and kept the cancellation run in progress for minutes.
|
||||
if: github.event_name != 'pull_request' || github.event.action != 'closed'
|
||||
# GitHub-hosted ubuntu-latest runs a recent kernel with io_uring and, unlike
|
||||
# a container, applies no seccomp filter that would block io_uring_setup — so
|
||||
# the probe succeeds and the tests exercise the real UringBackend/FdCache/
|
||||
@@ -746,7 +757,13 @@ jobs:
|
||||
# evaluates ILM within ~2s of the due time, well inside the poll window.
|
||||
s3-lifecycle-behavior-tests:
|
||||
name: S3 Lifecycle Behavior Tests
|
||||
needs: [ build-rustfs-debug-binary ]
|
||||
# Also gated on e2e-tests, matching s3-implemented-tests: when the e2e smoke
|
||||
# suite is already red this lane cannot tell us anything new, and it holds a
|
||||
# sm-standard-4 for up to 30 minutes doing so. Both lanes only download the
|
||||
# prebuilt debug binary (no cargo build), and s3-implemented-tests — which
|
||||
# already waits on e2e-tests — finishes later anyway, so a green PR's total
|
||||
# wall clock is unchanged.
|
||||
needs: [ build-rustfs-debug-binary, e2e-tests ]
|
||||
runs-on: sm-standard-4
|
||||
timeout-minutes: 30
|
||||
steps:
|
||||
|
||||
@@ -382,6 +382,7 @@ struct ManualTransitionRunReport {
|
||||
skipped_delete_marker: u64,
|
||||
skipped_directory: u64,
|
||||
skipped_replication: u64,
|
||||
skipped_already_transitioned: u64,
|
||||
skipped_already_in_flight: u64,
|
||||
skipped_queue_full: u64,
|
||||
skipped_queue_closed: u64,
|
||||
@@ -407,6 +408,41 @@ fn assert_completed_or_in_flight_partial(state: &str, report: &ManualTransitionR
|
||||
}
|
||||
}
|
||||
|
||||
fn assert_conflict_winner_report(state: &str, report: &ManualTransitionRunReport, expected_objects: u64, context: &str) {
|
||||
assert_completed_or_in_flight_partial(state, report, context);
|
||||
if report.skipped_already_in_flight > 0 {
|
||||
assert!(
|
||||
report.scanned <= expected_objects,
|
||||
"{context}: scanned more objects than the conflict scope contains: {report:#?}"
|
||||
);
|
||||
assert!(
|
||||
report.eligible <= expected_objects,
|
||||
"{context}: marked more objects eligible than the conflict scope contains: {report:#?}"
|
||||
);
|
||||
assert_eq!(
|
||||
report.enqueued + report.skipped_already_in_flight,
|
||||
report.eligible,
|
||||
"{context}: partial in-flight accounting must cover every eligible object: {report:#?}"
|
||||
);
|
||||
} else {
|
||||
assert_eq!(report.scanned, expected_objects, "{context}: {report:#?}");
|
||||
assert_eq!(
|
||||
report.eligible + report.skipped_already_transitioned,
|
||||
expected_objects,
|
||||
"{context}: {report:#?}"
|
||||
);
|
||||
assert_eq!(
|
||||
report.enqueued + report.skipped_already_in_flight,
|
||||
expected_objects,
|
||||
"{context}: {report:#?}"
|
||||
);
|
||||
}
|
||||
assert_eq!(
|
||||
report.transition_completed, report.enqueued,
|
||||
"{context}: winner must wait for all queued transitions: {report:#?}"
|
||||
);
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
struct ManualTransitionQueueSnapshot {
|
||||
queue_capacity: u64,
|
||||
@@ -1240,15 +1276,8 @@ async fn test_manual_transition_async_scope_conflicts_report_active_job() -> Tes
|
||||
cold_client.create_bucket().bucket(TIER_BUCKET).send().await?;
|
||||
|
||||
let mut hot = RustFSTestEnvironment::new().await?;
|
||||
hot.start_rustfs_server_with_env(
|
||||
vec![],
|
||||
&[
|
||||
("RUSTFS_SCANNER_ENABLED", "false"),
|
||||
("RUSTFS_SCANNER_CYCLE", "3600"),
|
||||
(MANUAL_TRANSITION_CANCEL_BARRIER_ENV, "1"),
|
||||
],
|
||||
)
|
||||
.await?;
|
||||
hot.start_rustfs_server_with_env(vec![], &[("RUSTFS_SCANNER_ENABLED", "false"), ("RUSTFS_SCANNER_CYCLE", "3600")])
|
||||
.await?;
|
||||
let hot_client = hot.create_s3_client();
|
||||
add_rustfs_tier(&hot, &cold).await?;
|
||||
|
||||
@@ -1318,21 +1347,20 @@ async fn test_manual_transition_async_scope_conflicts_report_active_job() -> Tes
|
||||
assert_eq!(conflict.cancel_endpoint, status_endpoint);
|
||||
assert!(!conflict.scope_key.is_empty());
|
||||
|
||||
manual_transition_job_cancel(&hot, cancel_endpoint).await?;
|
||||
|
||||
let terminal = wait_for_manual_transition_job_terminal(&hot, status_endpoint, MANUAL_ASYNC_CONFLICT_TERMINAL_TIMEOUT).await?;
|
||||
assert_eq!(terminal.job_id, job_id);
|
||||
assert_eq!(terminal.status, "cancelled", "terminal conflict winner response: {terminal:#?}");
|
||||
assert!(!terminal.report.dry_run);
|
||||
assert_eq!(terminal.report.bucket, MANUAL_ASYNC_CONFLICT_BUCKET);
|
||||
assert_eq!(terminal.report.prefix, accepted.report.prefix);
|
||||
assert!(terminal.report.cancelled, "terminal conflict winner response: {terminal:#?}");
|
||||
assert_eq!(terminal.report.scanned, 0, "terminal conflict winner response: {terminal:#?}");
|
||||
assert_eq!(terminal.report.enqueued, 0, "terminal conflict winner response: {terminal:#?}");
|
||||
assert_eq!(
|
||||
terminal.report.transition_completed, 0,
|
||||
"terminal conflict winner response: {terminal:#?}"
|
||||
assert_conflict_winner_report(
|
||||
&terminal.status,
|
||||
&terminal.report,
|
||||
MANUAL_ASYNC_CONFLICT_OBJECTS as u64,
|
||||
"terminal conflict winner response",
|
||||
);
|
||||
assert_eq!(terminal.report.dry_run_eligible, 0, "terminal conflict winner response: {terminal:#?}");
|
||||
assert_eq!(terminal.report.transition_failed, 0, "terminal conflict winner response: {terminal:#?}");
|
||||
assert_eq!(terminal.report.tier_failure, 0, "terminal conflict winner response: {terminal:#?}");
|
||||
let after_remote_count = cold_tier_object_count(&cold_client).await?;
|
||||
assert!(after_remote_count >= before_remote_count);
|
||||
assert!(after_remote_count <= before_remote_count + MANUAL_ASYNC_CONFLICT_OBJECTS);
|
||||
|
||||
@@ -75,7 +75,7 @@ const DATA_SCANNER_COMPACT_LEAST_OBJECT: usize = 500;
|
||||
const DATA_SCANNER_COMPACT_AT_CHILDREN: usize = 10000;
|
||||
const DATA_SCANNER_COMPACT_AT_FOLDERS: usize = DATA_SCANNER_COMPACT_AT_CHILDREN / 4;
|
||||
const DATA_SCANNER_FORCE_COMPACT_AT_FOLDERS: usize = 250_000;
|
||||
const SCANNER_LIST_PATH_RAW_STALL_TIMEOUT: Duration = Duration::from_secs(60);
|
||||
const SCANNER_LIST_PATH_RAW_TIMEOUT: Duration = Duration::from_secs(60);
|
||||
const SCANNER_ENTRY_PROGRESS_BATCH: u64 = 32;
|
||||
const SCANNER_ENTRY_PROGRESS_INTERVAL: Duration = Duration::from_secs(30);
|
||||
const DEFAULT_HEAL_OBJECT_SELECT_PROB: u32 = 1024;
|
||||
@@ -100,21 +100,6 @@ static SCANNER_INLINE_HEAL_WARN_ONCE: Once = Once::new();
|
||||
static SCANNER_INLINE_HEAL_METRICS_ONCE: Once = Once::new();
|
||||
static SCANNER_ALERT_METRICS_ONCE: Once = Once::new();
|
||||
|
||||
#[cfg(test)]
|
||||
type ListPathRawTimeoutSnapshot = (bool, Option<Duration>, Option<Duration>);
|
||||
|
||||
fn scanner_abandoned_child_list_options() -> ListPathRawOptions {
|
||||
// A complete heal walk scales with bucket size and may legitimately take
|
||||
// longer than a fixed wall-clock budget. Keep the total duration unbounded;
|
||||
// Retain the scanner's per-read stall budget and keep cancellation controlled
|
||||
// by the scanner cycle token.
|
||||
ListPathRawOptions {
|
||||
skip_walkdir_total_timeout: true,
|
||||
walkdir_stall_timeout: Some(SCANNER_LIST_PATH_RAW_STALL_TIMEOUT),
|
||||
..Default::default()
|
||||
}
|
||||
}
|
||||
|
||||
pub fn data_usage_update_dir_cycles() -> u32 {
|
||||
rustfs_utils::get_env_u32(ENV_DATA_USAGE_UPDATE_DIR_CYCLES, DATA_USAGE_UPDATE_DIR_CYCLES)
|
||||
}
|
||||
@@ -1296,8 +1281,6 @@ pub struct FolderScanner {
|
||||
skip_heal: Arc<std::sync::atomic::AtomicBool>,
|
||||
local_disk: Arc<Disk>,
|
||||
pending_heals_changed: bool,
|
||||
#[cfg(test)]
|
||||
list_path_raw_options_observer: Option<mpsc::UnboundedSender<ListPathRawTimeoutSnapshot>>,
|
||||
}
|
||||
|
||||
impl FolderScanner {
|
||||
@@ -2427,79 +2410,75 @@ impl FolderScanner {
|
||||
let bucket_clone = bucket.clone();
|
||||
let prefix_clone = prefix.clone();
|
||||
let child_ctx_clone = child_ctx.clone();
|
||||
#[cfg(test)]
|
||||
let list_path_raw_options_observer = self.list_path_raw_options_observer.clone();
|
||||
|
||||
tokio::spawn(async move {
|
||||
let options = ListPathRawOptions {
|
||||
disks,
|
||||
bucket: bucket_clone.clone(),
|
||||
path: prefix_clone.clone(),
|
||||
recursive: true,
|
||||
report_not_found: true,
|
||||
min_disks: disks_quorum,
|
||||
agreed: Some(Box::new(move |entry: MetaCacheEntry| {
|
||||
let entry_name = entry.name.clone();
|
||||
let agreed_tx = agreed_tx.clone();
|
||||
Box::pin(async move {
|
||||
if let Err(e) = agreed_tx.send(entry_name).await {
|
||||
error!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_FOLDER_STATE,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_FOLDER,
|
||||
entry = %entry.name,
|
||||
state = "list_path_agreed_send_failed",
|
||||
error = %e,
|
||||
"Scanner list_path_raw agreed callback failed"
|
||||
);
|
||||
}
|
||||
})
|
||||
})),
|
||||
partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option<DiskError>]| {
|
||||
let partial_tx = partial_tx.clone();
|
||||
Box::pin(async move {
|
||||
if let Err(e) = partial_tx.send(entries).await {
|
||||
error!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_FOLDER_STATE,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_FOLDER,
|
||||
state = "list_path_partial_send_failed",
|
||||
error = %e,
|
||||
"Scanner list_path_raw partial callback failed"
|
||||
);
|
||||
}
|
||||
})
|
||||
})),
|
||||
finished: Some(Box::new(move |errs: &[Option<DiskError>]| {
|
||||
let finished_tx = finished_tx.clone();
|
||||
let errs_clone = errs.to_vec();
|
||||
Box::pin(async move {
|
||||
if let Err(e) = finished_tx.send(errs_clone).await {
|
||||
error!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_FOLDER_STATE,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_FOLDER,
|
||||
state = "list_path_finished_send_failed",
|
||||
error = %e,
|
||||
"Scanner list_path_raw finished callback failed"
|
||||
);
|
||||
}
|
||||
})
|
||||
})),
|
||||
..scanner_abandoned_child_list_options()
|
||||
};
|
||||
#[cfg(test)]
|
||||
if let Some(observer) = list_path_raw_options_observer {
|
||||
let _ = observer.send((
|
||||
options.skip_walkdir_total_timeout,
|
||||
options.walkdir_timeout,
|
||||
options.walkdir_stall_timeout,
|
||||
));
|
||||
}
|
||||
if let Err(e) = list_path_raw(child_ctx_clone.clone(), options).await {
|
||||
if let Err(e) = list_path_raw(
|
||||
child_ctx_clone.clone(),
|
||||
ListPathRawOptions {
|
||||
disks,
|
||||
bucket: bucket_clone.clone(),
|
||||
path: prefix_clone.clone(),
|
||||
recursive: true,
|
||||
report_not_found: true,
|
||||
min_disks: disks_quorum,
|
||||
walkdir_timeout: Some(SCANNER_LIST_PATH_RAW_TIMEOUT),
|
||||
walkdir_stall_timeout: Some(SCANNER_LIST_PATH_RAW_TIMEOUT),
|
||||
agreed: Some(Box::new(move |entry: MetaCacheEntry| {
|
||||
let entry_name = entry.name.clone();
|
||||
let agreed_tx = agreed_tx.clone();
|
||||
Box::pin(async move {
|
||||
if let Err(e) = agreed_tx.send(entry_name).await {
|
||||
error!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_FOLDER_STATE,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_FOLDER,
|
||||
entry = %entry.name,
|
||||
state = "list_path_agreed_send_failed",
|
||||
error = %e,
|
||||
"Scanner list_path_raw agreed callback failed"
|
||||
);
|
||||
}
|
||||
})
|
||||
})),
|
||||
partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option<DiskError>]| {
|
||||
let partial_tx = partial_tx.clone();
|
||||
Box::pin(async move {
|
||||
if let Err(e) = partial_tx.send(entries).await {
|
||||
error!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_FOLDER_STATE,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_FOLDER,
|
||||
state = "list_path_partial_send_failed",
|
||||
error = %e,
|
||||
"Scanner list_path_raw partial callback failed"
|
||||
);
|
||||
}
|
||||
})
|
||||
})),
|
||||
finished: Some(Box::new(move |errs: &[Option<DiskError>]| {
|
||||
let finished_tx = finished_tx.clone();
|
||||
let errs_clone = errs.to_vec();
|
||||
Box::pin(async move {
|
||||
if let Err(e) = finished_tx.send(errs_clone).await {
|
||||
error!(
|
||||
target: "rustfs::scanner::folder",
|
||||
event = EVENT_SCANNER_FOLDER_STATE,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_FOLDER,
|
||||
state = "list_path_finished_send_failed",
|
||||
error = %e,
|
||||
"Scanner list_path_raw finished callback failed"
|
||||
);
|
||||
}
|
||||
})
|
||||
})),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
{
|
||||
if is_missing_path_disk_error(&e) {
|
||||
debug!(
|
||||
target: "rustfs::scanner::folder",
|
||||
@@ -2841,8 +2820,6 @@ pub async fn scan_data_folder(
|
||||
skip_heal,
|
||||
local_disk,
|
||||
pending_heals_changed: false,
|
||||
#[cfg(test)]
|
||||
list_path_raw_options_observer: None,
|
||||
};
|
||||
|
||||
// Check if context is cancelled
|
||||
@@ -3049,7 +3026,6 @@ mod tests {
|
||||
skip_heal: Arc::new(AtomicBool::new(false)),
|
||||
local_disk: disk,
|
||||
pending_heals_changed: false,
|
||||
list_path_raw_options_observer: None,
|
||||
};
|
||||
|
||||
(scanner, temp_dir)
|
||||
@@ -4234,8 +4210,6 @@ mod tests {
|
||||
scanner.heal_object_select = 1;
|
||||
scanner.disks = disks;
|
||||
scanner.disks_quorum = 2;
|
||||
let (options_tx, mut options_rx) = mpsc::unbounded_channel();
|
||||
scanner.list_path_raw_options_observer = Some(options_tx);
|
||||
scanner.old_cache.replace(
|
||||
&format!("{bucket}/{object}"),
|
||||
bucket,
|
||||
@@ -4257,11 +4231,6 @@ mod tests {
|
||||
.expect("scan_folder should not hang after list_path_raw finishes")
|
||||
.expect("scan_folder should finish successfully");
|
||||
|
||||
let observed_options = tokio::time::timeout(Duration::from_secs(1), options_rx.recv())
|
||||
.await
|
||||
.expect("abandoned-child listing options should be observed promptly")
|
||||
.expect("abandoned-child listing options channel should remain open");
|
||||
assert_eq!(observed_options, (true, None, Some(SCANNER_LIST_PATH_RAW_STALL_TIMEOUT)));
|
||||
let root = scanner
|
||||
.new_cache
|
||||
.checked_flatten(bucket)
|
||||
|
||||
+1
-36
@@ -12,32 +12,9 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
#[cfg(all(feature = "hotpath", feature = "hotpath-alloc"))]
|
||||
use std::alloc::{GlobalAlloc, Layout};
|
||||
|
||||
#[cfg(all(feature = "hotpath", feature = "hotpath-alloc"))]
|
||||
#[derive(Default)]
|
||||
struct DefaultMiMalloc;
|
||||
|
||||
#[cfg(all(feature = "hotpath", feature = "hotpath-alloc"))]
|
||||
// SAFETY: allocation and deallocation are forwarded unchanged to MiMalloc, so
|
||||
// MiMalloc's GlobalAlloc guarantees apply to every returned pointer and layout.
|
||||
#[allow(unsafe_code)]
|
||||
unsafe impl GlobalAlloc for DefaultMiMalloc {
|
||||
unsafe fn alloc(&self, layout: Layout) -> *mut u8 {
|
||||
// SAFETY: the caller upholds GlobalAlloc's contract for layout.
|
||||
unsafe { mimalloc::MiMalloc.alloc(layout) }
|
||||
}
|
||||
|
||||
unsafe fn dealloc(&self, ptr: *mut u8, layout: Layout) {
|
||||
// SAFETY: ptr and layout came from this allocator and are forwarded unchanged.
|
||||
unsafe { mimalloc::MiMalloc.dealloc(ptr, layout) }
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(all(feature = "hotpath", feature = "hotpath-alloc"))]
|
||||
#[global_allocator]
|
||||
static GLOBAL: hotpath::CountingAllocator<DefaultMiMalloc> = hotpath::CountingAllocator::new();
|
||||
static GLOBAL: hotpath::CountingAllocator = hotpath::CountingAllocator::new();
|
||||
|
||||
#[cfg(not(all(feature = "hotpath", feature = "hotpath-alloc")))]
|
||||
#[global_allocator]
|
||||
@@ -48,15 +25,3 @@ fn main() {
|
||||
|
||||
rustfs::startup_entrypoint::run_process();
|
||||
}
|
||||
|
||||
#[cfg(all(test, feature = "hotpath", feature = "hotpath-alloc"))]
|
||||
mod tests {
|
||||
#[test]
|
||||
#[allow(unsafe_code)]
|
||||
fn hotpath_allocator_uses_mimalloc() {
|
||||
let allocation = Box::new([0_u8; 64]);
|
||||
|
||||
// SAFETY: the live Box pointer is valid to inspect for heap ownership.
|
||||
assert!(unsafe { libmimalloc_sys::mi_is_in_heap_region(allocation.as_ptr().cast()) });
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user