diff --git a/config/config.exs b/config/config.exs index 0dc5d3bd4a..e0e1386fe9 100644 --- a/config/config.exs +++ b/config/config.exs @@ -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 diff --git a/config/runtime.exs b/config/runtime.exs index f83af81b74..23288e1788 100644 --- a/config/runtime.exs +++ b/config/runtime.exs @@ -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} ]} ] diff --git a/lib/logflare/key_values.ex b/lib/logflare/key_values.ex index 267f648dc6..604e0ed58f 100644 --- a/lib/logflare/key_values.ex +++ b/lib/logflare/key_values.ex @@ -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 @@ -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 diff --git a/lib/logflare/key_values/cache.ex b/lib/logflare/key_values/cache.ex index 3a55a9b0d3..3a4f5b191b 100644 --- a/lib/logflare/key_values/cache.ex +++ b/lib/logflare/key_values/cache.ex @@ -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) @@ -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) diff --git a/lib/logflare/key_values/cache_warmer.ex b/lib/logflare/key_values/cache_warmer.ex index d82162376d..b91badae9e 100644 --- a/lib/logflare/key_values/cache_warmer.ex +++ b/lib/logflare/key_values/cache_warmer.ex @@ -5,12 +5,14 @@ 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 @@ -18,7 +20,8 @@ defmodule Logflare.KeyValues.CacheWarmer 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 -> @@ -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 -> @@ -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 diff --git a/lib/logflare/key_values/key_value_usage.ex b/lib/logflare/key_values/key_value_usage.ex new file mode 100644 index 0000000000..d2fb9cebfd --- /dev/null +++ b/lib/logflare/key_values/key_value_usage.ex @@ -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 diff --git a/lib/logflare/key_values/usage_touch_worker.ex b/lib/logflare/key_values/usage_touch_worker.ex new file mode 100644 index 0000000000..1aa3191860 --- /dev/null +++ b/lib/logflare/key_values/usage_touch_worker.ex @@ -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 diff --git a/priv/repo/migrations/20260921000000_create_key_value_usages.exs b/priv/repo/migrations/20260921000000_create_key_value_usages.exs new file mode 100644 index 0000000000..02fa2785ad --- /dev/null +++ b/priv/repo/migrations/20260921000000_create_key_value_usages.exs @@ -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()") + end + + create unique_index(:key_value_usages, [:key_value_id]) + create index(:key_value_usages, [:last_used_at]) + end +end diff --git a/test/logflare/key_values/cache_test.exs b/test/logflare/key_values/cache_test.exs index 0e67d2ad40..1e96e9acde 100644 --- a/test/logflare/key_values/cache_test.exs +++ b/test/logflare/key_values/cache_test.exs @@ -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 @@ -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") diff --git a/test/logflare/key_values/cache_warmer_test.exs b/test/logflare/key_values/cache_warmer_test.exs index a2def0d5ac..3c62bd22cb 100644 --- a/test/logflare/key_values/cache_warmer_test.exs +++ b/test/logflare/key_values/cache_warmer_test.exs @@ -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"}) @@ -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) @@ -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) diff --git a/test/logflare/key_values_test.exs b/test/logflare/key_values_test.exs index 9f93b079c7..97320bbf31 100644 --- a/test/logflare/key_values_test.exs +++ b/test/logflare/key_values_test.exs @@ -3,6 +3,7 @@ defmodule Logflare.KeyValuesTest do use Logflare.DataCase alias Logflare.KeyValues + alias Logflare.KeyValues.KeyValueUsage setup do user = insert(:user) @@ -92,6 +93,68 @@ defmodule Logflare.KeyValuesTest do assert upserted.inserted_at == inserted_at end + test "prune_usages/1", %{user: user} do + insert(:key_value, user: user, key: "stale") + %{id: fresh_id} = insert(:key_value, user: user, key: "fresh") + KeyValues.bump_usages([{user.id, "stale"}], DateTime.add(DateTime.utc_now(), -31, :day)) + KeyValues.bump_usages([{user.id, "fresh"}], DateTime.add(DateTime.utc_now(), -1, :day)) + + assert {1, _} = KeyValues.prune_usages(DateTime.utc_now()) + assert [%KeyValueUsage{key_value_id: ^fresh_id}] = Repo.all(KeyValueUsage) + end + + describe "bump_usages/2" do + alias Logflare.KeyValues.KeyValueUsage + + test "same user keys", %{user: user} do + kv1 = insert(:key_value, user: user, key: "k1") + kv2 = insert(:key_value, user: user, key: "k2") + now = DateTime.utc_now() + + assert :ok = KeyValues.bump_usages([{user.id, "k1"}, {user.id, "k2"}], now) + + usages = Repo.all(KeyValueUsage) + assert MapSet.new(usages, & &1.key_value_id) == MapSet.new([kv1.id, kv2.id]) + end + + test "multiple users keys", %{user: user} do + other = insert(:user) + kv1 = insert(:key_value, user: user, key: "shared") + kv2 = insert(:key_value, user: other, key: "shared") + + assert :ok = + KeyValues.bump_usages( + [{user.id, "shared"}, {other.id, "shared"}], + DateTime.utc_now() + ) + + values = for %{key_value_id: id} <- Repo.all(KeyValueUsage), into: MapSet.new(), do: id + + assert values == MapSet.new([kv1.id, kv2.id]) + end + + test "ignores no matching key_value", %{user: user} do + assert :ok = KeyValues.bump_usages([{user.id, "missing"}], DateTime.utc_now()) + assert [] = Repo.all(KeyValueUsage) + end + + test "upserts last_used_at on repeat", %{user: user} do + kv = insert(:key_value, user: user, key: "k1") + first = DateTime.utc_now() + + assert :ok = KeyValues.bump_usages([{user.id, "k1"}], first) + + later = DateTime.add(first, 1, :minute) + assert :ok = KeyValues.bump_usages([{user.id, "k1"}], later) + + assert [%KeyValueUsage{key_value_id: kv_id, last_used_at: last_used_at}] = + Repo.all(KeyValueUsage) + + assert kv_id == kv.id + assert DateTime.compare(last_used_at, first) == :gt + end + end + test "count_key_values/1", %{user: user} do assert 0 = KeyValues.count_key_values(user.id) insert(:key_value, user: user, key: "k1") diff --git a/test/support/factory.ex b/test/support/factory.ex index 95e02a9b02..b1c7f22bee 100644 --- a/test/support/factory.ex +++ b/test/support/factory.ex @@ -23,6 +23,7 @@ defmodule Logflare.Factory do alias Logflare.Users.UserPreferences alias Logflare.Alerting.AlertQuery alias Logflare.KeyValues.KeyValue + alias Logflare.KeyValues.KeyValueUsage alias Logflare.Google.BigQuery.SchemaUtils alias PaperTrail.Version @@ -434,6 +435,13 @@ defmodule Logflare.Factory do } end + def key_value_usage_factory do + %KeyValueUsage{ + key_value: build(:key_value), + last_used_at: DateTime.utc_now() + } + end + def saved_search_factory(attrs \\ %{}) do source = build(:source, user: build(:user)) %{bigquery_schema: schema} = build(:source_schema)