mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-11 05:30:19 +08:00
新增 crate: - aether-runtime: 服务运行时基础设施(并发门控、分布式并发、指标、队列、优雅关闭、tracing) - aether-cache: 通用 TTL 缓存与命名空间抽象 - aether-data: 数据访问层(PostgreSQL/Redis 后端、repository 模式) - aether-http: HTTP 客户端封装(重试、配置) - aether-testkit: 集成测试工具集(gateway/executor/hub/proxy fixture、等待、负载测试) gateway 扩展: - 引入 audit 模块(shadow 执行审计、决策链路追踪、请求审计 bundle) - 引入 cache 模块(AuthContext 缓存、direct-plan bypass 缓存) - 引入 data 模块(auth/candidates/config/usage/video_tasks 数据访问) - 集成 ConcurrencyGate/DistributedConcurrencyGate 请求门控 - 新增本地 auth 拒绝、过载响应构建器 - 补充 control/auth_cache/video/concurrency 集成测试 aether-proxy 扩展: - AppState 集成 stream_gate / distributed_stream_gate 并发门控 - 新增 ProxyAdmissionError 及准入拒绝流程 - stream_handler 补充门控饱和/不可用场景测试 - 配置与注册客户端逻辑完善 aether-hub 扩展: - main.rs 引入运行时初始化、指标端点、健康检查 - local_relay 重构为 lib.rs 暴露公共接口
51 lines
1.2 KiB
Rust
51 lines
1.2 KiB
Rust
use std::future::Future;
|
|
use std::time::Duration;
|
|
|
|
pub async fn wait_until<F, Fut>(
|
|
timeout: Duration,
|
|
poll_interval: Duration,
|
|
mut predicate: F,
|
|
) -> bool
|
|
where
|
|
F: FnMut() -> Fut,
|
|
Fut: Future<Output = bool>,
|
|
{
|
|
let deadline = tokio::time::Instant::now() + timeout;
|
|
loop {
|
|
if predicate().await {
|
|
return true;
|
|
}
|
|
if tokio::time::Instant::now() >= deadline {
|
|
return false;
|
|
}
|
|
tokio::time::sleep(poll_interval).await;
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use std::sync::atomic::{AtomicBool, Ordering};
|
|
use std::sync::Arc;
|
|
use std::time::Duration;
|
|
|
|
use super::wait_until;
|
|
|
|
#[tokio::test]
|
|
async fn returns_true_when_predicate_eventually_passes() {
|
|
let flag = Arc::new(AtomicBool::new(false));
|
|
let background = flag.clone();
|
|
tokio::spawn(async move {
|
|
tokio::time::sleep(Duration::from_millis(10)).await;
|
|
background.store(true, Ordering::Release);
|
|
});
|
|
|
|
let ready = wait_until(Duration::from_millis(100), Duration::from_millis(5), || {
|
|
let flag = flag.clone();
|
|
async move { flag.load(Ordering::Acquire) }
|
|
})
|
|
.await;
|
|
|
|
assert!(ready);
|
|
}
|
|
}
|