Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion config/config.exs
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,10 @@ config :textbin,
generators: [timestamp_type: :utc_datetime],
allow_guest_pastes: false,
guest_paste_ttl: "6h",
max_paste_bytes: 1_048_576
max_paste_bytes: 1_048_576,
expiration_cleanup_enabled: true,
expiration_cleanup_interval_ms: :timer.minutes(15),
expiration_cleanup_batch_size: 500

# Configures the endpoint
config :textbin, TextbinWeb.Endpoint,
Expand Down
3 changes: 3 additions & 0 deletions config/test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,9 @@ config :textbin, TextbinWeb.Endpoint,
# In test we don't send emails
config :textbin, Textbin.Mailer, adapter: Swoosh.Adapters.Test

# Cleanup is exercised explicitly so the global worker does not race SQL sandbox tests.
config :textbin, expiration_cleanup_enabled: false

# Disable swoosh api client as it is only required for production adapters
config :swoosh, :api_client, false

Expand Down
25 changes: 15 additions & 10 deletions lib/textbin/application.ex
Original file line number Diff line number Diff line change
Expand Up @@ -7,16 +7,13 @@ defmodule Textbin.Application do

@impl true
def start(_type, _args) do
children = [
TextbinWeb.Telemetry,
Textbin.Repo,
{DNSCluster, query: Application.get_env(:textbin, :dns_cluster_query) || :ignore},
{Phoenix.PubSub, name: Textbin.PubSub},
# Start a worker by calling: Textbin.Worker.start_link(arg)
# {Textbin.Worker, arg},
# Start to serve requests, typically the last entry
TextbinWeb.Endpoint
]
children =
[
TextbinWeb.Telemetry,
Textbin.Repo,
{DNSCluster, query: Application.get_env(:textbin, :dns_cluster_query) || :ignore},
{Phoenix.PubSub, name: Textbin.PubSub}
] ++ expiration_cleanup_children() ++ [TextbinWeb.Endpoint]

# See https://hexdocs.pm/elixir/Supervisor.html
# for other strategies and supported options
Expand All @@ -31,4 +28,12 @@ defmodule Textbin.Application do
TextbinWeb.Endpoint.config_change(changed, removed)
:ok
end

defp expiration_cleanup_children do
if Application.get_env(:textbin, :expiration_cleanup_enabled, true) do
[Textbin.Pastes.ExpirationCleaner]
else
[]
end
end
end
31 changes: 31 additions & 0 deletions lib/textbin/pastes.ex
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,37 @@ defmodule Textbin.Pastes do

def delete_paste(%Scope{}, %Paste{}), do: {:error, :not_found}

@doc """
Hard-deletes one bounded batch of expired pastes and returns the number deleted.

Rows with an expiration equal to `:now` are considered expired. The oldest
expirations are selected first so a backlog is drained predictably.
"""
@spec delete_expired_pastes(keyword()) :: non_neg_integer()
def delete_expired_pastes(opts \\ []) do
now = Keyword.get_lazy(opts, :now, &Paste.utc_now_ms/0)
limit = Keyword.get(opts, :limit, 500)

if !is_integer(limit) or limit <= 0 do
raise ArgumentError, ":limit must be a positive integer"
end

expired_ids =
from p in Paste,
where: not is_nil(p.expires_at) and p.expires_at <= ^now,
order_by: [asc: p.expires_at, asc: p.id],
limit: ^limit,
select: p.id

{deleted_count, nil} =
Repo.delete_all(
from p in Paste,
where: p.id in subquery(expired_ids)
)

deleted_count
end

def change_paste(%Scope{user: %User{} = user}, %Paste{} = paste, attrs \\ %{}) do
paste = %{paste | user_id: paste.user_id || user.id}

Expand Down
128 changes: 128 additions & 0 deletions lib/textbin/pastes/expiration_cleaner.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,128 @@
defmodule Textbin.Pastes.ExpirationCleaner do
@moduledoc """
Periodically hard-deletes expired pastes in bounded batches.
"""

use GenServer

require Logger

alias Textbin.Pastes

@telemetry_event [:textbin, :pastes, :expiration_cleanup]
@default_interval_ms :timer.minutes(15)
@default_batch_size 500

def start_link(opts \\ []) do
case Keyword.get(opts, :name, __MODULE__) do
nil -> GenServer.start_link(__MODULE__, opts)
name -> GenServer.start_link(__MODULE__, opts, name: name)
end
end

@doc """
Runs one cleanup batch immediately.

Returns the number of deleted pastes, or `:error` after a cleanup failure has
been logged and reported through telemetry.
"""
@spec run_now(GenServer.server()) :: non_neg_integer() | :error
def run_now(server \\ __MODULE__) do
GenServer.call(server, :run_now)
end

@impl true
def init(opts) do
interval_ms =
opts
|> Keyword.get(
:interval_ms,
Application.get_env(
:textbin,
:expiration_cleanup_interval_ms,
@default_interval_ms
)
)
|> positive_integer!(:interval_ms)

batch_size =
opts
|> Keyword.get(
:batch_size,
Application.get_env(:textbin, :expiration_cleanup_batch_size, @default_batch_size)
)
|> positive_integer!(:batch_size)

initial_delay_ms =
opts
|> Keyword.get(:initial_delay_ms, 0)
|> non_negative_integer!(:initial_delay_ms)

state = %{interval_ms: interval_ms, batch_size: batch_size}
{:ok, schedule_cleanup(state, initial_delay_ms)}
end

@impl true
def handle_call(:run_now, _from, state) do
{:reply, clean_batch(state), state}
end

@impl true
def handle_info(:cleanup, state) do
next_delay_ms =
case clean_batch(state) do
deleted_count when deleted_count == state.batch_size -> 0
_deleted_count_or_error -> state.interval_ms
end

{:noreply, schedule_cleanup(state, next_delay_ms)}
end

defp clean_batch(state) do
started_at = System.monotonic_time()

try do
deleted_count = Pastes.delete_expired_pastes(limit: state.batch_size)

:telemetry.execute(
@telemetry_event,
%{duration: System.monotonic_time() - started_at, deleted_count: deleted_count},
%{batch_size: state.batch_size, result: :ok}
)

if deleted_count > 0 do
Logger.info("Deleted #{deleted_count} expired pastes")
end

deleted_count
rescue
exception ->
:telemetry.execute(
@telemetry_event,
%{duration: System.monotonic_time() - started_at, deleted_count: 0},
%{batch_size: state.batch_size, result: :error}
)

Logger.error("Failed to delete expired pastes: #{Exception.message(exception)}")
:error
end
end

defp schedule_cleanup(state, delay_ms) do
Process.send_after(self(), :cleanup, delay_ms)
state
end

defp positive_integer!(value, _name) when is_integer(value) and value > 0, do: value

defp positive_integer!(value, name) do
raise ArgumentError, "#{inspect(name)} must be a positive integer, got: #{inspect(value)}"
end

defp non_negative_integer!(value, _name) when is_integer(value) and value >= 0, do: value

defp non_negative_integer!(value, name) do
raise ArgumentError,
"#{inspect(name)} must be a non-negative integer, got: #{inspect(value)}"
end
end
11 changes: 11 additions & 0 deletions lib/textbin_web/telemetry.ex
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,17 @@ defmodule TextbinWeb.Telemetry do
"The time the connection spent waiting before being checked out for the query"
),

# Paste expiration cleanup
sum("textbin.pastes.expiration_cleanup.deleted_count",
tags: [:result],
description: "The number of expired pastes permanently deleted"
),
summary("textbin.pastes.expiration_cleanup.duration",
tags: [:result],
unit: {:native, :millisecond},
description: "The time spent deleting one batch of expired pastes"
),

# VM Metrics
summary("vm.memory.total", unit: {:byte, :kilobyte}),
summary("vm.total_run_queue_lengths.total"),
Expand Down
52 changes: 52 additions & 0 deletions test/textbin/pastes/expiration_cleaner_test.exs
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
defmodule Textbin.Pastes.ExpirationCleanerTest do
use Textbin.DataCase, async: false

alias Textbin.Pastes.ExpirationCleaner
alias Textbin.Pastes.Paste

import Textbin.AccountsFixtures

test "run_now/1 deletes one batch and emits telemetry" do
user = user_fixture()

expired_paste =
Repo.insert!(%Paste{
data: "expired",
syntax_highlight: "plain",
visibility: "private",
user_id: user.id,
expires_at: DateTime.add(Paste.utc_now_ms(), -1, :second)
})

handler_id = "expiration-cleaner-test-#{System.unique_integer([:positive])}"

:telemetry.attach(
handler_id,
[:textbin, :pastes, :expiration_cleanup],
fn event, measurements, metadata, test_pid ->
send(test_pid, {:cleanup_telemetry, event, measurements, metadata})
end,
self()
)

on_exit(fn -> :telemetry.detach(handler_id) end)

cleaner =
start_supervised!(
{ExpirationCleaner,
name: nil, initial_delay_ms: :timer.hours(1), interval_ms: :timer.hours(1), batch_size: 1}
)

Ecto.Adapters.SQL.Sandbox.allow(Textbin.Repo, self(), cleaner)

assert ExpirationCleaner.run_now(cleaner) == 1
refute Repo.get(Paste, expired_paste.id)

assert_receive {:cleanup_telemetry, [:textbin, :pastes, :expiration_cleanup], measurements,
metadata}

assert measurements.deleted_count == 1
assert is_integer(measurements.duration)
assert metadata == %{batch_size: 1, result: :ok}
end
end
40 changes: 40 additions & 0 deletions test/textbin/pastes_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -218,8 +218,48 @@ defmodule Textbin.PastesTest do
assert Pastes.get_paste!(other_scope, paste.id)
end

test "delete_expired_pastes/1 deletes the oldest expired rows in bounded batches", %{
scope: scope
} do
now = Paste.utc_now_ms()

oldest_expired =
insert_paste!(scope, "oldest expired", DateTime.add(now, -2, :second))

newly_expired = insert_paste!(scope, "newly expired", now)
active = insert_paste!(scope, "active", DateTime.add(now, 1, :hour))
never_expires = insert_paste!(scope, "never", nil)

assert Pastes.delete_expired_pastes(now: now, limit: 1) == 1
refute Repo.get(Paste, oldest_expired.id)
assert Repo.get(Paste, newly_expired.id)

assert Pastes.delete_expired_pastes(now: now, limit: 1) == 1
refute Repo.get(Paste, newly_expired.id)
assert Repo.get(Paste, active.id)
assert Repo.get(Paste, never_expires.id)

assert Pastes.delete_expired_pastes(now: now, limit: 1) == 0
end

test "delete_expired_pastes/1 rejects invalid batch limits" do
assert_raise ArgumentError, ":limit must be a positive integer", fn ->
Pastes.delete_expired_pastes(limit: 0)
end
end

test "change_paste/3 returns a paste changeset", %{scope: scope} do
assert %Ecto.Changeset{} = Pastes.change_paste(scope, %Paste{})
end

defp insert_paste!(scope, data, expires_at) do
Repo.insert!(%Paste{
data: data,
syntax_highlight: "plain",
visibility: "private",
user_id: scope.user.id,
expires_at: expires_at
})
end
end
end
Loading