Bedrock.DataPlane.Demux.PersistenceQueue (bedrock v0.5.2)

View Source

Bounded queue primitives for async shard persistence.

This queue tracks pending, in-flight, and retry-scheduled entries. It is intentionally side-effect free aside from telemetry emission.

Summary

Types

counts()

@type counts() :: %{
  pending: non_neg_integer(),
  scheduled: non_neg_integer(),
  in_flight: non_neg_integer(),
  lag: non_neg_integer()
}

entry()

@type entry() :: %{
  payload: payload(),
  attempt: non_neg_integer(),
  enqueued_at_ms: non_neg_integer()
}

payload()

@type payload() :: term()

scheduled_entry()

@type scheduled_entry() :: %{
  payload: payload(),
  attempt: non_neg_integer(),
  enqueued_at_ms: non_neg_integer(),
  due_at_ms: non_neg_integer()
}

t()

@type t() :: %Bedrock.DataPlane.Demux.PersistenceQueue{
  capacity: pos_integer(),
  in_flight: %{required(token()) => entry()},
  max_retries: non_neg_integer(),
  next_token: token(),
  pending: :queue.queue(entry()),
  retry_base_backoff_ms: pos_integer(),
  scheduled: [scheduled_entry()]
}

token()

@type token() :: pos_integer()

Functions

ack(queue, token)

@spec ack(t(), token()) :: {:ok, t()} | {:error, :unknown_token, t()}

counts(queue)

@spec counts(t()) :: counts()

dequeue(queue, opts \\ [])

@spec dequeue(
  t(),
  keyword()
) :: {:ok, token(), payload(), t()} | :empty

enqueue(queue, payload, opts \\ [])

@spec enqueue(t(), payload(), keyword()) :: {:ok, t()} | {:error, :full, t()}

lag(queue)

@spec lag(t()) :: non_neg_integer()

nack(queue, token, reason, opts \\ [])

@spec nack(t(), token(), term(), keyword()) ::
  {:rescheduled, t()} | {:dropped, t()} | {:error, :unknown_token, t()}

new(opts \\ [])

@spec new(keyword()) :: t()

next_retry_due_ms(arg1)

@spec next_retry_due_ms(t()) :: non_neg_integer() | nil