471 lines
12 KiB
Elixir
471 lines
12 KiB
Elixir
defmodule Elektrine.DNS.EdgeSites do
|
|
@moduledoc """
|
|
Instance-admin registry of multi-node edge sites and their heartbeat status.
|
|
|
|
## Defaults (no rows)
|
|
|
|
When `dns_edge_sites` is empty, proxied answers use the env/hostname edge
|
|
addresses (`DNS_EDGE_PROXY_*`) as a single synthetic **up** site — no
|
|
heartbeat required.
|
|
|
|
## Heartbeats
|
|
|
|
Edge agents POST to `/_edge/site/v1/heartbeat` with the per-site bearer
|
|
(shown once at mint). Defaults:
|
|
|
|
- interval hint: `DNS_EDGE_SITE_HEARTBEAT_INTERVAL_MS` = 15_000
|
|
- stale → down: `DNS_EDGE_SITE_STALE_AFTER_MS` = 45_000
|
|
|
|
Staleness is evaluated at read time (and optionally persisted by
|
|
`mark_stale_sites/0`) so a missed background sweep cannot leave a site
|
|
permanently "up".
|
|
|
|
## Fail-open
|
|
|
|
If registered **proxy** sites exist but none are effectively up, the full
|
|
proxy-role site IP pool is returned so DNS never empties the edge answer
|
|
set because of a control-plane/heartbeat outage. Sites with only the `dns`
|
|
role are excluded from proxied answers (authority / multi-NS use). Origin
|
|
failover remains `Elektrine.DNS.HealthMonitor` and is intentionally separate.
|
|
Instance-admin operations for edge site registry and per-site bearer auth.
|
|
|
|
Soft dependency for multi-node config distribution (KD-15). Sites are not
|
|
zone-owned; mint/rotate is never available via user PAT scopes.
|
|
"""
|
|
|
|
import Ecto.Query, warn: false
|
|
|
|
alias Elektrine.DNS.EdgeSite
|
|
alias Elektrine.DNS.EdgeSiteCache
|
|
alias Elektrine.Repo
|
|
|
|
@default_heartbeat_interval_ms 15_000
|
|
@default_stale_after_ms 45_000
|
|
@cache_ttl_ms 1_000
|
|
@proxy_role "proxy"
|
|
|
|
@doc "Configured agent heartbeat interval (informational; agents use this)."
|
|
def heartbeat_interval_ms do
|
|
Application.get_env(:elektrine, :dns, [])
|
|
|> Keyword.get(:edge_site_heartbeat_interval_ms, @default_heartbeat_interval_ms)
|
|
end
|
|
|
|
@doc "Age after which a previously-up site is treated as down without heartbeat."
|
|
def stale_after_ms do
|
|
Application.get_env(:elektrine, :dns, [])
|
|
|> Keyword.get(:edge_site_stale_after_ms, @default_stale_after_ms)
|
|
end
|
|
|
|
def list_sites do
|
|
EdgeSite
|
|
|> order_by([s], asc: s.name)
|
|
|> Repo.all()
|
|
rescue
|
|
# Query path must not hard-fail when Repo is unavailable (e.g. unit tests
|
|
# without sandbox ownership). Empty registry → env singleton behaviour.
|
|
_ -> []
|
|
end
|
|
|
|
def get_site(id) when is_integer(id), do: Repo.get(EdgeSite, id)
|
|
def get_site(_), do: nil
|
|
|
|
def get_site!(id), do: Repo.get!(EdgeSite, id)
|
|
|
|
def create_site(attrs) when is_map(attrs) do
|
|
%EdgeSite{}
|
|
|> EdgeSite.changeset(attrs)
|
|
|> Repo.insert()
|
|
|> tap_ok(fn _ -> invalidate_cache() end)
|
|
end
|
|
|
|
def update_site(%EdgeSite{} = site, attrs) when is_map(attrs) do
|
|
site
|
|
|> EdgeSite.changeset(attrs)
|
|
|> Repo.update()
|
|
|> tap_ok(fn _ -> invalidate_cache() end)
|
|
end
|
|
|
|
def delete_site(%EdgeSite{} = site) do
|
|
case Repo.delete(site) do
|
|
{:ok, deleted} ->
|
|
invalidate_cache()
|
|
{:ok, deleted}
|
|
|
|
error ->
|
|
error
|
|
end
|
|
end
|
|
|
|
@doc """
|
|
Mint or rotate the site bearer token.
|
|
|
|
Returns `{:ok, site}` with the raw token in the virtual `:token` field.
|
|
The raw value is never stored; only `token_hash` / `token_prefix` persist.
|
|
|
|
Minting is instance-admin only — never expose this via user PAT scopes.
|
|
"""
|
|
def mint_token(%EdgeSite{} = site) do
|
|
{raw, hash, prefix} = EdgeSite.generate_token()
|
|
|
|
case site
|
|
|> EdgeSite.token_changeset(hash, prefix)
|
|
|> Repo.update() do
|
|
{:ok, updated} ->
|
|
invalidate_cache()
|
|
{:ok, %{updated | token: raw}}
|
|
|
|
error ->
|
|
error
|
|
end
|
|
end
|
|
|
|
def mint_token(_), do: {:error, :not_found}
|
|
|
|
@doc """
|
|
Record a successful heartbeat for the site identified by the raw bearer.
|
|
"""
|
|
def record_heartbeat(raw_token) when is_binary(raw_token) do
|
|
case get_site_by_token(raw_token) do
|
|
%EdgeSite{} = site ->
|
|
case site
|
|
|> EdgeSite.heartbeat_changeset()
|
|
|> Repo.update() do
|
|
{:ok, updated} ->
|
|
invalidate_cache()
|
|
{:ok, updated}
|
|
|
|
error ->
|
|
error
|
|
end
|
|
|
|
nil ->
|
|
{:error, :unauthorized}
|
|
end
|
|
end
|
|
|
|
def record_heartbeat(_), do: {:error, :unauthorized}
|
|
|
|
def get_site_by_token(raw_token) when is_binary(raw_token) and raw_token != "" do
|
|
hash = EdgeSite.hash_token(raw_token)
|
|
|
|
EdgeSite
|
|
|> where([s], s.token_hash == ^hash)
|
|
|> Repo.one()
|
|
end
|
|
|
|
def get_site_by_token(_), do: nil
|
|
|
|
@doc """
|
|
Persist `status=down` for sites whose last heartbeat is older than
|
|
`stale_after_ms/0`. Returns the number of rows updated.
|
|
"""
|
|
def mark_stale_sites(now \\ DateTime.utc_now()) do
|
|
cutoff =
|
|
now
|
|
|> DateTime.add(-div(stale_after_ms(), 1000), :second)
|
|
|> DateTime.truncate(:second)
|
|
|
|
{count, _} =
|
|
from(s in EdgeSite,
|
|
where: s.status == "up",
|
|
where: is_nil(s.last_heartbeat_at) or s.last_heartbeat_at < ^cutoff
|
|
)
|
|
|> Repo.update_all(set: [status: "down", updated_at: DateTime.truncate(now, :second)])
|
|
|
|
if count > 0, do: invalidate_cache()
|
|
count
|
|
end
|
|
|
|
@doc """
|
|
Whether a site should contribute IPs to proxied answers right now.
|
|
|
|
Uses wall-clock freshness against `last_heartbeat_at`, not only the
|
|
stored `status` column (covers the window before `mark_stale_sites/0`).
|
|
"""
|
|
def effectively_up?(site, now \\ DateTime.utc_now())
|
|
|
|
def effectively_up?(%EdgeSite{} = site, now) do
|
|
case site do
|
|
%{status: "up", last_heartbeat_at: %DateTime{} = at} ->
|
|
DateTime.diff(now, at, :millisecond) < stale_after_ms()
|
|
|
|
_ ->
|
|
false
|
|
end
|
|
end
|
|
|
|
def effectively_up?(_, _), do: false
|
|
|
|
@doc """
|
|
Filter env edge addresses through the site registry.
|
|
|
|
Only sites with the `"proxy"` role contribute IPs to proxied answers.
|
|
Pure `"dns"` sites (authority-only) are ignored here.
|
|
|
|
- No proxy-role registry rows → return `env_addresses` unchanged (synthetic up site).
|
|
- Some proxy sites up → union of those addresses for `family` (`:ipv4` / `:ipv6`).
|
|
- All proxy sites down/stale → fail-open with the full proxy-role pool for `family`.
|
|
- Filtered pool empty for this family → fall back to `env_addresses`.
|
|
"""
|
|
def filter_edge_addresses(family, env_addresses)
|
|
when family in [:ipv4, :ipv6] and is_list(env_addresses) do
|
|
proxy_sites = cached_sites() |> Enum.filter(&proxy_role?/1)
|
|
|
|
if proxy_sites == [] do
|
|
env_addresses
|
|
else
|
|
now = DateTime.utc_now()
|
|
up = Enum.filter(proxy_sites, &effectively_up?(&1, now))
|
|
source = if up == [], do: proxy_sites, else: up
|
|
|
|
addresses =
|
|
source
|
|
|> Enum.flat_map(fn site ->
|
|
case family do
|
|
:ipv4 -> site.ipv4
|
|
:ipv6 -> site.ipv6
|
|
end
|
|
end)
|
|
|> Enum.uniq()
|
|
|
|
if addresses == [], do: env_addresses, else: addresses
|
|
end
|
|
end
|
|
|
|
def filter_edge_addresses(_family, env_addresses), do: List.wrap(env_addresses)
|
|
|
|
@doc false
|
|
def proxy_role?(%EdgeSite{roles: roles}) when is_list(roles), do: @proxy_role in roles
|
|
def proxy_role?(_), do: false
|
|
|
|
@doc false
|
|
def invalidate_cache do
|
|
EdgeSiteCache.ensure_table()
|
|
:ets.delete(EdgeSiteCache.table(), :sites)
|
|
:ok
|
|
rescue
|
|
ArgumentError -> :ok
|
|
end
|
|
|
|
defp cached_sites do
|
|
EdgeSiteCache.ensure_table()
|
|
now = System.monotonic_time(:millisecond)
|
|
table = EdgeSiteCache.table()
|
|
|
|
case :ets.lookup(table, :sites) do
|
|
[{:sites, %{expires_at: exp, list: list}}] when is_integer(exp) and exp > now ->
|
|
list
|
|
|
|
_ ->
|
|
list = list_sites()
|
|
:ets.insert(table, {:sites, %{list: list, expires_at: now + @cache_ttl_ms}})
|
|
list
|
|
end
|
|
rescue
|
|
_ ->
|
|
[]
|
|
end
|
|
|
|
defp tap_ok({:ok, value}, fun) do
|
|
fun.(value)
|
|
{:ok, value}
|
|
end
|
|
|
|
defp tap_ok(other, _fun), do: other
|
|
alias Elektrine.Repo
|
|
|
|
@doc """
|
|
List non-revoked edge sites (newest first by name).
|
|
"""
|
|
def list_sites do
|
|
EdgeSite
|
|
|> where([s], s.status != "revoked")
|
|
|> order_by([s], asc: s.name)
|
|
|> Repo.all()
|
|
end
|
|
|
|
def get_site(id) when is_integer(id), do: Repo.get(EdgeSite, id)
|
|
|
|
def get_site(id) when is_binary(id) do
|
|
case Integer.parse(id) do
|
|
{int, ""} -> Repo.get(EdgeSite, int)
|
|
_ -> nil
|
|
end
|
|
end
|
|
|
|
def get_site(_), do: nil
|
|
|
|
@doc """
|
|
Create a site and return the raw bearer token + HMAC secret once.
|
|
"""
|
|
def create_site(attrs) when is_map(attrs) do
|
|
{raw, hash, prefix} = EdgeSite.generate_token()
|
|
hmac = EdgeSite.generate_hmac_secret()
|
|
|
|
attrs =
|
|
attrs
|
|
|> stringify_keys()
|
|
|> Map.put("token_hash", hash)
|
|
|> Map.put("token_prefix", prefix)
|
|
|> Map.put("hmac_secret", hmac)
|
|
|> Map.put_new("status", "unknown")
|
|
|
|
%EdgeSite{}
|
|
|> EdgeSite.changeset(attrs)
|
|
|> Repo.insert()
|
|
|> case do
|
|
{:ok, site} ->
|
|
{:ok, %{site | token: raw, hmac_secret_plain: hmac}}
|
|
|
|
other ->
|
|
other
|
|
end
|
|
end
|
|
|
|
def update_site(%EdgeSite{} = site, attrs) when is_map(attrs) do
|
|
if EdgeSite.revoked?(site) do
|
|
{:error, :revoked}
|
|
else
|
|
site
|
|
|> EdgeSite.changeset(stringify_keys(attrs))
|
|
|> Repo.update()
|
|
end
|
|
end
|
|
|
|
@doc """
|
|
Rotate the site bearer (and HMAC secret). Returns raw secrets once.
|
|
"""
|
|
def rotate_token(%EdgeSite{} = site) do
|
|
if EdgeSite.revoked?(site) do
|
|
{:error, :revoked}
|
|
else
|
|
{raw, hash, prefix} = EdgeSite.generate_token()
|
|
hmac = EdgeSite.generate_hmac_secret()
|
|
|
|
site
|
|
|> EdgeSite.changeset(%{
|
|
"token_hash" => hash,
|
|
"token_prefix" => prefix,
|
|
"hmac_secret" => hmac
|
|
})
|
|
|> Repo.update()
|
|
|> case do
|
|
{:ok, updated} ->
|
|
{:ok, %{updated | token: raw, hmac_secret_plain: hmac}}
|
|
|
|
other ->
|
|
other
|
|
end
|
|
end
|
|
end
|
|
|
|
def revoke_site(%EdgeSite{} = site) do
|
|
site
|
|
|> EdgeSite.revoke_changeset()
|
|
|> Repo.update()
|
|
end
|
|
|
|
@doc """
|
|
Resolve a long-lived site bearer (`ess_...`) to an active site row.
|
|
"""
|
|
def get_site_by_token(raw_token) when is_binary(raw_token) do
|
|
hash = EdgeSite.hash_token(String.trim(raw_token))
|
|
|
|
case Repo.get_by(EdgeSite, token_hash: hash) do
|
|
%EdgeSite{} = site ->
|
|
if EdgeSite.revoked?(site), do: {:error, :revoked}, else: {:ok, site}
|
|
|
|
nil ->
|
|
{:error, :not_found}
|
|
end
|
|
end
|
|
|
|
def get_site_by_token(_), do: {:error, :not_found}
|
|
|
|
def record_heartbeat(%EdgeSite{} = site) do
|
|
site
|
|
|> EdgeSite.heartbeat_changeset()
|
|
|> Repo.update()
|
|
end
|
|
|
|
@doc """
|
|
Public representation without secrets.
|
|
"""
|
|
def public_site(%EdgeSite{} = site) do
|
|
%{
|
|
id: site.id,
|
|
name: site.name,
|
|
ipv4: site.ipv4 || [],
|
|
ipv6: site.ipv6 || [],
|
|
roles: site.roles || ["proxy"],
|
|
status: site.status,
|
|
last_heartbeat_at: site.last_heartbeat_at,
|
|
token_prefix: site.token_prefix,
|
|
metadata: site.metadata || %{},
|
|
revoked_at: site.revoked_at,
|
|
inserted_at: site.inserted_at,
|
|
updated_at: site.updated_at
|
|
}
|
|
end
|
|
|
|
def public_site(%EdgeSite{} = site, raw_token) when is_binary(raw_token) do
|
|
site
|
|
|> public_site()
|
|
|> Map.put(:token, raw_token)
|
|
|> Map.put(:hmac_secret, site.hmac_secret_plain)
|
|
end
|
|
|
|
@doc """
|
|
Edge IPs for the pool: union of up sites, or env singleton when no rows.
|
|
|
|
Fail-open: if every registered site is down/unknown, return all site IPs plus
|
|
env addresses so proxied answers do not empty.
|
|
"""
|
|
def edge_pool_ips do
|
|
sites = list_sites()
|
|
|
|
case sites do
|
|
[] ->
|
|
env_edge_ips()
|
|
|
|
list ->
|
|
up = Enum.filter(list, &(&1.status == "up"))
|
|
chosen = if up == [], do: list, else: up
|
|
|
|
%{
|
|
ipv4: chosen |> Enum.flat_map(&(&1.ipv4 || [])) |> Enum.uniq(),
|
|
ipv6: chosen |> Enum.flat_map(&(&1.ipv6 || [])) |> Enum.uniq()
|
|
}
|
|
|> merge_env_if_empty()
|
|
end
|
|
end
|
|
|
|
def env_edge_ips do
|
|
dns = Application.get_env(:elektrine, :dns, [])
|
|
|
|
%{
|
|
ipv4: Keyword.get(dns, :edge_proxy_ipv4_addresses, []) || [],
|
|
ipv6: Keyword.get(dns, :edge_proxy_ipv6_addresses, []) || []
|
|
}
|
|
end
|
|
|
|
def stale_after_ms do
|
|
Application.get_env(:elektrine, :dns, [])
|
|
|> Keyword.get(:edge_site_stale_after_ms, 45_000)
|
|
end
|
|
|
|
def heartbeat_interval_ms do
|
|
Application.get_env(:elektrine, :dns, [])
|
|
|> Keyword.get(:edge_site_heartbeat_interval_ms, 15_000)
|
|
end
|
|
|
|
defp merge_env_if_empty(%{ipv4: [], ipv6: []} = _pool), do: env_edge_ips()
|
|
defp merge_env_if_empty(pool), do: pool
|
|
|
|
defp stringify_keys(attrs) when is_map(attrs) do
|
|
Map.new(attrs, fn
|
|
{k, v} when is_atom(k) -> {Atom.to_string(k), v}
|
|
{k, v} -> {k, v}
|
|
end)
|
|
end
|
|
end
|