Minga.Buffer.State.Swap (Minga v0.1.0)

Copy Markdown View Source

Buffer-owned admission state for crash-recovery swap writes.

The state bounds work to one active preparation and one latest pending snapshot. Its generation advances on every admission and invalidation, so a prepared older snapshot cannot publish after newer work becomes authoritative.

Summary

Types

The active preparation worker: snapshot, worker pid, and monitor.

Monotonic swap admission generation.

An admitted buffer snapshot: generation, source path, and content.

t()

Function that starts the debounce timer for an admitted swap.

Functions

Admits a snapshot, starting it immediately or replacing the latest pending work.

Returns the configured storage backend.

Returns backend options for a generation, including the swap directory.

Completes matching active work and returns the latest pending snapshot, if any.

Completes matching active work identified by its process monitor.

Returns true when this buffer has a configured swap directory.

Consumes the matching timer token and rejects stale timer messages.

Invalidates every admitted write and returns resources the Buffer process must revoke.

Builds swap admission state from Buffer start options.

Checks whether a worker's prepared result is still the newest admitted snapshot.

Records the current debounce timer and its stale-message token.

Clears the current timer and returns its timer reference for cancellation.

Returns the configured debounce timer starter.

Records the monitored worker that is preparing an admitted snapshot.

Types

active_work()

@type active_work() :: {snapshot(), pid(), reference()}

The active preparation worker: snapshot, worker pid, and monitor.

generation()

@type generation() :: non_neg_integer()

Monotonic swap admission generation.

snapshot()

@type snapshot() :: {generation(), String.t(), binary()}

An admitted buffer snapshot: generation, source path, and content.

t()

@type t() :: %Minga.Buffer.State.Swap{
  active: active_work() | nil,
  backend: module() | nil,
  backend_options: keyword(),
  directory: String.t() | nil,
  generation: generation(),
  pending: snapshot() | nil,
  timer: reference() | nil,
  timer_start: timer_start(),
  timer_token: reference() | nil
}

timer_start()

@type timer_start() :: (pid(), term(), non_neg_integer() -> reference())

Function that starts the debounce timer for an admitted swap.

Functions

admit(state, path, content)

@spec admit(t(), String.t(), binary()) :: {:start, snapshot(), t()} | {:pending, t()}

Admits a snapshot, starting it immediately or replacing the latest pending work.

backend(swap)

@spec backend(t()) :: module()

Returns the configured storage backend.

backend_options(state, generation)

@spec backend_options(t(), generation()) :: keyword()

Returns backend options for a generation, including the swap directory.

complete(state, pid, generation)

@spec complete(t(), pid(), generation()) ::
  {:ok, reference(), snapshot() | nil, t()} | :unknown

Completes matching active work and returns the latest pending snapshot, if any.

complete_monitor(state, pid, monitor)

@spec complete_monitor(t(), pid(), reference()) ::
  {:ok, generation(), snapshot() | nil, t()} | :unknown

Completes matching active work identified by its process monitor.

configured?(swap)

@spec configured?(t()) :: boolean()

Returns true when this buffer has a configured swap directory.

consume_timer(state, token)

@spec consume_timer(t(), reference()) :: {:ok, t()} | :stale

Consumes the matching timer token and rejects stale timer messages.

invalidate(state)

@spec invalidate(t()) :: {reference() | nil, active_work() | nil, t()}

Invalidates every admitted write and returns resources the Buffer process must revoke.

new(opts)

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

Builds swap admission state from Buffer start options.

publication_status(swap, pid, generation)

@spec publication_status(t(), pid(), generation()) :: :current | :obsolete | :unknown

Checks whether a worker's prepared result is still the newest admitted snapshot.

schedule(state, timer, token)

@spec schedule(t(), reference(), reference()) :: t()

Records the current debounce timer and its stale-message token.

take_timer(state)

@spec take_timer(t()) :: {reference() | nil, t()}

Clears the current timer and returns its timer reference for cancellation.

timer_start(swap)

@spec timer_start(t()) :: timer_start()

Returns the configured debounce timer starter.

worker_started(state, snapshot, pid, monitor)

@spec worker_started(t(), snapshot(), pid(), reference()) :: t()

Records the monitored worker that is preparing an admitted snapshot.