Bedrock.Internal.WaitingList (bedrock v0.6.0)

View Source

Unified waiting list for version-based out-of-order request handling.

Supports both single-waiter (Resolver) and multi-waiter (LongPulls) patterns using a map of deadline-sorted lists.

Structure: %{version => [{deadline, reply_fn, data}, ...]} Lists are sorted by deadline (earliest first).

Summary

Functions

Expire entries past their deadline from waiting list. Returns {new_map, expired_entries}.

Expire entries past their deadline from waiting list using provided time function. Returns {new_map, expired_entries}.

Find first entry for a version without removing it. Returns entry or nil.

Add entry to waiting list with deadline. Lists are kept sorted by deadline (earliest first). Returns {new_map, timeout_for_next_deadline}.

Add entry to waiting list with deadline using provided time function. Lists are kept sorted by deadline (earliest first). Returns {new_map, timeout_for_next_deadline}.

Calculate timeout for next deadline in waiting list.

Calculate timeout for next deadline in waiting list using provided time function.

Remove first entry for a version from waiting list. For single-waiter patterns (Resolver). Returns {new_map, removed_entry | nil}.

Remove all entries for a version from waiting list. For multi-waiter patterns (LongPulls). Returns {new_map, removed_entries}.

Remove all entries where version is less than the given threshold. For range-match patterns (Demux/ShardServer long-pull). Returns {new_map, removed_entries} where entries are flattened across all matching versions.

Types

entry()

@type entry() :: {deadline :: integer(), reply_fn :: reply_fn(), data :: any()}

reply_fn()

@type reply_fn() :: (any() -> :ok)

t()

@type t() :: %{required(version()) => [entry()]}

timeout_ms()

@type timeout_ms() :: non_neg_integer()

version()

@type version() :: Bedrock.version()

Functions

expire(map)

@spec expire(t()) :: {t(), [entry()]}

Expire entries past their deadline from waiting list. Returns {new_map, expired_entries}.

expire(map, time_fn)

@spec expire(t(), (-> integer())) :: {t(), [entry()]}

Expire entries past their deadline from waiting list using provided time function. Returns {new_map, expired_entries}.

find(map, version)

@spec find(t(), version()) :: entry() | nil

Find first entry for a version without removing it. Returns entry or nil.

insert(map, version, data, reply_fn, timeout_ms)

@spec insert(t(), version(), any(), reply_fn(), timeout_ms()) :: {t(), timeout()}

Add entry to waiting list with deadline. Lists are kept sorted by deadline (earliest first). Returns {new_map, timeout_for_next_deadline}.

insert(map, version, data, reply_fn, timeout_ms, time_fn)

@spec insert(t(), version(), any(), reply_fn(), timeout_ms(), (-> integer())) ::
  {t(), timeout()}

Add entry to waiting list with deadline using provided time function. Lists are kept sorted by deadline (earliest first). Returns {new_map, timeout_for_next_deadline}.

next_timeout(map)

@spec next_timeout(t()) :: timeout()

Calculate timeout for next deadline in waiting list.

next_timeout(map, time_fn)

@spec next_timeout(t(), (-> integer())) :: timeout()

Calculate timeout for next deadline in waiting list using provided time function.

remove(map, version)

@spec remove(t(), version()) :: {t(), entry() | nil}

Remove first entry for a version from waiting list. For single-waiter patterns (Resolver). Returns {new_map, removed_entry | nil}.

remove_all(map, version)

@spec remove_all(t(), version()) :: {t(), [entry()]}

Remove all entries for a version from waiting list. For multi-waiter patterns (LongPulls). Returns {new_map, removed_entries}.

remove_all_less_than(map, threshold_version)

@spec remove_all_less_than(t(), version()) :: {t(), [entry()]}

Remove all entries where version is less than the given threshold. For range-match patterns (Demux/ShardServer long-pull). Returns {new_map, removed_entries} where entries are flattened across all matching versions.

Example

iex> map = %{5 => [{...}], 10 => [{...}], 15 => [{...}]}
iex> {new_map, removed} = WaitingList.remove_all_less_than(map, 12)
iex> Map.keys(new_map)
[15]  # only version 15 remains

reply_to_expired(expired_entries, error_response \\ {:error, :waiting_timeout})

@spec reply_to_expired([entry()], any()) :: :ok

Reply to expired entries with error response.