Skip to content
Open
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
2 changes: 2 additions & 0 deletions config/config.exs
Original file line number Diff line number Diff line change
Expand Up @@ -173,6 +173,8 @@ config :logflare, Logflare.ContextCache.CacheBuster,
replication_slot: :temporary,
publications: ["logflare_pub"]

config :logflare, Logflare.KeyValues.CacheWarmer, warm_limit: 500_000

config :open_api_spex, :cache_adapter, OpenApiSpex.Plug.PersistentTermCache

config :logflare, Logflare.Cluster.Utils, min_cluster_size: 1
Expand Down
3 changes: 2 additions & 1 deletion config/runtime.exs
Original file line number Diff line number Diff line change
Expand Up @@ -584,7 +584,8 @@ config :logflare, Oban,
{Oban.Plugins.Cron,
crontab: [
{"* * * * *", Logflare.Alerting.AlertSchedulerWorker},
{"*/15 * * * *", Logflare.Sources.RecentEventsTouchWorker}
{"*/15 * * * *", Logflare.Sources.RecentEventsTouchWorker},
{"*/15 * * * *", Logflare.KeyValues.UsageTouchWorker}
]}
]

Expand Down
28 changes: 28 additions & 0 deletions lib/logflare/key_values.ex
Original file line number Diff line number Diff line change
Expand Up @@ -4,9 +4,11 @@ defmodule Logflare.KeyValues do
import Ecto.Query

alias Logflare.KeyValues.KeyValue
alias Logflare.KeyValues.KeyValueUsage
alias Logflare.Repo

@list_limit 500
@usage_retention_days 30

@spec list_key_values_query(keyword()) :: Ecto.Query.t()
def list_key_values_query(kw) do
Expand Down Expand Up @@ -125,6 +127,32 @@ defmodule Logflare.KeyValues do
)
end

@spec bump_usages([{integer(), String.t()}], DateTime.t()) :: :ok
def bump_usages(pairs, now) when is_list(pairs) do
pairs
|> Enum.group_by(fn {uid, _} -> uid end, fn {_, k} -> k end)
|> Enum.each(fn {user_id, keys} ->
src =
KeyValue
|> where(user_id: ^user_id)
|> where([kv], kv.key in ^keys)
|> select([kv], %{key_value_id: kv.id, last_used_at: type(^now, :utc_datetime_usec)})

Repo.insert_all(KeyValueUsage, src,
on_conflict: {:replace, [:last_used_at]},
conflict_target: [:key_value_id]
)
end)

:ok
end

@spec prune_usages(DateTime.t()) :: {non_neg_integer(), nil}
def prune_usages(now \\ DateTime.utc_now()) do
cutoff = DateTime.add(now, -@usage_retention_days, :day)
KeyValueUsage |> where([u], u.last_used_at < ^cutoff) |> Repo.delete_all()
end

@spec bulk_delete_by_keys(integer(), [String.t()]) :: {non_neg_integer(), nil}
def bulk_delete_by_keys(user_id, keys) when is_list(keys) do
KeyValue
Expand Down
46 changes: 44 additions & 2 deletions lib/logflare/key_values/cache.ex
Original file line number Diff line number Diff line change
@@ -1,15 +1,19 @@
defmodule Logflare.KeyValues.Cache do
@moduledoc false

import Cachex.Spec

alias Logflare.ContextCache
alias Logflare.KeyValues
alias Logflare.KeyValues.CacheWarmer
alias Logflare.Repo
alias Logflare.Utils

import Cachex.Spec

@behaviour ContextCache

@usage_touch_window :timer.minutes(20)
@usage_touch_chunk 1_000

def child_spec(_) do
stats = Application.get_env(:logflare, :cache_stats, false)

Expand Down Expand Up @@ -85,6 +89,44 @@ defmodule Logflare.KeyValues.Cache do
end
end

@doc """
Records which key-values were recently in use so the next full warm can
prioritise them.

Usage is inferred from entries that entered this node's cache since the last
run, excluding those loaded by the warmer itself. The signal is approximate:
it captures that a key was needed, not how often it was read.

Returns the number of `{user_id, key}` pairs whose usage was recorded.
"""
@spec touch_recent_usages(DateTime.t()) :: {:ok, non_neg_integer()}
def touch_recent_usages(now \\ DateTime.utc_now()) do
query = Cachex.Query.build(where: {:>=, :modified, touch_cutoff(now)}, output: :key)

__MODULE__
|> Cachex.stream!(query)
|> Stream.flat_map(fn
{:lookup, [user_id, key, _accessor]} -> [{user_id, key}]
_ -> []
end)
|> Stream.uniq()
|> Stream.chunk_every(@usage_touch_chunk)
|> Enum.reduce(0, fn pairs, acc ->
KeyValues.bump_usages(pairs, now)
acc + length(pairs)
end)
|> then(&{:ok, &1})
end

defp touch_cutoff(now) do
window_start = DateTime.to_unix(now, :millisecond) - @usage_touch_window

case CacheWarmer.warmed_at() do
nil -> window_start
warmed_at -> max(window_start, DateTime.to_unix(warmed_at, :millisecond) + 1)
end
end

@impl ContextCache
def bust_by(kw) do
entries = bust_entries(kw)
Expand Down
31 changes: 28 additions & 3 deletions lib/logflare/key_values/cache_warmer.ex
Original file line number Diff line number Diff line change
Expand Up @@ -5,20 +5,23 @@ defmodule Logflare.KeyValues.CacheWarmer do

alias Logflare.KeyValues.Cache
alias Logflare.KeyValues.KeyValue
alias Logflare.KeyValues.KeyValueUsage
alias Logflare.Repo

require Logger
import Ecto.Query

@pt_key {__MODULE__, :initialized}
@pt_warmed_at {__MODULE__, :warmed_at}

@impl true
def execute(_state) do
if initialized?() do
Repo.apply_with_replica(__MODULE__, :warm_recent, [])
else
try do
Repo.apply_with_replica(__MODULE__, :warm_full, [])
Repo.apply_with_replica(__MODULE__, :warm_top_n, [])
:persistent_term.put(@pt_warmed_at, DateTime.utc_now())
:persistent_term.put(@pt_key, true)
rescue
e ->
Expand All @@ -29,9 +32,19 @@ defmodule Logflare.KeyValues.CacheWarmer do
:ignore
end

def warm_full do
Repo.transaction(fn ->
def warm_top_n do
limit =
Application.get_env(:logflare, __MODULE__, [])
|> Keyword.get(:warm_limit, 500_000)

ordered =
KeyValue
|> join(:left, [kv], u in KeyValueUsage, on: u.key_value_id == kv.id)
|> order_by([kv, u], desc_nulls_last: u.last_used_at, desc: kv.updated_at)
|> limit(^limit)

Repo.transaction(fn ->
ordered
|> Repo.stream()
|> Stream.chunk_every(500)
|> Enum.each(fn chunk ->
Expand All @@ -55,6 +68,18 @@ defmodule Logflare.KeyValues.CacheWarmer do
{{:lookup, [kv.user_id, kv.key, nil]}, {:cached, kv.value}}
end

@doc """
When the last full warm completed on this node, or `nil` if none has.

Entries `put` by the warmer were written at or before this moment, so consumers
of write-recency (like `Cache.touch_recent_usages/1`) can use it to tell warmed
entries apart from genuine cache misses.
"""
@spec warmed_at() :: DateTime.t() | nil
def warmed_at do
:persistent_term.get(@pt_warmed_at, nil)
end

defp initialized? do
:persistent_term.get(@pt_key, false)
end
Expand Down
20 changes: 20 additions & 0 deletions lib/logflare/key_values/key_value_usage.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
defmodule Logflare.KeyValues.KeyValueUsage do
@moduledoc """
A helper table storing the last-touched timestamp for `KeyValue` entries.

This table is intentionally **excluded from the `logflare_pub` publication**, so
upserting `last_used_at` does not flow through `ContextCache.CacheBuster` and does
not evict entries from `KeyValues.Cache` — which is precisely the cache whose access
patterns drive the usage signal. Storing usage on the `key_values` row itself would
cause every touch to bust the cache entry it is trying to track.
"""
use TypedEctoSchema

alias Logflare.KeyValues.KeyValue

typed_schema "key_value_usages" do
field :last_used_at, :utc_datetime_usec

belongs_to :key_value, KeyValue
end
end
15 changes: 15 additions & 0 deletions lib/logflare/key_values/usage_touch_worker.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
defmodule Logflare.KeyValues.UsageTouchWorker do
@moduledoc false
use Oban.Worker, queue: :default, max_attempts: 1

alias Logflare.KeyValues
alias Logflare.KeyValues.Cache

@impl Oban.Worker
@spec perform(Oban.Job.t()) :: :ok
def perform(_job) do
Cache.touch_recent_usages()
KeyValues.prune_usages()
:ok
end
end
13 changes: 13 additions & 0 deletions priv/repo/migrations/20260921000000_create_key_value_usages.exs
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
defmodule Logflare.Repo.Migrations.CreateKeyValueUsages do
use Ecto.Migration

def change do
create table(:key_value_usages) do
add :key_value_id, references(:key_values, on_delete: :delete_all), null: false
add :last_used_at, :utc_datetime_usec, null: false, default: fragment("now()")
Comment on lines +6 to +7

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

hmm this also implies that if all 1 mil records are in use, then 1 mil records will also be created in this table. why would we want a separate table as compared to having these two fields directly on the :key_values table?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'm guessing that the separate table means that the table isn't tracked it isn't tracked as part of ContextCache since it is not included in the published tables, so bumping it would not trigger cache invalidations.
should make this clear in the KeyValuesUsage module

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

should probably have some periodic truncation logic to keep the size of this table down

end

create unique_index(:key_value_usages, [:key_value_id])
create index(:key_value_usages, [:last_used_at])
end
end
65 changes: 65 additions & 0 deletions test/logflare/key_values/cache_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -3,8 +3,11 @@ defmodule Logflare.KeyValues.CacheTest do
use Logflare.DataCase, async: false

alias Logflare.KeyValues
alias Logflare.KeyValues.CacheWarmer

setup do
:persistent_term.erase({CacheWarmer, :initialized})
:persistent_term.erase({CacheWarmer, :warmed_at})
user = insert(:user)
[user: user]
end
Expand Down Expand Up @@ -66,6 +69,68 @@ defmodule Logflare.KeyValues.CacheTest do
assert {:ok, 0} = KeyValues.Cache.bust_by(user_id: user.id, key: "nonexistent")
end

describe "touch_recent_usages/1" do
alias Logflare.KeyValues.KeyValueUsage

setup %{user: user} do
kv1 = insert(:key_value, user: user, key: "k1")
kv2 = insert(:key_value, user: user, key: "k2")
[kv1: kv1, kv2: kv2]
end

test "duplicated accessor-path", %{user: user, kv1: %{id: kv1_id}} do
for path <- [nil, "org.id", "$.org.id"] do
{:ok, true} =
Cachex.put(KeyValues.Cache, {:lookup, [user.id, "k1", path]}, {:cached, "v"})
end

assert {:ok, 1} = KeyValues.Cache.touch_recent_usages(DateTime.utc_now())
assert [%KeyValueUsage{key_value_id: ^kv1_id}] = Repo.all(KeyValueUsage)
end

test "count key", %{user: user} do
{:ok, true} = Cachex.put(KeyValues.Cache, {:count, user.id}, {:cached, 5})

assert {:ok, 0} = KeyValues.Cache.touch_recent_usages(DateTime.utc_now())
assert [] = Repo.all(KeyValueUsage)
end

test "outside window", %{user: user} do
{:ok, true} = Cachex.put(KeyValues.Cache, {:lookup, [user.id, "k1", nil]}, {:cached, "v"})

# Passing a future `now` makes cutoff = future - 20min, which is after the
# entry's modified time (≈ real now), so the entry falls outside the window.
future = DateTime.add(DateTime.utc_now(), 25, :minute)
assert {:ok, 0} = KeyValues.Cache.touch_recent_usages(future)
assert [] = Repo.all(KeyValueUsage)
end

test "entries written by the initial warm", %{user: user, kv1: %{id: kv1_id, key: key}} do
CacheWarmer.execute(nil)
assert {:ok, 0} = KeyValues.Cache.touch_recent_usages(DateTime.utc_now())

Process.sleep(2)
{:ok, true} = Cachex.put(KeyValues.Cache, {:lookup, [user.id, key, nil]}, {:cached, "v"})

assert {:ok, 1} = KeyValues.Cache.touch_recent_usages(DateTime.utc_now())
assert [%KeyValueUsage{key_value_id: ^kv1_id}] = Repo.all(KeyValueUsage)
end

test "multiple users", %{user: user, kv1: kv1, kv2: kv2} do
other = insert(:user)
kv_other = insert(:key_value, user: other, key: "k3")

for {uid, key} <- [{user.id, "k1"}, {user.id, "k2"}, {other.id, "k3"}] do
{:ok, true} = Cachex.put(KeyValues.Cache, {:lookup, [uid, key, nil]}, {:cached, "v"})
end

assert {:ok, 3} = KeyValues.Cache.touch_recent_usages(DateTime.utc_now())

ids = Repo.all(KeyValueUsage) |> MapSet.new(& &1.key_value_id)
assert ids == MapSet.new([kv1.id, kv2.id, kv_other.id])
end
end

describe "count/1" do
test "caches the count for a user", %{user: user} do
insert(:key_value, user: user, key: "k1")
Expand Down
43 changes: 41 additions & 2 deletions test/logflare/key_values/cache_warmer_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ defmodule Logflare.KeyValues.CacheWarmerTest do
[user: user]
end

describe "initial warm (full table stream)" do
describe "initial warm (top-N by recency)" do
test "populates cache with all key_values", %{user: user} do
kv1 = insert(:key_value, user: user, key: "k1", value: %{"v" => "1"})
kv2 = insert(:key_value, user: user, key: "k2", value: %{"v" => "2"})
Expand All @@ -26,6 +26,45 @@ defmodule Logflare.KeyValues.CacheWarmerTest do
assert {:cached, kv2.value} == Cachex.get!(Cache, {:lookup, [user.id, "k2", nil]})
end

test "warms only the top-N rows ordered by recency of use", %{user: user} do
kv_used = insert(:key_value, user: user, key: "used", value: %{"v" => "used"})
insert(:key_value, user: user, key: "unused", value: %{"v" => "unused"})

insert(:key_value_usage, key_value: kv_used, last_used_at: DateTime.utc_now())

prev = Application.get_env(:logflare, CacheWarmer, [])
Application.put_env(:logflare, CacheWarmer, Keyword.put(prev, :warm_limit, 1))
on_exit(fn -> Application.put_env(:logflare, CacheWarmer, prev) end)

CacheWarmer.execute(nil)

assert {:cached, kv_used.value} == Cachex.get!(Cache, {:lookup, [user.id, "used", nil]})
assert is_nil(Cachex.get!(Cache, {:lookup, [user.id, "unused", nil]}))
end

test "falls back to updated_at recency when no usage rows exist", %{user: user} do
old_time = DateTime.add(DateTime.utc_now(), -2, :hour)

Repo.insert!(%Logflare.KeyValues.KeyValue{
user_id: user.id,
key: "old",
value: %{"v" => "old"},
inserted_at: old_time,
updated_at: old_time
})

kv_new = insert(:key_value, user: user, key: "new", value: %{"v" => "new"})

prev = Application.get_env(:logflare, CacheWarmer, [])
Application.put_env(:logflare, CacheWarmer, Keyword.put(prev, :warm_limit, 1))
on_exit(fn -> Application.put_env(:logflare, CacheWarmer, prev) end)

CacheWarmer.execute(nil)

assert {:cached, kv_new.value} == Cachex.get!(Cache, {:lookup, [user.id, "new", nil]})
assert is_nil(Cachex.get!(Cache, {:lookup, [user.id, "old", nil]}))
end

test "marks itself as initialized after first run" do
refute :persistent_term.get(@pt_key, false)

Expand All @@ -35,7 +74,7 @@ defmodule Logflare.KeyValues.CacheWarmerTest do
end

test "if cache warmer fails, does not mark itself as initialized" do
stub(Logflare.KeyValues.CacheWarmer, :warm_full, fn ->
stub(Logflare.KeyValues.CacheWarmer, :warm_top_n, fn ->
raise RuntimeError, "test"
end)

Expand Down
Loading
Loading