#![allow(clippy::module_inception)]
#[allow(unused_imports)]
pub use super::*;
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::AtomicUsize;
fn manager() -> OffloadManager {
OffloadManager::new(8)
}
#[tokio::test]
async fn spawn_deduplicates_same_key_while_in_flight() {
static EXECUTED: AtomicUsize = AtomicUsize::new(0);
let mgr = manager();
let first = mgr.spawn("k", async {
EXECUTED.fetch_add(1, Ordering::SeqCst);
tokio::time::sleep(Duration::from_millis(80)).await;
});
assert!(first, "first spawn must run");
let deadline = Instant::now() + Duration::from_secs(1);
while EXECUTED.load(Ordering::SeqCst) == 0 && Instant::now() < deadline {
tokio::time::sleep(Duration::from_millis(2)).await;
}
let second = mgr.spawn("k", async {
EXECUTED.fetch_add(1, Ordering::SeqCst);
});
assert!(
!second,
"same-key spawn while in flight must be deduplicated"
);
assert_eq!(
EXECUTED.load(Ordering::SeqCst),
1,
"deduped task must not run"
);
assert!(mgr.is_in_flight("k"));
let completed = mgr.wait_all(Duration::from_secs(2)).await;
assert_eq!(completed, 1);
assert!(!mgr.is_in_flight("k"), "completion must clear in_flight");
let third = mgr.spawn("k", async {});
assert!(third, "spawn after completion must run");
}
#[tokio::test]
async fn distinct_keys_spawn_concurrently() {
let mgr = manager();
assert!(mgr.spawn("a", async {
tokio::time::sleep(Duration::from_millis(50)).await
}));
assert!(mgr.spawn("b", async {
tokio::time::sleep(Duration::from_millis(50)).await
}));
assert_eq!(mgr.in_flight_count(), 2);
assert!(mgr.wait_all(Duration::from_secs(2)).await >= 2);
}
#[tokio::test]
async fn concurrency_limit_rejects_overflow() {
let mgr = OffloadManager::new(1);
assert!(mgr.spawn("slot", async {
tokio::time::sleep(Duration::from_millis(80)).await
}));
let rejected = mgr.spawn("other", async {});
assert!(
!rejected,
"spawn beyond max_concurrent_tasks must be rejected"
);
assert!(mgr.wait_all(Duration::from_secs(2)).await >= 1);
assert!(mgr.spawn("other", async {}));
mgr.wait_all(Duration::from_secs(2)).await;
}
#[tokio::test]
async fn cancel_all_clears_registry_only() {
let mgr = manager();
assert!(mgr.spawn("c1", async {
tokio::time::sleep(Duration::from_millis(60)).await
}));
assert_eq!(mgr.cancel_all(), 1);
assert_eq!(mgr.in_flight_count(), 0);
mgr.wait_all(Duration::from_secs(2)).await;
}
#[tokio::test]
async fn panic_in_task_does_not_leak_registry() {
static AFTER_PANIC: AtomicUsize = AtomicUsize::new(0);
let mgr = manager();
let _ = mgr.spawn("p", async {
AFTER_PANIC.fetch_add(1, Ordering::SeqCst);
panic!("simulated offload panic");
});
tokio::time::sleep(Duration::from_millis(120)).await;
assert!(
!mgr.is_in_flight("p"),
"panic must not leak the in_flight registry"
);
assert!(
mgr.spawn("p", async {}),
"spawn after panic must be accepted"
);
mgr.wait_all(Duration::from_secs(2)).await;
}
#[cfg(feature = "metrics")]
#[tokio::test]
#[serial_test::serial]
async fn metrics_counters_record_lifecycle() {
let mgr = manager();
let metrics = &crate::infra::metrics::unified::GLOBAL_UNIFIED_METRICS;
let spawned_before = dynamic_counter(metrics, "oxcache_offload_spawned_total");
let dedup_before = dynamic_counter(metrics, "oxcache_offload_deduplicated_total");
assert!(mgr.spawn("m", async {
tokio::time::sleep(Duration::from_millis(60)).await
}));
assert!(!mgr.spawn("m", async {}));
mgr.wait_all(Duration::from_secs(2)).await;
assert!(
dynamic_counter(metrics, "oxcache_offload_spawned_total") > spawned_before,
"spawn must increment its counter"
);
assert!(
dynamic_counter(metrics, "oxcache_offload_deduplicated_total") > dedup_before,
"dedup must increment its counter"
);
}
#[cfg(feature = "metrics")]
fn dynamic_counter(metrics: &crate::infra::metrics::unified::UnifiedMetrics, key: &str) -> u64 {
metrics
.get_dynamic_metrics()
.get(key)
.and_then(|v| match v {
crate::infra::metrics::unified::MetricValue::Counter(c) => Some(*c),
_ => None,
})
.unwrap_or(0)
}
#[tokio::test]
#[serial_test::serial]
async fn cancel_policy_counts_timeout() {
let mgr = OffloadManager::with_policy(2, TimeoutPolicy::Cancel(Duration::from_millis(30)));
#[cfg(feature = "metrics")]
let metrics = &crate::infra::metrics::unified::GLOBAL_UNIFIED_METRICS;
#[cfg(feature = "metrics")]
let timeout_before = dynamic_counter(metrics, "oxcache_offload_timeout_total");
assert!(mgr.spawn("slow", async {
tokio::time::sleep(Duration::from_secs(5)).await;
}));
let completed = mgr.wait_all(Duration::from_secs(2)).await;
assert_eq!(completed, 1, "cancelled task must release its slot");
#[cfg(feature = "metrics")]
assert!(
dynamic_counter(metrics, "oxcache_offload_timeout_total") > timeout_before,
"cancel-policy elapse must increment timeout counter"
);
}
#[tokio::test]
async fn accessors_report_normalized_limit_and_policy() {
let mgr = OffloadManager::with_policy(0, TimeoutPolicy::None);
assert_eq!(mgr.max_concurrent_tasks(), 1, "0 并发归一为 1");
assert!(matches!(mgr.timeout_policy(), TimeoutPolicy::None));
let mgr = OffloadManager::with_policy(3, TimeoutPolicy::Warn(Duration::from_millis(10)));
assert_eq!(mgr.max_concurrent_tasks(), 3);
assert!(matches!(mgr.timeout_policy(), TimeoutPolicy::Warn(_)));
}
#[tokio::test]
async fn none_policy_runs_task_without_timeout_wrap() {
let mgr = OffloadManager::with_policy(1, TimeoutPolicy::None);
assert!(mgr.spawn("plain", async { }));
assert_eq!(mgr.wait_all(Duration::from_secs(2)).await, 1);
}
#[tokio::test]
async fn warn_policy_records_slow_but_uninterrupted_task() {
let mgr = OffloadManager::with_policy(1, TimeoutPolicy::Warn(Duration::from_millis(10)));
assert!(mgr.spawn("slow-but-finishes", async {
tokio::time::sleep(Duration::from_millis(40)).await;
}));
let completed = mgr.wait_all(Duration::from_secs(2)).await;
assert_eq!(completed, 1, "warn 策略不打断任务,任务自然跑完");
}
}