Skip to content
Draft
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
1 change: 0 additions & 1 deletion .dialyzer_ignore.exs
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,6 @@
{"lib/logflare_web/controllers/api/endpoint_controller.ex", :pattern_match},
{"lib/logflare_web/controllers/billing_controller.ex", :unused_fun},
{"lib/logflare_web/controllers/billing_controller.ex", :pattern_match_cov},
{"lib/logflare_web/controllers/health_check_controller.ex", :pattern_match},
{"lib/logflare_web/controllers/source_controller.ex", :pattern_match},
{"lib/logflare_web/live/monaco_editor_component.ex", :pattern_match},
{"lib/logflare_web/live/query_live.ex", :no_return},
Expand Down
5 changes: 3 additions & 2 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -84,8 +84,9 @@ gate to go green - fix the code.
that have pre-existing findings. They are off so the gate is green on existing
code. Fixing a backlog check's findings and removing it from that list is a
welcome change on its own. Never add a check to the list to go green.
- `mix test.slop` runs `ex_dna --max-clones 29`, a duplication ratchet. A new
clone fails CI. When you remove clones, lower the number. Never raise it.
- `mix test.slop` runs `ex_dna --max-clones <Current clone number>`,
a duplication ratchet. A new clone fails CI. When you remove clones,
lower the number. Never raise it.
- `mix test.structure` runs `reach.check --smells` against
`.reach.baseline.json` (194 accepted findings), so only **new** structural
smells fail. To accept a
Expand Down
28 changes: 8 additions & 20 deletions lib/logflare/auth/cache.ex
Original file line number Diff line number Diff line change
Expand Up @@ -4,31 +4,19 @@ defmodule Logflare.Auth.Cache do
Cachex `expiration`.
"""

use Logflare.ContextCache

alias Logflare.Auth
alias Logflare.Cache.CachexOps
alias Logflare.OauthAccessTokens.OauthAccessToken
alias Logflare.User
alias Logflare.Utils

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

%{
id: __MODULE__,
start:
{Cachex, :start_link,
[
__MODULE__,
[
hooks:
[
if(stats, do: Utils.cache_stats()),
Utils.cache_limit(100_000)
]
|> Enum.filter(& &1),
expiration: Utils.cache_expiration_min(5, 2)
]
]}
}
CachexOps.child_spec(__MODULE__,
limit: 100_000,
ttl: to_timeout(minute: 5),
purge_interval: to_timeout(minute: 2)
)
end

@spec verify_access_token(OauthAccessToken.t() | String.t()) ::
Expand Down
28 changes: 4 additions & 24 deletions lib/logflare/backends/cache.ex
Original file line number Diff line number Diff line change
@@ -1,33 +1,13 @@
defmodule Logflare.Backends.Cache do
@moduledoc false

use Logflare.ContextCache

alias Logflare.Backends
alias Logflare.Utils
import Cachex.Spec
alias Logflare.Cache.CachexOps

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

%{
id: __MODULE__,
start:
{Cachex, :start_link,
[
__MODULE__,
[
warmers: [
warmer(required: false, module: Backends.CacheWarmer, name: Backends.CacheWarmer)
],
hooks:
[
if(stats, do: Utils.cache_stats()),
Utils.cache_limit(100_000)
]
|> Enum.filter(& &1),
expiration: Utils.cache_expiration_min()
]
]}
}
CachexOps.child_spec(__MODULE__, limit: 100_000, warmer: Backends.CacheWarmer)
end

def list_backends(arg), do: apply_repo_fun(__ENV__.function, [arg])
Expand Down
28 changes: 8 additions & 20 deletions lib/logflare/billing/cache.ex
Original file line number Diff line number Diff line change
@@ -1,29 +1,17 @@
defmodule Logflare.Billing.Cache do
@moduledoc false

use Logflare.ContextCache

alias Logflare.Billing
alias Logflare.Utils
alias Logflare.Cache.CachexOps

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

%{
id: __MODULE__,
start:
{Cachex, :start_link,
[
__MODULE__,
[
hooks:
[
if(stats, do: Utils.cache_stats()),
Utils.cache_limit(100_000)
]
|> Enum.filter(& &1),
expiration: Utils.cache_expiration_min(180, 10)
]
]}
}
CachexOps.child_spec(__MODULE__,
limit: 100_000,
ttl: to_timeout(hour: 3),
purge_interval: to_timeout(minute: 10)
)
end

def get_billing_account_by(keyword) do
Expand Down
52 changes: 52 additions & 0 deletions lib/logflare/cache.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
defmodule Logflare.Cache do
@moduledoc """
Operational contract of an application cache, independent of its storage backend.

`use Logflare.Cache` injects default implementations of every callback, delegating to
`Logflare.Cache.CachexOps`. A cache on another backend overrides them, or passes
`impl: module` naming a `Logflare.Cache.Ops` implementation.
"""

alias Logflare.Cache.CachexOps

@typedoc """
Counters since the last `c:reset/0`. Rates are percentages (0-100); `total_heap_size` is in bytes.
"""
@type stats() :: %{
evictions: non_neg_integer(),
expirations: non_neg_integer(),
operations: non_neg_integer(),
hits: non_neg_integer(),
misses: non_neg_integer(),
hit_rate: number(),
miss_rate: number(),
total_heap_size: non_neg_integer()
}

@doc "Whether the cache on this node can serve requests."
@callback healthy?() :: boolean()

@callback stats() :: stats()

@doc "Clears all entries and statistics."
@callback reset() :: :ok

defmacro __using__(opts) do
impl = Keyword.get(opts, :impl, CachexOps)

quote do
@behaviour Logflare.Cache

@impl Logflare.Cache
def healthy?, do: unquote(impl).healthy?(__MODULE__)

@impl Logflare.Cache
def stats, do: unquote(impl).stats(__MODULE__)

@impl Logflare.Cache
def reset, do: unquote(impl).reset(__MODULE__)

defoverridable healthy?: 0, stats: 0, reset: 0
end
end
end
207 changes: 207 additions & 0 deletions lib/logflare/cache/cachex_ops.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,207 @@
defmodule Logflare.Cache.CachexOps do
@moduledoc """
Cachex implementations of the `Logflare.Cache` and `Logflare.ContextCache` callbacks, and builders
for Cachex start options.

Context cache values are stored as `{:cached, value}`, because Cachex treats a stored `nil` as a miss.
"""

@behaviour Logflare.Cache.Ops
@behaviour Logflare.ContextCache.Ops

import Cachex.Spec

alias Logflare.ContextCache.Gossip

@stat_keys [
:evictions,
:expirations,
:operations,
:hits,
:misses,
:hit_rate,
:miss_rate
]

@type warmer() :: module() | {module(), interval: pos_integer()}
@type cache_opt() ::
{:limit, pos_integer() | nil}
| {:ttl, pos_integer()}
| {:purge_interval, pos_integer()}
| {:warmer, warmer() | nil}
| {:compressed, boolean()}

@doc """
Child spec for a Cachex cache registered as `name`, with the stats hook when `stats_enabled?/0`.

Options:
* `:limit` (required) - maximum number of entries, `nil` for no limit
* `:ttl` - default entry time-to-live, a timeout in ms, defaults to 20 minutes
* `:purge_interval` - timeout in ms between purges of expired entries, defaults to 5 minutes
* `:warmer` - a non-required warmer module, or `{module, interval: ms}`, registered under the module name
* `:compressed` - defaults to `false`
"""
@spec child_spec(atom(), [cache_opt()]) :: Supervisor.child_spec()
def child_spec(name, opts) do
opts = Keyword.update!(cachex_opts(opts), :hooks, &(stats_hooks() ++ &1))
Supervisor.child_spec({Cachex, [name, opts]}, id: name)
end

@doc """
Cachex start options for `t:cache_opt/0`, without the stats hook.
"""
@spec cachex_opts([cache_opt()]) :: keyword()
def cachex_opts(opts) do
[
hooks: limit_hooks(Keyword.fetch!(opts, :limit)),
expiration:
expiration(
default: Keyword.get(opts, :ttl, to_timeout(minute: 20)),
interval: Keyword.get(opts, :purge_interval, to_timeout(minute: 5)),
lazy: true
),
warmers: warmers(opts[:warmer]),
compressed: Keyword.get(opts, :compressed, false)
]
end

@spec stats_enabled?() :: boolean()
def stats_enabled?, do: Application.get_env(:logflare, :cache_stats, false)

@impl Logflare.Cache.Ops
@spec healthy?(Cachex.t()) :: boolean()
def healthy?(cache), do: match?({:ok, _}, Cachex.size(cache))

@doc """
Raises when the cache runs without the stats hook.
"""
@impl Logflare.Cache.Ops
@spec stats(Cachex.t()) :: Logflare.Cache.stats()
def stats(cache) do
{:ok, stats} = Cachex.stats(cache)

{:total_heap_size, heap_words} =
cache
|> Process.whereis()
|> Process.info(:total_heap_size)

heap_bytes = heap_words * :erlang.system_info(:wordsize)

for key <- @stat_keys, into: %{total_heap_size: heap_bytes} do
{key, Map.get(stats, key, 0)}
end
end

@impl Logflare.Cache.Ops
@spec reset(Cachex.t()) :: :ok
def reset(cache) do
{:ok, true} = Cachex.reset(cache, hooks: [Cachex.Stats])
:ok
end

@doc """
With `[id: pkey]`, busts every entry whose cached value is a map with that `:id`, an `{:ok, map}`
with that `:id`, or a list containing such a map. Raises `ArgumentError` for any other keyword.

It scans the cache with a match spec instead of keeping a reverse index of primary keys.
"""
@impl Logflare.ContextCache.Ops
@spec bust_by(Cachex.t(), keyword()) :: {:ok, non_neg_integer()}
def bust_by(cache, id: pkey) do
filter =
{
# use orelse to prevent 2nd condition failing as value is not a map
:orelse,
{
:orelse,
# handle lists
{:is_list, {:element, 2, :value}},
# handle :ok tuples when struct with id is in 2nd element pos.
{:andalso, {:is_tuple, {:element, 2, :value}},
{:andalso, {:==, {:element, 1, {:element, 2, :value}}, :ok},
{:andalso, {:is_map, {:element, 2, {:element, 2, :value}}},
{:==, {:map_get, :id, {:element, 2, {:element, 2, :value}}}, pkey}}}}
},
# handle single maps
{:andalso, {:is_map, {:element, 2, :value}},
{:==, {:map_get, :id, {:element, 2, :value}}, pkey}}
}

query = Cachex.Query.build(where: filter, output: {:key, :value})

keys =
cache
|> Cachex.stream!(query)
|> Stream.filter(fn
{_k, {:cached, v}} when is_list(v) -> Enum.any?(v, &(&1.id == pkey))
{_k, _v} -> true
end)
|> Stream.map(fn {k, _v} -> k end)

delete_keys(cache, keys)
end

def bust_by(cache, kw) do
raise ArgumentError, "#{inspect(cache)} does not support busting by #{inspect(kw)}"
end

@doc """
Returns the cached value for `key`, calling `getter` and caching its result on a miss.

A miss is also multicast to peer nodes, see `Logflare.ContextCache.Gossip`. Accepts a Cachex
worker, so calls can be batched in `Cachex.execute/2`.
"""
@impl Logflare.ContextCache.Ops
@spec fetch(Cachex.t(), term(), (-> term())) :: term()
def fetch(cache, key, getter) do
case Cachex.fetch(cache, key, fn _key -> {:commit, {:cached, getter.()}} end) do
{:commit, {:cached, value}} ->
Gossip.multicast(cache, key, value)
value

{:ok, {:cached, value}} ->
value
end
end

@impl Logflare.ContextCache.Ops
@spec update(Cachex.t(), term(), term()) :: :ok
def update(cache, key, value) do
{:ok, _updated?} = Cachex.update(cache, key, {:cached, value})
:ok
end

@doc """
Deletes `keys` and returns how many of them were present.
"""
@spec delete_keys(Cachex.t(), Enumerable.t()) :: {:ok, non_neg_integer()}
def delete_keys(cache, keys) do
Cachex.execute(cache, fn worker ->
Enum.reduce(keys, 0, fn key, acc -> acc + take_count(worker, key) end)
end)
end

defp take_count(cache, key) do
case Cachex.take(cache, key) do
{:ok, nil} -> 0
{:ok, _value} -> 1
end
end

defp stats_hooks, do: if(stats_enabled?(), do: [hook(module: Cachex.Stats)], else: [])

defp limit_hooks(nil), do: []

defp limit_hooks(limit) when is_integer(limit) do
[hook(module: Cachex.Limit.Scheduled, args: {limit, [], []})]
end

defp warmers(nil), do: []

defp warmers({module, opts}) do
[interval: interval] = Keyword.validate!(opts, interval: nil)
[warmer(module: module, name: module, required: false, interval: interval)]
end

defp warmers(module), do: warmers({module, []})
end
Loading
Loading