Skip to main content

TimerWheel

Struct TimerWheel 

Source
pub struct TimerWheel { /* private fields */ }
Expand description

Timing Wheel Timer Manager

时间轮定时器管理器

Implementations§

Source§

impl TimerWheel

Source

pub fn new( config: WheelConfig, batch_config: BatchConfig, ) -> Result<Self, TimerError>

Create a new timer manager

§Parameters
  • config: Timing wheel configuration
  • batch_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();
}
Source

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();
}
Source

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);
    }
}
Source

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);
}
Source

pub fn allocate_handle(&self) -> TaskHandle

Allocate a handle from DeferredMap

§Returns

A unique handle for later insertion

§返回值

用于后续插入的唯一 handle

Source

pub fn allocate_handles(&self, count: usize) -> Vec<TaskHandle>

Batch allocate handles from DeferredMap

§Parameters
  • count: Number of handles to allocate
§Returns

Vector of unique handles for later batch insertion

§参数
  • count: 要分配的 handle 数量
§返回值

用于后续批量插入的唯一 handles 向量

Source

pub fn register( &self, handle: TaskHandle, task: TimerTask, ) -> Result<TimerHandleWithCompletion, TimerError>

Register timer task to timing wheel (registration phase)

§Parameters
  • task: Task created via create_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();
        },
        _ => {}
    }
}
Source

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 tasks
  • tasks: List of timer tasks
§Returns
  • Ok(BatchHandleWithCompletion) if all tasks are successfully registered
  • Err(TimerError::BatchLengthMismatch) if handles and tasks lengths don’t match
  • Err(TimerError::WrongWheel) if any handle belongs to another wheel
  • Err(TimerError::Shutdown) if the wheel is closed

批量注册定时器任务到时间轮 (注册阶段)

§参数
  • handles: 任务的预分配 handles
  • tasks: 定时器任务列表
§返回值
  • 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());
}
Source

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);
}
Source

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);
}
Source

pub fn postpone( &self, task_id: TaskId, new_delay: Duration, callback: Option<CallbackWrapper>, ) -> Result<bool, TimerError>

Postpone timer

§Parameters
  • task_id: Task ID to postpone
  • new_delay: New delay duration, recalculated from current time
  • callback: New callback function, pass None to keep original callback, pass Some to 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: 要推迟的任务 ID
  • new_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);
}
Source

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_callbacks if 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);
}
Source

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);
}
Source

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

当 TimerWheel 被销毁时关闭时间轮并中止后台 tick 任务

Source§

fn drop(&mut self)

Executes the destructor for this type. Read more
Source§

fn pin_drop(self: Pin<&mut Self>)

🔬This is a nightly-only experimental API. (pin_ergonomics)
Execute the destructor for this type, but different to Drop::drop, it requires self to be pinned. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.