Packages

A distributed cache library for Elixir with pluggable backends, topologies, and near-cache support.

Current section

Files

Jump to
vela_cache lib vela topology partitioned.ex
Raw

lib/vela/topology/partitioned.ex

defmodule Vela.Topology.Partitioned do
@moduledoc """
Partitioned topology: each key lives on exactly one node.
Uses a consistent hash ring to determine which node owns a key.
If the local node is the owner, the operation goes directly to the
local backend. If a remote node owns the key, the operation is
forwarded via `:erpc`.
Reads require a network hop when the key lives on another node.
Scales linearly — more nodes = more total capacity.
"""
@behaviour Vela.Topology
alias Vela.Cache.Entry
alias Vela.Distributed.Ring
alias Vela.Distributed.RPC
defstruct [:backend, :backend_state, :cache_name, :ring]
@impl true
def init(config) do
{:ok, backend_state} = config.backend.init(config)
# Build ring from currently connected nodes + self
nodes = [node() | Node.list()]
ring = Ring.new(nodes)
state = %__MODULE__{
backend: config.backend,
backend_state: backend_state,
cache_name: config.name,
ring: ring
}
{:ok, state}
end
@impl true
def get(%__MODULE__{} = state, key) do
case Ring.owner(state.ring, key) do
owner when owner == node() ->
state.backend.get(state.backend_state, key)
remote_node ->
RPC.call(remote_node, __MODULE__, :handle_remote_get, [state.cache_name, key])
|> unwrap_rpc()
end
end
@impl true
def put(%__MODULE__{} = state, %Entry{} = entry) do
case Ring.owner(state.ring, entry.key) do
owner when owner == node() ->
{:ok, new_bs} = state.backend.put(state.backend_state, entry)
{:ok, %{state | backend_state: new_bs}}
remote_node ->
case RPC.call(remote_node, __MODULE__, :handle_remote_put, [state.cache_name, entry]) do
{:ok, :ok} -> {:ok, state}
{:error, _} = err -> err
end
end
end
@impl true
def delete(%__MODULE__{} = state, key) do
case Ring.owner(state.ring, key) do
owner when owner == node() ->
{:ok, new_bs} = state.backend.delete(state.backend_state, key)
{:ok, %{state | backend_state: new_bs}}
remote_node ->
RPC.call(remote_node, __MODULE__, :handle_remote_delete, [state.cache_name, key])
{:ok, state}
end
end
@impl true
def get_many(%__MODULE__{} = state, keys) do
# Group keys by owning node
grouped = Enum.group_by(keys, &Ring.owner(state.ring, &1))
results =
Enum.reduce(grouped, %{}, fn {owner, node_keys}, acc ->
Map.merge(acc, fetch_from_node(state, owner, node_keys))
end)
{:ok, results}
end
@impl true
def put_many(%__MODULE__{} = state, entries) do
grouped = Enum.group_by(entries, &Ring.owner(state.ring, &1.key))
new_state =
Enum.reduce(grouped, state, fn {owner, node_entries}, acc ->
put_to_node(acc, owner, node_entries)
end)
{:ok, new_state}
end
defp fetch_from_node(state, owner, keys) do
if owner == node() do
{:ok, entries} = state.backend.get_many(state.backend_state, keys)
entries
else
case RPC.call(owner, __MODULE__, :handle_remote_get_many, [state.cache_name, keys]) do
{:ok, entries} -> entries
{:error, _} -> %{}
end
end
end
defp put_to_node(state, owner, entries) do
if owner == node() do
{:ok, new_bs} =
Enum.reduce(entries, {:ok, state.backend_state}, fn entry, {:ok, bs} ->
state.backend.put(bs, entry)
end)
%{state | backend_state: new_bs}
else
RPC.call(owner, __MODULE__, :handle_remote_put_many, [state.cache_name, entries])
state
end
end
@impl true
def flush(%__MODULE__{} = state) do
# Flush local backend
{:ok, new_bs} = state.backend.flush(state.backend_state)
# Tell all remote nodes to flush their partition
other_nodes = Node.list()
Enum.each(other_nodes, fn remote_node ->
RPC.cast(remote_node, __MODULE__, :handle_remote_flush, [state.cache_name])
end)
{:ok, %{state | backend_state: new_bs}}
end
@impl true
def invalidate_tag(%__MODULE__{} = state, tag) do
# Invalidate locally
{:ok, local_count, new_bs} = state.backend.delete_by_tag(state.backend_state, tag)
# Broadcast to all remote nodes
remote_counts =
Node.list()
|> Enum.map(fn remote_node ->
case RPC.call(remote_node, __MODULE__, :handle_remote_invalidate_tag, [
state.cache_name,
tag
]) do
{:ok, count} -> count
{:error, _} -> 0
end
end)
|> Enum.sum()
{:ok, local_count + remote_counts, %{state | backend_state: new_bs}}
end
@impl true
def size(%__MODULE__{} = state) do
local_size = state.backend.size(state.backend_state)
remote_sizes =
Node.list()
|> Enum.map(fn remote_node ->
case RPC.call(remote_node, __MODULE__, :handle_remote_size, [state.cache_name]) do
{:ok, size} -> size
{:error, _} -> 0
end
end)
|> Enum.sum()
local_size + remote_sizes
end
# ---- Remote Handlers ----
# These are called via RPC on the node that owns the key.
# They read topology state from persistent_term.
@doc false
def handle_remote_get(cache_name, key) do
ts = get_topology_state(cache_name)
ts.backend.get(ts.backend_state, key)
end
@doc false
def handle_remote_put(cache_name, entry) do
ts = get_topology_state(cache_name)
{:ok, _new_bs} = ts.backend.put(ts.backend_state, entry)
:ok
end
@doc false
def handle_remote_delete(cache_name, key) do
ts = get_topology_state(cache_name)
ts.backend.delete(ts.backend_state, key)
end
@doc false
def handle_remote_get_many(cache_name, keys) do
ts = get_topology_state(cache_name)
ts.backend.get_many(ts.backend_state, keys)
end
@doc false
def handle_remote_put_many(cache_name, entries) do
ts = get_topology_state(cache_name)
Enum.each(entries, fn entry ->
ts.backend.put(ts.backend_state, entry)
end)
:ok
end
@doc false
def handle_remote_flush(cache_name) do
ts = get_topology_state(cache_name)
ts.backend.flush(ts.backend_state)
end
@doc false
def handle_remote_invalidate_tag(cache_name, tag) do
ts = get_topology_state(cache_name)
{:ok, count, _new_bs} = ts.backend.delete_by_tag(ts.backend_state, tag)
count
end
@doc false
def handle_remote_size(cache_name) do
ts = get_topology_state(cache_name)
ts.backend.size(ts.backend_state)
end
defp get_topology_state(cache_name) do
:persistent_term.get({:vela_topology_state, cache_name})
end
defp unwrap_rpc({:ok, result}), do: result
defp unwrap_rpc({:error, _} = err), do: err
end