Packages
electric
1.1.13
1.7.8
1.7.7
1.7.6
1.7.5
1.7.4
1.7.3
1.7.2
1.7.1
1.7.0
1.6.10
1.6.9
1.6.8
1.6.7
1.6.6
1.6.5
1.6.4
1.6.3
1.6.2
1.6.1
1.6.0
1.5.1
1.5.0
1.4.16
1.4.16-beta-1
1.4.15
1.4.14
1.4.13
1.4.12
1.4.11
1.4.10
1.4.8
1.4.7
1.4.6
1.4.5
1.4.4
1.4.3
1.4.2
1.4.1
1.4.0
1.3.4
1.3.3
1.3.2
1.2.4
1.2.3
1.2.2
1.2.1
1.2.0
1.1.14
1.1.13
1.1.12
1.1.11
1.1.10
1.1.9
1.1.8
1.1.7
1.1.6
retired
1.1.5
retired
1.1.4
retired
1.1.3
retired
1.1.2
1.1.1
1.1.0
1.0.24
1.0.23
1.0.22
1.0.21
1.0.20
1.0.19
1.0.18
1.0.17
1.0.15
1.0.13
1.0.12
1.0.11
1.0.10
1.0.9
1.0.5
1.0.4
1.0.3
1.0.2
1.0.1
1.0.0
1.0.0-beta.23
1.0.0-beta.22
1.0.0-beta.20
1.0.0-beta.19
1.0.0-beta.18
1.0.0-beta.17
1.0.0-beta.16
1.0.0-beta.15
1.0.0-beta.14
1.0.0-beta.13
1.0.0-beta.12
1.0.0-beta.11
1.0.0-beta.10
1.0.0-beta.9
1.0.0-beta.8
1.0.0-beta.7
1.0.0-beta.6
1.0.0-beta.5
1.0.0-beta.4
1.0.0-beta.3
1.0.0-beta.2
1.0.0-beta.1
0.9.5
0.9.4
0.9.3
0.9.2
0.9.1
0.9.0
0.8.1
0.8.0
0.7.7
0.7.6
0.7.5
0.7.4
0.7.3
0.7.2
0.7.1
0.7.0
0.6.3
0.6.2
0.6.1
0.5.2
0.4.4
Postgres sync engine. Sync little subsets of your Postgres data into local apps and services.
Current section
Files
Jump to
Current section
Files
lib/electric/shape_cache/expiry_manager.ex
defmodule Electric.ShapeCache.ExpiryManager do
use GenServer
alias Electric.Telemetry.OpenTelemetry
require Logger
@schema NimbleOptions.new!(
max_shapes: [type: {:or, [:non_neg_integer, nil]}, default: nil],
expiry_batch_size: [type: :pos_integer],
period: [type: :non_neg_integer, default: 60_000],
stack_id: [type: :string, required: true],
shape_status: [type: :mod_arg, required: true]
)
def name(stack_id) when not is_map(stack_id) and not is_list(stack_id) do
Electric.ProcessRegistry.name(stack_id, __MODULE__)
end
def name(opts) do
stack_id = Access.fetch!(opts, :stack_id)
name(stack_id)
end
def start_link(opts) do
with {:ok, opts} <- NimbleOptions.validate(opts, @schema) do
GenServer.start_link(__MODULE__, opts, name: name(opts))
end
end
def init(opts) do
stack_id = Keyword.fetch!(opts, :stack_id)
Process.set_label({:shape_expiry_manager, stack_id})
Logger.metadata(stack_id: stack_id)
Electric.Telemetry.Sentry.set_tags_context(stack_id: stack_id)
state =
%{
stack_id: stack_id,
max_shapes: Keyword.fetch!(opts, :max_shapes),
expiry_batch_size: Keyword.fetch!(opts, :expiry_batch_size),
period: Keyword.fetch!(opts, :period),
shape_status: Keyword.fetch!(opts, :shape_status)
}
if not is_nil(state.max_shapes), do: schedule_next_check(state)
{:ok, state}
end
defp schedule_next_check(state) do
Process.send_after(self(), :maybe_expire_shapes, state.period)
end
def handle_info(:maybe_expire_shapes, state) do
maybe_expire_shapes(state)
schedule_next_check(state)
{:noreply, state}
end
defp maybe_expire_shapes(%{max_shapes: nil}), do: :ok
defp maybe_expire_shapes(%{max_shapes: max_shapes} = state) do
shape_count = shape_count(state)
if shape_count > max_shapes do
expire_shapes(shape_count, state)
end
end
defp expire_shapes(shape_count, state) do
shapes_to_expire = least_recently_used(state, state.expiry_batch_size)
Logger.info(
"Expiring #{length(shapes_to_expire)} shapes as the number of shapes " <>
"has exceeded the limit (#{state.max_shapes})"
)
OpenTelemetry.with_span(
"expiry_manager.expire_shapes",
[
max_shapes: state.max_shapes,
shape_count: shape_count,
number_to_expire: state.expiry_batch_size
],
fn -> Enum.each(shapes_to_expire, &expire_shape(&1, state)) end
)
end
defp expire_shape(shape, state) do
OpenTelemetry.with_span(
"expiry_manager.expire_shape",
[
shape_handle: shape.shape_handle,
elapsed_minutes_since_use: shape.elapsed_minutes_since_use
],
fn ->
Electric.ShapeCache.ShapeCleaner.remove_shape(shape.shape_handle,
stack_id: state.stack_id
)
end
)
end
defp least_recently_used(%{shape_status: {shape_status, shape_status_state}}, number_to_expire) do
OpenTelemetry.with_span("expiry_manager.get_least_recently_used", [], fn ->
shape_status.least_recently_used(shape_status_state, number_to_expire)
end)
end
defp shape_count(%{shape_status: {shape_status, shape_status_state}}) do
OpenTelemetry.with_span("expiry_manager.get_shape_count", [], fn ->
shape_status.count_shapes(shape_status_state)
end)
end
end