Rheo.Inflight (rheo v1.0.0)

Copy Markdown View Source

Inflight-lease bookkeeping shared by Rheo.Group and Rheo.Producer.

An inflight set maps an opaque key (a task monitor reference for the Group, a lease_id for the Producer) to a metadata map that always contains :lease. All functions are pure; the owning process performs backend calls and telemetry around them.

Examples

iex> event = %Rheo.Event{id: "e1", stream: "s", partition: 0, sequence: 1,
...>   timestamp: ~U[2026-01-01 00:00:00.000Z], payload: %{}}
iex> lease = %Rheo.Lease{lease_id: "l1", stream: "s", group: "g", event_id: "e1",
...>   event: event, consumer_id: "c", attempt: 1,
...>   leased_at: ~U[2026-01-01 00:00:00.000Z], expires_at: ~U[2026-01-01 00:00:30.000Z]}
iex> inflight = Rheo.Inflight.put(Rheo.Inflight.new(), "l1", lease)
iex> {Rheo.Inflight.size(inflight), Rheo.Inflight.capacity(inflight, 10)}
{1, 9}
iex> {:ok, %{lease: ^lease}, inflight} = Rheo.Inflight.pop(inflight, "l1")
iex> Rheo.Inflight.size(inflight)
0

Summary

Types

Opaque key for an inflight entry (usually a monitor reference).

Metadata stored with an inflight lease (must include :lease).

Outcome of renewing one tracked lease.

t()

Map of inflight keys to lease metadata.

Functions

Slots left before max_demand unsettled leases are held.

Drops one key.

Drops many keys.

Fetches the lease tracked under key.

All tracked leases.

Empty inflight set.

Removes key, returning its metadata when it was tracked.

Tracks lease under key, merging meta into the stored metadata.

Renews every tracked lease with renew_fun.

Number of unsettled leases.

Replaces the lease for an existing key; unknown keys are ignored.

Types

key()

@type key() :: term()

Opaque key for an inflight entry (usually a monitor reference).

meta()

@type meta() :: %{:lease => Rheo.Lease.t(), optional(atom()) => term()}

Metadata stored with an inflight lease (must include :lease).

renew_result()

@type renew_result() :: {key(), Rheo.Lease.t(), :ok | Rheo.Settle.reason()}

Outcome of renewing one tracked lease.

t()

@type t() :: %{optional(key()) => meta()}

Map of inflight keys to lease metadata.

Functions

capacity(inflight, max_demand)

@spec capacity(t(), non_neg_integer()) :: non_neg_integer()

Slots left before max_demand unsettled leases are held.

delete(inflight, key)

@spec delete(t(), key()) :: t()

Drops one key.

drop(inflight, keys)

@spec drop(t(), [key()]) :: t()

Drops many keys.

fetch_lease(inflight, key)

@spec fetch_lease(t(), key()) :: {:ok, Rheo.Lease.t()} | :error

Fetches the lease tracked under key.

leases(inflight)

@spec leases(t()) :: [Rheo.Lease.t()]

All tracked leases.

new()

@spec new() :: t()

Empty inflight set.

pop(inflight, key)

@spec pop(t(), key()) :: {:ok, meta(), t()} | :error

Removes key, returning its metadata when it was tracked.

put(inflight, key, lease, meta \\ %{})

@spec put(t(), key(), Rheo.Lease.t(), map()) :: t()

Tracks lease under key, merging meta into the stored metadata.

renew_all(inflight, renew_fun)

@spec renew_all(t(), (Rheo.Lease.t() -> {:ok, Rheo.Lease.t()} | {:error, term()})) ::
  {t(), [renew_result()]}

Renews every tracked lease with renew_fun.

Renewed leases replace the stored lease. Leases the holder has lost (Rheo.Settle.lost?/1) are dropped so the backend can redeliver them; other failures keep the lease tracked for the next renewal round. Returns the updated set and one renew_result/0 per lease so the caller can emit telemetry.

size(inflight)

@spec size(t()) :: non_neg_integer()

Number of unsettled leases.

update_lease(inflight, key, lease)

@spec update_lease(t(), key(), Rheo.Lease.t()) :: t()

Replaces the lease for an existing key; unknown keys are ignored.