Packages
electric
1.7.5
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.ShapeCache.ShapeStatus
alias Electric.StatusMonitor
alias Electric.Telemetry.OpenTelemetry
require Logger
@schema NimbleOptions.new!(
max_shapes: [type: {:or, [:non_neg_integer, nil]}, default: nil],
period: [type: :non_neg_integer, default: 60_000],
stack_id: [type: :string, required: true]
)
def name(stack_ref) do
Electric.ProcessRegistry.name(stack_ref, __MODULE__)
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),
period: Keyword.fetch!(opts, :period)
}
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
case StatusMonitor.status(state.stack_id) do
%{shape: :up} ->
shape_count = shape_count(state)
if shape_count > max_shapes do
expire_shapes(shape_count, state)
end
status ->
# We do not expire shapes if the stack is not active since this may mean that
# shapes have not fully restored yet and we don't want to expire while restoring
# as this may cause race conditions.
Logger.debug("Expiry check skipped due to inactive stack: #{inspect(status)}")
end
end
defp expire_shapes(shape_count, state) do
number_to_expire = shape_count - state.max_shapes
{handles_to_expire, min_age} = least_recently_used(state, number_to_expire)
Logger.info(
"Expiring shapes as the number of shapes has exceeded the limit",
number_to_expire: number_to_expire,
max_shapes: state.max_shapes
)
OpenTelemetry.with_span(
"expiry_manager.expire_shapes",
[
max_shapes: state.max_shapes,
shape_count: shape_count,
number_to_expire: number_to_expire,
elapsed_minutes_since_use: min_age
],
fn ->
Electric.ShapeCache.ShapeCleaner.remove_shapes(state.stack_id, handles_to_expire)
end
)
end
defp least_recently_used(%{stack_id: stack_id}, number_to_expire) do
ShapeStatus.least_recently_used(stack_id, number_to_expire)
end
defp shape_count(%{stack_id: stack_id}) do
ShapeStatus.count_shapes(stack_id)
end
end