pub struct TimerWheel { /* private fields */ }Expand description
Timing Wheel Timer Manager
时间轮定时器管理器
Implementations§
Source§impl TimerWheel
impl TimerWheel
Sourcepub fn new(
config: WheelConfig,
batch_config: BatchConfig,
) -> Result<Self, TimerError>
pub fn new( config: WheelConfig, batch_config: BatchConfig, ) -> Result<Self, TimerError>
Create a new timer manager
§Parameters
config: Timing wheel configurationbatch_config: Batch operation configuration
创建新的定时器管理器
§参数
config: 时间轮配置batch_config: 批量操作配置
§Examples (示例)
use kestrel_timer::{TimerWheel, config::WheelConfig, TimerTask, config::BatchConfig};
use std::time::Duration;
#[tokio::main]
async fn main() {
let config = WheelConfig::builder()
.l0_tick_duration(Duration::from_millis(10))
.l0_slot_count(512)
.l1_tick_duration(Duration::from_secs(1))
.l1_slot_count(64)
.build()
.unwrap();
let timer = TimerWheel::new(config, BatchConfig::default()).unwrap();
// Use two-step API: allocate handle first, then register
// 使用两步 API:先分配 handle,再注册
let handle = timer.allocate_handle();
let task = TimerTask::new_oneshot(Duration::from_secs(1), None);
let _timer_handle = timer.register(handle, task).unwrap();
}Sourcepub fn with_defaults() -> Self
pub fn with_defaults() -> Self
Create timer manager with default configuration, hierarchical mode
- L0 layer tick duration: 10ms, slot count: 512
- L1 layer tick duration: 1s, slot count: 64
§Returns
Timer manager instance
使用默认配置创建定时器管理器,分层模式
- L0 层 tick 持续时间:10ms,槽数量:512
- L1 层 tick 持续时间:1s,槽数量:64
§返回值
定时器管理器实例
§Examples (示例)
use kestrel_timer::TimerWheel;
#[tokio::main]
async fn main() {
let timer = TimerWheel::with_defaults();
}Sourcepub fn create_service(&self, service_config: ServiceConfig) -> TimerService
pub fn create_service(&self, service_config: ServiceConfig) -> TimerService
Create TimerService bound to this timing wheel with default configuration
§Parameters
service_config: Service configuration
§Returns
TimerService instance bound to this timing wheel
创建绑定到此时间轮的 TimerService,使用默认配置
§参数
service_config: 服务配置
§返回值
绑定到此时间轮的 TimerService 实例
§Examples (示例)
use kestrel_timer::{TimerWheel, TimerService, CallbackWrapper, config::ServiceConfig};
use std::time::Duration;
#[tokio::main]
async fn main() {
let timer = TimerWheel::with_defaults();
let mut service = timer.create_service(ServiceConfig::default());
// Use two-step API to batch schedule timers through service
// 使用两步 API 通过服务批量调度定时器
// Step 1: Allocate handles
let handles = service.allocate_handles(5);
// Step 2: Create tasks
let tasks: Vec<_> = (0..5)
.map(|_| {
use kestrel_timer::TimerTask;
TimerTask::new_oneshot(Duration::from_millis(100), Some(CallbackWrapper::new(|| async {})))
})
.collect();
// Step 3: Register batch
service.register_batch(handles, tasks).unwrap();
// Receive timeout notifications
// 接收超时通知
let mut rx = service.take_receiver().unwrap();
while let Some(task_id) = rx.recv().await {
println!("Task {:?} completed", task_id);
}
}Sourcepub fn create_service_with_config(&self, config: ServiceConfig) -> TimerService
pub fn create_service_with_config(&self, config: ServiceConfig) -> TimerService
Create TimerService bound to this timing wheel with custom configuration
§Parameters
config: Service configuration
§Returns
TimerService instance bound to this timing wheel
创建绑定到此时间轮的 TimerService,使用自定义配置
§参数
config: 服务配置
§返回值
绑定到此时间轮的 TimerService 实例
§Examples (示例)
use kestrel_timer::{TimerWheel, config::ServiceConfig, TimerTask};
use std::num::NonZeroUsize;
#[tokio::main]
async fn main() {
let timer = TimerWheel::with_defaults();
let config = ServiceConfig::builder()
.command_channel_capacity(NonZeroUsize::new(1024).unwrap())
.timeout_channel_capacity(NonZeroUsize::new(2000).unwrap())
.build();
let service = timer.create_service_with_config(config);
}Sourcepub fn allocate_handle(&self) -> TaskHandle
pub fn allocate_handle(&self) -> TaskHandle
Sourcepub fn allocate_handles(&self, count: usize) -> Vec<TaskHandle>
pub fn allocate_handles(&self, count: usize) -> Vec<TaskHandle>
Sourcepub fn register(
&self,
handle: TaskHandle,
task: TimerTask,
) -> Result<TimerHandleWithCompletion, TimerError>
pub fn register( &self, handle: TaskHandle, task: TimerTask, ) -> Result<TimerHandleWithCompletion, TimerError>
Register timer task to timing wheel (registration phase)
§Parameters
task: Task created viacreate_task()
§Returns
Return Ok with a timer handle and completion receiver. Returns
Err(TimerError::WrongWheel) when the handle belongs to another wheel,
or Err(TimerError::Shutdown) when the wheel is closed.
注册定时器任务到时间轮 (注册阶段)
§参数
task: 通过create_task()创建的任务
§返回值
成功时返回包含完成通知接收器的定时器句柄;handle 属于其他时间轮时
返回 Err(TimerError::WrongWheel);时间轮关闭时返回
Err(TimerError::Shutdown)。
§Examples (示例)
use kestrel_timer::{TimerWheel, TimerTask, CallbackWrapper};
use std::time::Duration;
#[tokio::main]
async fn main() {
let timer = TimerWheel::with_defaults();
// Step 1: Allocate handle
let allocated_handle = timer.allocate_handle();
let task_id = allocated_handle.task_id();
// Step 2: Create task
let task = TimerTask::new_oneshot(Duration::from_secs(1), Some(CallbackWrapper::new(|| async {
println!("Timer fired!");
})));
// Step 3: Register task
let handle = timer.register(allocated_handle, task).unwrap();
// Wait for timer completion
// 等待定时器完成
use kestrel_timer::CompletionReceiver;
let (rx, _handle) = handle.into_parts();
match rx {
CompletionReceiver::OneShot(receiver) => {
receiver.recv().await.unwrap();
},
_ => {}
}
}Sourcepub fn register_batch(
&self,
handles: Vec<TaskHandle>,
tasks: Vec<TimerTask>,
) -> Result<BatchHandleWithCompletion, TimerError>
pub fn register_batch( &self, handles: Vec<TaskHandle>, tasks: Vec<TimerTask>, ) -> Result<BatchHandleWithCompletion, TimerError>
Batch register timer tasks to timing wheel (registration phase)
§Parameters
handles: Pre-allocated handles for taskstasks: List of timer tasks
§Returns
Ok(BatchHandleWithCompletion)if all tasks are successfully registeredErr(TimerError::BatchLengthMismatch)if handles and tasks lengths don’t matchErr(TimerError::WrongWheel)if any handle belongs to another wheelErr(TimerError::Shutdown)if the wheel is closed
批量注册定时器任务到时间轮 (注册阶段)
§参数
handles: 任务的预分配 handlestasks: 定时器任务列表
§返回值
Ok(BatchHandleWithCompletion)如果所有任务成功注册Err(TimerError::BatchLengthMismatch)如果 handles 和 tasks 长度不匹配Err(TimerError::WrongWheel)如果任一 handle 属于其他时间轮Err(TimerError::Shutdown)如果时间轮已关闭
§Examples (示例)
use kestrel_timer::{TimerWheel, TimerTask};
use std::time::Duration;
#[tokio::main]
async fn main() {
let timer = TimerWheel::with_defaults();
// Step 1: Allocate handles
let handles = timer.allocate_handles(3);
// Step 2: Create tasks
let tasks: Vec<_> = (0..3)
.map(|_| TimerTask::new_oneshot(Duration::from_secs(1), None))
.collect();
// Step 3: Batch register
let batch = timer.register_batch(handles, tasks)
.expect("register_batch should succeed");
println!("Registered {} timers", batch.len());
}Sourcepub fn cancel(&self, task_id: TaskId) -> Result<bool, TimerError>
pub fn cancel(&self, task_id: TaskId) -> Result<bool, TimerError>
Cancel timer
§Parameters
task_id: Task ID
§Returns
Returns Ok(true) when cancelled, Ok(false) when absent, or
Err(TimerError::WrongWheel) for an ID from another wheel.
取消定时器
§参数
task_id: 任务 ID
§返回值
任务存在且取消成功时返回 Ok(true),任务不存在时返回 Ok(false),
ID 属于其他时间轮时返回 Err(TimerError::WrongWheel)。
§Examples (示例)
use kestrel_timer::{TimerWheel, TimerTask, CallbackWrapper};
use std::time::Duration;
#[tokio::main]
async fn main() {
let timer = TimerWheel::with_defaults();
// Step 1: Allocate handle
let allocated_handle = timer.allocate_handle();
let task_id = allocated_handle.task_id();
// Step 2: Create and register task
let task = TimerTask::new_oneshot(Duration::from_secs(10), Some(CallbackWrapper::new(|| async {
println!("Timer fired!");
})));
let _handle = timer.register(allocated_handle, task).unwrap();
// Cancel task using task ID
// 使用任务 ID 取消任务
let cancelled = timer.cancel(task_id).unwrap();
println!("Canceled successfully: {}", cancelled);
}Sourcepub fn cancel_batch(&self, task_ids: &[TaskId]) -> Result<usize, TimerError>
pub fn cancel_batch(&self, task_ids: &[TaskId]) -> Result<usize, TimerError>
Batch cancel timers
§Parameters
task_ids: List of task IDs to cancel
§Returns
Number of successfully cancelled tasks, or
Err(TimerError::WrongWheel) if any ID belongs to another wheel.
批量取消定时器
§参数
task_ids: 要取消的任务 ID 列表
§返回值
成功取消的任务数量;如果任一 ID 属于其他时间轮则返回
Err(TimerError::WrongWheel),且不修改任何任务。
§Performance Advantages
- Batch processing reduces lock contention
- Internally optimized batch cancellation operation
§Examples (示例)
use kestrel_timer::{TimerWheel, TimerTask};
use std::time::Duration;
#[tokio::main]
async fn main() {
let timer = TimerWheel::with_defaults();
// Create multiple timers
// 创建多个定时器
let task1 = TimerTask::new_oneshot(Duration::from_secs(10), None);
let task2 = TimerTask::new_oneshot(Duration::from_secs(10), None);
let task3 = TimerTask::new_oneshot(Duration::from_secs(10), None);
// Allocate handles and get task IDs
let h1 = timer.allocate_handle();
let h2 = timer.allocate_handle();
let h3 = timer.allocate_handle();
let task_ids = vec![h1.task_id(), h2.task_id(), h3.task_id()];
let _h1 = timer.register(h1, task1).unwrap();
let _h2 = timer.register(h2, task2).unwrap();
let _h3 = timer.register(h3, task3).unwrap();
// Batch cancel
// 批量取消
let cancelled = timer.cancel_batch(&task_ids).unwrap();
println!("Canceled {} timers", cancelled);
}Sourcepub fn postpone(
&self,
task_id: TaskId,
new_delay: Duration,
callback: Option<CallbackWrapper>,
) -> Result<bool, TimerError>
pub fn postpone( &self, task_id: TaskId, new_delay: Duration, callback: Option<CallbackWrapper>, ) -> Result<bool, TimerError>
Postpone timer
§Parameters
task_id: Task ID to postponenew_delay: New delay duration, recalculated from current timecallback: New callback function, passNoneto keep original callback, passSometo replace with new callback
§Returns
Returns Ok(true) when postponed, Ok(false) when absent, or
Err(TimerError::WrongWheel) for an ID from another wheel.
推迟定时器
§参数
task_id: 要推迟的任务 IDnew_delay: 新的延迟时间,从当前时间重新计算callback: 新的回调函数,传递None保持原始回调,传递Some替换为新的回调
§返回值
任务存在且延期成功时返回 Ok(true),任务不存在时返回 Ok(false),
ID 属于其他时间轮时返回 Err(TimerError::WrongWheel)。
§Note
- Task ID remains unchanged after postponement
- Original completion_receiver remains valid
§注意
- 任务 ID 在推迟后保持不变
- 原始 completion_receiver 保持有效
§Examples (示例)
§Keep original callback (保持原始回调)
use kestrel_timer::{TimerWheel, TimerTask, CallbackWrapper};
use std::time::Duration;
#[tokio::main]
async fn main() {
let timer = TimerWheel::with_defaults();
// Allocate handle first
let allocated_handle = timer.allocate_handle();
let task_id = allocated_handle.task_id();
let task = TimerTask::new_oneshot(Duration::from_secs(5), Some(CallbackWrapper::new(|| async {
println!("Timer fired!");
})));
let _handle = timer.register(allocated_handle, task).unwrap();
// Postpone to 10 seconds after triggering, and keep original callback
// 推迟到 10 秒后触发,并保持原始回调
let success = timer
.postpone(task_id, Duration::from_secs(10), None)
.unwrap();
println!("Postponed successfully: {}", success);
}§Replace with new callback (替换为新的回调)
use kestrel_timer::{TimerWheel, TimerTask, CallbackWrapper};
use std::time::Duration;
#[tokio::main]
async fn main() {
let timer = TimerWheel::with_defaults();
// Allocate handle first
let allocated_handle = timer.allocate_handle();
let task_id = allocated_handle.task_id();
let task = TimerTask::new_oneshot(Duration::from_secs(5), Some(CallbackWrapper::new(|| async {
println!("Original callback!");
})));
let _handle = timer.register(allocated_handle, task).unwrap();
// Postpone to 10 seconds after triggering, and replace with new callback
// 推迟到 10 秒后触发,并替换为新的回调
let success = timer.postpone(task_id, Duration::from_secs(10), Some(CallbackWrapper::new(|| async {
println!("New callback!");
}))).unwrap();
println!("Postponed successfully: {}", success);
}Sourcepub fn postpone_batch(
&self,
updates: Vec<(TaskId, Duration)>,
) -> Result<usize, TimerError>
pub fn postpone_batch( &self, updates: Vec<(TaskId, Duration)>, ) -> Result<usize, TimerError>
Batch postpone timers (keep original callbacks)
§Parameters
updates: List of tuples of (task ID, new delay)
§Returns
Number of successfully postponed tasks, or
Err(TimerError::WrongWheel) if any ID belongs to another wheel.
批量推迟定时器 (保持原始回调)
§参数
updates: (任务 ID, 新延迟) 元组列表
§返回值
成功推迟的任务数量;如果任一 ID 属于其他时间轮则返回
Err(TimerError::WrongWheel),且不修改任何任务。
§Note
- This method keeps all tasks’ original callbacks unchanged
- Use
postpone_batch_with_callbacksif you need to replace callbacks
§注意
- 此方法保持所有任务的原始回调不变
- 如果需要替换回调,请使用
postpone_batch_with_callbacks
§Performance Advantages
- Batch processing reduces lock contention
- Internally optimized batch postponement operation
§Examples (示例)
use kestrel_timer::{TimerWheel, TimerTask, CallbackWrapper};
use std::time::Duration;
#[tokio::main]
async fn main() {
let timer = TimerWheel::with_defaults();
// Create multiple tasks with callbacks
// 创建多个带有回调的任务
let task1 = TimerTask::new_oneshot(Duration::from_secs(5), Some(CallbackWrapper::new(|| async {
println!("Task 1 fired!");
})));
let task2 = TimerTask::new_oneshot(Duration::from_secs(5), Some(CallbackWrapper::new(|| async {
println!("Task 2 fired!");
})));
let task3 = TimerTask::new_oneshot(Duration::from_secs(5), Some(CallbackWrapper::new(|| async {
println!("Task 3 fired!");
})));
// Allocate handles and register
let h1 = timer.allocate_handle();
let h2 = timer.allocate_handle();
let h3 = timer.allocate_handle();
let task_ids = vec![
(h1.task_id(), Duration::from_secs(10)),
(h2.task_id(), Duration::from_secs(15)),
(h3.task_id(), Duration::from_secs(20)),
];
timer.register(h1, task1).unwrap();
timer.register(h2, task2).unwrap();
timer.register(h3, task3).unwrap();
// Batch postpone (keep original callbacks)
// 批量推迟 (保持原始回调)
let postponed = timer.postpone_batch(task_ids).unwrap();
println!("Postponed {} timers", postponed);
}Sourcepub fn postpone_batch_with_callbacks(
&self,
updates: Vec<(TaskId, Duration, Option<CallbackWrapper>)>,
) -> Result<usize, TimerError>
pub fn postpone_batch_with_callbacks( &self, updates: Vec<(TaskId, Duration, Option<CallbackWrapper>)>, ) -> Result<usize, TimerError>
Batch postpone timers (replace callbacks)
§Parameters
updates: List of tuples of (task ID, new delay, new callback)
§Returns
Number of successfully postponed tasks, or
Err(TimerError::WrongWheel) if any ID belongs to another wheel.
批量推迟定时器 (替换回调)
§参数
updates: (任务 ID, 新延迟, 新回调) 元组列表
§返回值
成功推迟的任务数量;如果任一 ID 属于其他时间轮则返回
Err(TimerError::WrongWheel),且不修改任何任务。
§Performance Advantages
- Batch processing reduces lock contention
- Internally optimized batch postponement operation
§Examples (示例)
use kestrel_timer::{TimerWheel, TimerTask, CallbackWrapper};
use std::time::Duration;
use std::sync::Arc;
use std::sync::atomic::{AtomicU32, Ordering};
#[tokio::main]
async fn main() {
let timer = TimerWheel::with_defaults();
let counter = Arc::new(AtomicU32::new(0));
// Create multiple timers
// 创建多个定时器
let task1 = TimerTask::new_oneshot(Duration::from_secs(5), None);
let task2 = TimerTask::new_oneshot(Duration::from_secs(5), None);
// Allocate handles first
let h1 = timer.allocate_handle();
let h2 = timer.allocate_handle();
let id1 = h1.task_id();
let id2 = h2.task_id();
timer.register(h1, task1).unwrap();
timer.register(h2, task2).unwrap();
// Batch postpone and replace callbacks
// 批量推迟并替换回调
let updates: Vec<_> = vec![id1, id2]
.into_iter()
.map(|id| {
let counter = Arc::clone(&counter);
(id, Duration::from_secs(10), Some(CallbackWrapper::new(move || {
let counter = Arc::clone(&counter);
async move { counter.fetch_add(1, Ordering::SeqCst); }
})))
})
.collect();
let postponed = timer.postpone_batch_with_callbacks(updates).unwrap();
println!("Postponed {} timers", postponed);
}Sourcepub async fn shutdown(self)
pub async fn shutdown(self)
Graceful shutdown of TimerWheel
Cancels every task still owned by this wheel, sends Cancelled to
retained completion receivers, and stops the background driver. The
wheel rejects all later registrations.
优雅关闭 TimerWheel
取消此时间轮中仍存在的所有任务,向仍保留的完成接收器发送
Cancelled,并停止后台驱动。关闭后不再接受新的注册。
§Examples (示例)
let timer = TimerWheel::with_defaults();
// Use timer... (使用定时器...)
timer.shutdown().await;Trait Implementations§
Source§impl Drop for TimerWheel
Close the wheel and abort the background tick task when TimerWheel is dropped
impl Drop for TimerWheel
Close the wheel and abort the background tick task when TimerWheel is dropped
当 TimerWheel 被销毁时关闭时间轮并中止后台 tick 任务