Data-HierTimingWheel-Shared
view release on metacpan or search on metacpan
NAME
Data::HierTimingWheel::Shared - shared-memory hierarchical timing wheel
(O(1) timers at any delay)
SYNOPSIS
use Data::HierTimingWheel::Shared;
# 4 levels of 256 slots -> schedules any delay up to 256**4 - 1 ticks
my $tw = Data::HierTimingWheel::Shared->new(undef, 256, 4, 100_000);
my $id = $tw->add(30, $job_id); # fire in 30 ticks
my $id2 = $tw->add(5_000_000, $late); # ...or millions of ticks away, still O(1)
$tw->cancel($id); # cancel before it fires
# advance the clock; each call returns the payloads that came due
for (1 .. $ticks_elapsed) {
my @due = $tw->advance(1);
handle($_) for @due;
}
# share the wheel across processes via a backing file
my $shared = Data::HierTimingWheel::Shared->new("/tmp/timers.htw", 256, 4, 100_000);
DESCRIPTION
A hierarchical timing wheel in shared memory (the Varghese-Lauck design
behind Linux kernel and Kafka/Netty timers): a timer scheduler where
scheduling and cancelling are O(1) at any delay, from one tick to
billions, in a fixed amount of memory. It is the multi-level
generalisation of Data::TimingWheel::Shared: where the single-level wheel
revisits a far-future timer once per rotation, this one parks it in a
coarse level and only touches it as its time approaches.
Time advances in integer ticks. There are "num_levels" cascading wheels of
"num_slots" (= S) buckets each; a level-"k" slot spans "S**k" ticks, so
the whole structure schedules any delay in "[1, S**num_levels)". A timer
is placed in the lowest level whose range covers its delay. On each tick,
level 0 fires the timers in its current slot; when level 0 completes a
rotation, the next level's current slot cascades down -- its timers are
redistributed into finer levels by their remaining delay -- recursively up
the levels. A timer with delay "D" fires at exactly tick "D", just like
the single-level wheel, but a far-future timer costs O(1) instead of one
visit per rotation.
Each timer carries an arbitrary 64-bit payload (e.g. a job id) returned
when it fires. Timers live in a fixed pool of "capacity" slots; scheduling
beyond it, or with a delay at or beyond "S**num_levels", croaks.
Because the wheels live in a shared mapping, several processes schedule
into and advance one clock: any process that opens the same backing file,
inherits the anonymous mapping across "fork", or reopens a passed memfd
shares the same timers. A write-preferring futex rwlock with dead-process
recovery guards mutation. Linux-only. Requires 64-bit Perl.
METHODS
Constructors
my $tw = Data::HierTimingWheel::Shared->new($path, $num_slots, $num_levels, $capacity, $mode);
my $tw = Data::HierTimingWheel::Shared->new(undef, $num_slots, $num_levels, $capacity);
my $tw = Data::HierTimingWheel::Shared->new_memfd($name, $num_slots, $num_levels, $capacity);
my $tw = Data::HierTimingWheel::Shared->new_from_fd($fd);
$num_slots (S, 2..2^16, default 256) is the number of buckets per level,
and $num_levels (L, 1..16, default 4) the number of cascading wheels;
together they set the maximum schedulable delay to "S**L - 1" ticks
("num_slots ** num_levels" must fit in 64 bits, else the constructor
croaks). $capacity is the maximum number of concurrent timers (1..2^24).
Memory is "num_levels * num_slots * 4 + capacity * 32" bytes plus a fixed
header. Choose S and L so that "S**L" exceeds your longest delay -- e.g.
"256, 4" covers ~4.3 billion ticks, "64, 8" covers ~281 trillion. When
reopening an existing file or memfd the stored geometry wins and the
caller's arguments do not resize it -- but they are still range-checked,
so an out-of-range value croaks. An optional file mode may be passed as
the last argument to "new" (e.g. 0660) for cross-user sharing; it defaults
to 0600 (owner-only).
Scheduling
my $id = $tw->add($delay, $payload); # returns a timer id
my $id = $tw->schedule($delay, $payload); # alias for add
my $ok = $tw->cancel($id); # 1 if cancelled, 0 if already fired/invalid
"add" schedules a timer to fire $delay ticks from now (a delay below 1 is
treated as 1) carrying the integer $payload, and returns a timer id; it
croaks if the timer pool is full or $delay is at or beyond the wheel's
range ("max_delay + 1"). "cancel" removes a still-pending timer by its id,
returning 1 if it was cancelled or 0 if it had already fired or the id is
not active.
Advancing the clock
my @due = $tw->advance($ticks); # advance by $ticks (default 1)
my @due = $tw->advance; # advance by one tick
"advance" moves the wheel forward by $ticks ticks (default 1) and returns
the list of payloads of every timer that came due during those ticks, in
fire order. Timers that fire are removed automatically. Cost is O(ticks +
fired) amortised; cascades happen only when a level rolls over.
Introspection and lifecycle
$tw->now; # absolute tick count since creation (or last clear)
$tw->count; # number of pending timers
$tw->num_slots; # slots per level (S)
$tw->num_levels; # number of levels (L)
$tw->max_delay; # largest schedulable delay (S**L - 1)
$tw->capacity; # maximum concurrent timers
$tw->clear; # cancel all timers and reset the clock to 0
$tw->stats; # { now, count, num_slots, num_levels, max_delay, capacity, ops, mmap_size }
$tw->path; $tw->memfd; $tw->sync; $tw->unlink;
"clear" cancels every timer and resets the tick counter. "sync" flushes
the mapping to its backing store (a no-op for anonymous and memfd wheels);
"unlink" removes the backing file (also callable as
"Class->unlink($path)"); "path" returns the backing path ("undef" for
anonymous, memfd, or fd-reopened wheels) and "memfd" the backing
descriptor.
SHARING ACROSS PROCESSES
The wheels live in a shared mapping, shared the same three ways as the
rest of the family: a backing file, an anonymous mapping inherited across
"fork", or a memfd passed to an unrelated process and reopened with
new_from_fd($fd). The descriptor you pass is duplicated
("F_DUPFD_CLOEXEC"), so it stays yours to close and closing it does not
disturb the handle. Any process can schedule timers; typically one process
owns advancing the clock and dispatches the fired payloads, while others
schedule and cancel. The tick counter is shared, so all processes agree on
"now".
SECURITY
Backing files are created with mode 0600 (owner-only) by default; pass an
explicit octal mode (e.g. 0660) as the last argument to "new" for
cross-user sharing. The file is opened with "O_NOFOLLOW" and "O_EXCL", and
the header is validated on attach. Any process granted write access is
trusted not to corrupt the mapping.
CRASH SAFETY
Mutation is guarded by a futex-based write-preferring rwlock with
PID-encoded ownership and dead-owner recovery. Scheduling, cancelling, and
each tick of an advance are short bounded list operations, so a crash
leaves the wheel consistent up to the last completed operation.
Limitation: PID reuse is not detected (very unlikely in practice).
Reader-slot exhaustion (slotless readers): dead-process recovery
attributes a crashed lock holder's contribution through its reader-slot.
The slot table holds 1024 entries (one per concurrent reader process). If
more than that many reader processes share one mapping at once, a reader
that cannot claim a slot proceeds "slotless" -- it still takes the read
lock but leaves no per-process record. If such a slotless reader is then
killed while holding the read lock, its share of the lock cannot be
attributed to a dead process, so writer recovery cannot reclaim it and
writers may block until the mapping is recreated. Reaching this needs more
than 1024 concurrent reader processes on one mapping plus a crash in the
brief read-lock window; the dead-process slot reclaim keeps the table from
filling with stale entries, so in practice it is very unlikely.
An interrupted create is recovered too. A creator killed after the backing
file is sized but before its header is committed leaves a full-size,
all-zero file. "new" re-initializes such a file automatically, but only
( run in 0.718 second using v1.01-cache-2.11-cpan-e7c6538aa59 )