merge: pr-11c-featdns-secondary-ingest-answer-from-transferred-d

This commit is contained in:
maxfield 2026-08-02 01:54:37 -04:00
commit f315c10ff4
15 changed files with 1590 additions and 7 deletions

View file

@ -0,0 +1,11 @@
defmodule Elektrine.Repo.Migrations.AddSecondaryZoneMasters do
use Ecto.Migration
def change do
alter table(:dns_zones) do
# Master endpoints for kind=secondary (host or host:port, one per entry).
add :masters, {:array, :string}, null: false, default: []
add :secondary_last_transferred_at, :utc_datetime
end
end
end

View file

@ -17,6 +17,7 @@ defmodule Elektrine.DNS do
alias Elektrine.DNS.QueryStat
alias Elektrine.DNS.Record
alias Elektrine.DNS.Tunnels
alias Elektrine.DNS.SecondaryStore
alias Elektrine.DNS.Zone
alias Elektrine.DNS.ZoneCache
alias Elektrine.DNS.ZoneServiceConfig
@ -424,7 +425,8 @@ defmodule Elektrine.DNS do
def create_record(%Zone{} = zone, attrs) when is_map(attrs) do
normalized_attrs = normalize_record_attrs(attrs, zone.domain)
with :ok <- validate_zone_record_write(zone, normalized_attrs) do
with :ok <- validate_secondary_zone_not_writable(zone),
:ok <- validate_zone_record_write(zone, normalized_attrs) do
zone_domain = zone.domain
%Record{}
@ -452,12 +454,19 @@ defmodule Elektrine.DNS do
{:error, add_error(change(zone), :domain, "is managed by Elektrine")}
else
with :ok <- validate_axfr_tsig_key_for_zone(zone, public_attrs) do
previous = zone
zone
|> Zone.changeset(public_attrs)
|> Repo.update()
|> case do
{:ok, zone} -> {:ok, Repo.preload(zone, [:records, :service_configs, :axfr_tsig_key])}
error -> error
{:ok, updated} ->
clear_secondary_store_after_zone_change(previous, updated)
{:ok, Repo.preload(updated, [:records, :service_configs, :axfr_tsig_key])}
error ->
error
end
|> refresh_authority_cache_after_write(touch_zone: true)
end
@ -483,6 +492,10 @@ defmodule Elektrine.DNS do
{:error,
add_error(change(zone), :domain, "is the built-in profile subdomain and cannot be deleted")}
else
# Drop any transferred secondary snapshot so a recreate cannot answer
# from stale ETS data until the next successful AXFR.
clear_secondary_store(zone.domain)
zone
|> Repo.delete()
|> refresh_authority_cache_after_write()
@ -504,7 +517,8 @@ defmodule Elektrine.DNS do
zone = Repo.get(Zone, record.zone_id)
normalized_attrs = normalize_record_attrs(public_record_attrs(attrs), zone_domain)
with :ok <- validate_record_mutation(record, :update),
with :ok <- validate_secondary_zone_not_writable(zone),
:ok <- validate_record_mutation(record, :update),
:ok <-
validate_zone_record_write(
zone,
@ -519,7 +533,10 @@ defmodule Elektrine.DNS do
end
def delete_record(%Record{} = record) do
with :ok <- validate_record_mutation(record, :delete) do
zone = Repo.get(Zone, record.zone_id)
with :ok <- validate_secondary_zone_not_writable(zone),
:ok <- validate_record_mutation(record, :delete) do
record
|> Repo.delete()
|> refresh_authority_cache_after_write(touch_zone: true)
@ -949,6 +966,27 @@ defmodule Elektrine.DNS do
|> Keyword.get(:secondary_axfr_enabled, false)
end
@doc """
Global switch for inbound secondary AXFR ingest.
Defaults false (`DNS_SECONDARY_INGEST_ENABLED`). When true, the
`SecondaryIngest` worker transfers `kind=secondary` zones from their
masters into `SecondaryStore` and authority answers from that data
without writing `dns_records`.
"""
def secondary_ingest_enabled? do
Application.get_env(:elektrine, :dns, [])
|> Keyword.get(:secondary_ingest_enabled, false)
end
@doc """
Poll interval for secondary ingest cycles (milliseconds).
"""
def secondary_ingest_poll_interval_ms do
Application.get_env(:elektrine, :dns, [])
|> Keyword.get(:secondary_ingest_poll_interval_ms, 60_000)
end
def alias_resolver do
Application.get_env(:elektrine, :dns, [])
|> Keyword.get(:alias_resolver, :inet_res)
@ -2323,7 +2361,10 @@ defmodule Elektrine.DNS do
"axfr_enabled",
"axfr_allow_cidrs",
"axfr_require_tsig",
"axfr_tsig_key_id"
"axfr_tsig_key_id",
# masters for kind=secondary; serial / secondary_last_transferred_at are
# ingest-owned (SecondaryIngest.persist_serial/4) and not public attrs.
"masters"
])
|> normalize_public_zone_attrs()
end
@ -2334,6 +2375,27 @@ defmodule Elektrine.DNS do
|> normalize_boolean_attr("axfr_require_tsig")
|> normalize_axfr_allow_cidrs_attr()
|> normalize_axfr_tsig_key_id_attr()
|> normalize_masters_attr()
end
defp normalize_masters_attr(attrs) do
case Map.fetch(attrs, "masters") do
{:ok, value} when is_binary(value) ->
Map.put(
attrs,
"masters",
String.split(value, [",", "\n", " ", "\t"], trim: true)
)
{:ok, value} when is_list(value) ->
Map.put(attrs, "masters", value)
{:ok, nil} ->
Map.put(attrs, "masters", [])
_ ->
attrs
end
end
defp normalize_boolean_attr(attrs, key) do
@ -2459,6 +2521,48 @@ defmodule Elektrine.DNS do
defp normalize_record_name(name, _zone_domain), do: name
defp validate_secondary_zone_not_writable(%Zone{kind: "secondary"}) do
{:error,
add_error(
change(%Record{}),
:zone_id,
"secondary zones answer from transferred data; records are not writable"
)}
end
defp validate_secondary_zone_not_writable(_), do: :ok
defp clear_secondary_store(domain) when is_binary(domain) do
SecondaryStore.delete(domain)
drop_zone_cache_entry(domain)
:ok
end
defp clear_secondary_store(_), do: :ok
# When a zone leaves secondary kind or renames, drop the old snapshot so
# authority cannot keep answering from stale transferred data.
defp clear_secondary_store_after_zone_change(%Zone{} = previous, %Zone{} = updated) do
if previous.kind == "secondary" and
(updated.kind != "secondary" or
String.downcase(previous.domain || "") != String.downcase(updated.domain || "")) do
clear_secondary_store(previous.domain)
end
:ok
end
defp drop_zone_cache_entry(domain) when is_binary(domain) do
key = domain |> String.downcase() |> String.trim_trailing(".")
try do
:ets.delete(ZoneCache, key)
:ok
rescue
ArgumentError -> :ok
end
end
defp validate_zone_record_write(%Zone{} = zone, attrs) when is_map(attrs) do
with :ok <- validate_cname_exclusivity(zone, attrs) do
if builtin_user_zone?(zone) and builtin_user_zone_hosted_by_platform?(zone) do

View file

@ -24,6 +24,7 @@ Elektrine.DNS.EdgeCache,
defp authority_children do
if Elektrine.DNS.authority_enabled?() do
[
Elektrine.DNS.SecondaryStore,
Elektrine.DNS.ZoneCache,
Elektrine.DNS.ZoneChangeListener,
Elektrine.DNS.HealthMonitor,
@ -33,7 +34,15 @@ Elektrine.DNS.EdgeCache,
{Elektrine.DNS.UDPServer, name: Elektrine.DNS.UDPServerV6, family: :inet6},
{Elektrine.DNS.TCPServer, name: Elektrine.DNS.TCPServer, family: :inet},
{Elektrine.DNS.TCPServer, name: Elektrine.DNS.TCPServerV6, family: :inet6}
]
] ++ secondary_ingest_children()
else
[]
end
end
defp secondary_ingest_children do
if Elektrine.DNS.secondary_ingest_enabled?() do
[Elektrine.DNS.SecondaryIngest]
else
[]
end

View file

@ -0,0 +1,466 @@
defmodule Elektrine.DNS.AxfrClient do
@moduledoc """
Inbound AXFR client used by secondary ingest (RFC 5936).
Connects to a master over TCP, streams the zone transfer, and returns
decoded resource records. Unsigned first no TSIG MAC on the wire yet.
"""
import Bitwise
alias Elektrine.DNS.Packet
@default_port 53
@default_timeout_ms 10_000
# Hard cap on messages to bound memory on a hostile/malformed stream.
@max_messages 10_000
@type master :: String.t() | {tuple() | String.t(), pos_integer()}
@type transfer_result :: %{
serial: non_neg_integer(),
soa: map(),
records: [map()],
messages: non_neg_integer()
}
@doc """
Perform an AXFR of `domain` from `master`.
`master` is `"host"`, `"host:port"`, or `{ip_tuple_or_host, port}`.
"""
@spec transfer(String.t(), master(), keyword()) ::
{:ok, transfer_result()} | {:error, term()}
def transfer(domain, master, opts \\ []) when is_binary(domain) do
timeout = Keyword.get(opts, :timeout, @default_timeout_ms)
with {:ok, {address, port}} <- parse_master(master),
{:ok, ip} <- resolve_address(address) do
case :gen_tcp.connect(ip, port, [:binary, active: false, packet: 0], timeout) do
{:ok, socket} ->
# Always close after connect so send/recv/parse failures cannot leak FDs.
try do
do_transfer(socket, domain, timeout)
after
:gen_tcp.close(socket)
end
{:error, reason} ->
{:error, reason}
end
end
rescue
error -> {:error, error}
end
defp do_transfer(socket, domain, timeout) do
query = build_axfr_query(domain)
packet = Packet.encode_query(query)
case :gen_tcp.send(socket, <<byte_size(packet)::16, packet::binary>>) do
:ok -> recv_axfr_stream(socket, domain, timeout)
{:error, reason} -> {:error, reason}
end
end
@doc """
Parse a multi-message AXFR byte stream (list of full DNS messages) into
transfer records. Used by the live client and by tests with fixtures.
"""
@spec parse_axfr_messages([binary()], String.t()) ::
{:ok, transfer_result()} | {:error, term()}
def parse_axfr_messages(messages, domain) when is_list(messages) and is_binary(domain) do
domain = normalize_name(domain)
case collect_answers(messages, domain) do
{:ok, answers} ->
finalize_transfer(answers, domain)
{:error, _} = error ->
error
end
end
@doc """
Normalize a master endpoint string to `{host, port}`.
"""
@spec parse_master(master()) :: {:ok, {String.t() | tuple(), pos_integer()}} | {:error, term()}
def parse_master({host, port}) when is_integer(port) and port > 0 and port <= 65_535 do
{:ok, {host, port}}
end
def parse_master(endpoint) when is_binary(endpoint) do
trimmed = String.trim(endpoint)
cond do
trimmed == "" ->
{:error, :empty_master}
match?([_, _], Regex.run(~r/^\[(.+)\]:(\d+)$/, trimmed)) ->
[_, host, port] = Regex.run(~r/^\[(.+)\]:(\d+)$/, trimmed)
parse_port_pair(host, port)
match?([_, _], Regex.run(~r/^\[(.+)\]$/, trimmed)) ->
[_, host] = Regex.run(~r/^\[(.+)\]$/, trimmed)
{:ok, {host, @default_port}}
String.contains?(trimmed, ":") and not String.contains?(trimmed, "::") ->
case String.split(trimmed, ":", parts: 2) do
[host, port] -> parse_port_pair(host, port)
_ -> {:error, :invalid_master}
end
true ->
{:ok, {trimmed, @default_port}}
end
end
def parse_master(_), do: {:error, :invalid_master}
defp parse_port_pair(host, port_str) do
case Integer.parse(to_string(port_str)) do
{port, ""} when port > 0 and port <= 65_535 ->
{:ok, {host, port}}
_ ->
{:error, :invalid_port}
end
end
defp resolve_address(ip) when is_tuple(ip), do: {:ok, ip}
defp resolve_address(host) when is_binary(host) do
case :inet.parse_address(String.to_charlist(host)) do
{:ok, ip} ->
{:ok, ip}
{:error, _} ->
case :inet.getaddr(String.to_charlist(host), :inet) do
{:ok, ip} -> {:ok, ip}
{:error, reason} -> {:error, reason}
end
end
end
defp build_axfr_query(domain) do
%{
id: :erlang.unique_integer([:positive]) &&& 0xFFFF,
rd: 0,
qname: normalize_name(domain),
qtype: :axfr,
edns: false,
udp_size: 512
}
end
defp recv_axfr_stream(socket, domain, timeout) do
recv_messages(socket, domain, timeout, [], 0, false, nil)
end
defp recv_messages(_socket, _domain, _timeout, _acc, count, _seen_soa, _serial)
when count >= @max_messages do
{:error, :too_many_messages}
end
defp recv_messages(socket, domain, timeout, acc, count, seen_soa, serial) do
with {:ok, <<length::16>>} <- :gen_tcp.recv(socket, 2, timeout),
true <- length > 0 and length <= 65_535,
{:ok, message} <- :gen_tcp.recv(socket, length, timeout) do
case :inet_dns.decode(message) do
{:ok, decoded} ->
handle_decoded(socket, domain, timeout, acc, count, seen_soa, serial, message, decoded)
_ ->
{:error, :decode_failed}
end
else
false -> {:error, :invalid_message_length}
{:error, reason} -> {:error, reason}
end
end
defp handle_decoded(socket, domain, timeout, acc, count, seen_soa, serial, message, decoded) do
rcode = response_rcode(decoded)
if rcode != 0 do
{:error, {:rcode, rcode}}
else
answers = answer_records(decoded)
{new_seen, new_serial, finished?} = track_soa(answers, seen_soa, serial)
acc = acc ++ [message]
if finished? do
parse_axfr_messages(acc, domain)
else
recv_messages(socket, domain, timeout, acc, count + 1, new_seen, new_serial)
end
end
end
defp track_soa(answers, seen_soa, serial) do
Enum.reduce(answers, {seen_soa, serial, false}, fn rr, {seen, ser, done} ->
if done do
{seen, ser, true}
else
case rr_type(rr) do
:soa ->
this_serial = soa_serial(rr)
cond do
not seen ->
{true, this_serial, false}
ser == this_serial ->
{true, ser, true}
true ->
# Different serial mid-stream is unexpected; keep reading.
{true, ser, false}
end
_ ->
{seen, ser, false}
end
end
end)
end
defp collect_answers(messages, _domain) do
Enum.reduce_while(messages, {:ok, []}, fn message, {:ok, acc} ->
case :inet_dns.decode(message) do
{:ok, decoded} ->
rcode = response_rcode(decoded)
if rcode == 0 do
{:cont, {:ok, acc ++ answer_records(decoded)}}
else
{:halt, {:error, {:rcode, rcode}}}
end
_ ->
{:halt, {:error, :decode_failed}}
end
end)
end
defp finalize_transfer([], _domain), do: {:error, :empty_transfer}
defp finalize_transfer(answers, domain) do
with {:ok, first_soa, middle, last_soa} <- split_soa_bookends(answers),
serial when is_integer(serial) <- soa_serial(first_soa),
true <- serial == soa_serial(last_soa) do
soa = convert_soa(first_soa, domain)
records = Enum.map(middle, &convert_rr(&1, domain)) |> Enum.reject(&is_nil/1)
{:ok,
%{
serial: serial,
soa: soa,
records: records,
messages: 1
}}
else
false -> {:error, :soa_serial_mismatch}
{:error, _} = error -> error
_ -> {:error, :invalid_transfer}
end
end
defp split_soa_bookends(answers) do
case answers do
[first | rest] ->
if rr_type(first) != :soa do
{:error, :missing_leading_soa}
else
case Enum.split_while(rest, &(rr_type(&1) != :soa)) do
{middle, [last | _]} ->
{:ok, first, middle, last}
{_middle, []} ->
{:error, :missing_trailing_soa}
end
end
[] ->
{:error, :empty_transfer}
end
end
defp convert_soa(rr, domain) do
data = rr_data(rr)
{mname, rname, serial, refresh, retry, expire, minimum} = extract_soa_tuple(data)
%{
type: "SOA",
host: normalize_name(rr_domain(rr) || domain),
ttl: rr_ttl(rr),
mname: normalize_name(mname),
rname: normalize_name(rname),
serial: serial,
refresh: refresh,
retry: retry,
expire: expire,
minimum: minimum
}
end
defp convert_rr(rr, domain) do
type = rr_type(rr)
host = normalize_name(rr_domain(rr) || domain)
ttl = rr_ttl(rr)
data = rr_data(rr)
base = %{host: host, type: type_string(type), ttl: ttl}
case type do
:a ->
Map.put(base, :content, format_ip(data))
:aaaa ->
Map.put(base, :content, format_ip(data))
:ns ->
Map.put(base, :value, normalize_name(data))
:cname ->
Map.put(base, :content, normalize_name(data))
:mx ->
{pref, exchange} = extract_mx(data)
base
|> Map.put(:content, normalize_name(exchange))
|> Map.put(:priority, pref)
:txt ->
Map.put(base, :content, format_txt(data))
:srv ->
{priority, weight, port, target} = extract_srv(data)
base
|> Map.put(:content, normalize_name(target))
|> Map.put(:priority, priority)
|> Map.put(:weight, weight)
|> Map.put(:port, port)
:caa ->
{flags, tag, value} = extract_caa(data)
base
|> Map.put(:flags, flags)
|> Map.put(:tag, tag)
|> Map.put(:content, value)
:soa ->
# Middle SOA should not appear; drop if it does.
nil
other when is_atom(other) ->
Map.put(base, :content, to_string(data))
_ ->
nil
end
end
defp extract_soa_tuple(data) do
case data do
{mname, rname, serial, refresh, retry, expire, minimum} ->
{mname, rname, serial, refresh, retry, expire, minimum}
[mname, rname, serial, refresh, retry, expire, minimum] ->
{mname, rname, serial, refresh, retry, expire, minimum}
other when is_tuple(other) and tuple_size(other) >= 7 ->
{elem(other, 0), elem(other, 1), elem(other, 2), elem(other, 3), elem(other, 4),
elem(other, 5), elem(other, 6)}
_ ->
{"ns", "hostmaster", 0, 3600, 600, 1_209_600, 300}
end
end
defp extract_mx({pref, exchange}), do: {pref, exchange}
defp extract_mx([pref, exchange]), do: {pref, exchange}
defp extract_mx(other), do: {10, other}
defp extract_srv({priority, weight, port, target}), do: {priority, weight, port, target}
defp extract_srv([priority, weight, port, target]), do: {priority, weight, port, target}
defp extract_srv(other), do: {0, 0, 0, other}
defp extract_caa({flags, tag, value}), do: {flags, to_string(tag), to_string(value)}
defp extract_caa([flags, tag, value]), do: {flags, to_string(tag), to_string(value)}
defp extract_caa(other), do: {0, "issue", to_string(other)}
defp format_ip({a, b, c, d}), do: "#{a}.#{b}.#{c}.#{d}"
defp format_ip(tuple) when is_tuple(tuple) and tuple_size(tuple) == 8 do
tuple |> :inet.ntoa() |> to_string()
end
defp format_ip(bin) when is_binary(bin) and byte_size(bin) == 4 do
<<a, b, c, d>> = bin
"#{a}.#{b}.#{c}.#{d}"
end
defp format_ip(other), do: to_string(other)
defp format_txt(list) when is_list(list) do
Enum.map_join(list, fn
bin when is_binary(bin) -> bin
chars when is_list(chars) -> List.to_string(chars)
other -> to_string(other)
end)
end
defp format_txt(bin) when is_binary(bin), do: bin
defp format_txt(other), do: to_string(other)
defp type_string(type) when is_atom(type), do: type |> Atom.to_string() |> String.upcase()
defp type_string(type), do: to_string(type)
# :inet_dns record accessors (OTP version tolerant).
defp answer_records({:dns_rec, _header, _qd, an, _ns, _ar}), do: an || []
defp answer_records(rec) when is_tuple(rec) and tuple_size(rec) >= 4, do: elem(rec, 3) || []
defp answer_records(_), do: []
defp response_rcode({:dns_rec, header, _, _, _, _}), do: header_rcode(header)
defp response_rcode(rec) when is_tuple(rec), do: header_rcode(elem(rec, 1))
defp response_rcode(_), do: 2
defp header_rcode({:dns_header, _, _, _, _, _, _, _, _, rcode}), do: rcode
defp header_rcode(header) when is_tuple(header) and tuple_size(header) >= 10,
do: elem(header, 9)
defp header_rcode(_), do: 2
defp rr_type(rr) when is_tuple(rr) and tuple_size(rr) >= 3, do: elem(rr, 2)
defp rr_type(_), do: nil
defp rr_domain(rr) when is_tuple(rr) and tuple_size(rr) >= 2, do: elem(rr, 1)
defp rr_domain(_), do: nil
defp rr_ttl(rr) when is_tuple(rr) and tuple_size(rr) >= 5, do: elem(rr, 4)
defp rr_ttl(_), do: 300
defp rr_data(rr) when is_tuple(rr) and tuple_size(rr) >= 7, do: elem(rr, 6)
defp rr_data(_), do: nil
defp soa_serial(rr) do
case extract_soa_tuple(rr_data(rr)) do
{_, _, serial, _, _, _, _} when is_integer(serial) -> serial
_ -> nil
end
end
defp normalize_name(nil), do: ""
defp normalize_name(name) do
name
|> to_string()
|> String.downcase()
|> String.trim_trailing(".")
end
end

View file

@ -576,6 +576,12 @@ defmodule Elektrine.DNS.Query do
|> Enum.reject(&(Elektrine.DNS.Record.private?(&1) and not private_query?(opts)))
end
defp all_zone_records(%{kind: "secondary"} = zone) do
# Secondary answers only from transferred RRs (SecondaryStore), never
# primary onboarding bootstrap or live primary dns_records writes.
zone.records || []
end
defp all_zone_records(zone) do
bootstrap = DNS.zone_onboarding_records(zone)
persisted = zone.records || []

View file

@ -0,0 +1,319 @@
defmodule Elektrine.DNS.SecondaryIngest do
@moduledoc """
Secondary DNS ingest worker.
For each `kind=secondary` zone with configured masters, performs an inbound
AXFR and stores the result in `Elektrine.DNS.SecondaryStore` so the authority
can answer from transferred data **without writing `dns_records`** (primary
DB RR rows). Only zone metadata (`serial`, SOA timers,
`secondary_last_transferred_at`) is updated in Postgres.
Gated by `DNS_SECONDARY_INGEST_ENABLED` (default false) and requires
authority to be enabled so the secondary can serve answers.
"""
use GenServer
import Ecto.Query, only: [from: 2]
require Logger
alias Elektrine.DNS
alias Elektrine.DNS.AxfrClient
alias Elektrine.DNS.SecondaryStore
alias Elektrine.DNS.Zone
alias Elektrine.DNS.ZoneCache
alias Elektrine.Repo
@default_poll_ms 60_000
@min_poll_ms 5_000
def start_link(opts \\ []) do
GenServer.start_link(__MODULE__, opts, name: __MODULE__)
end
@doc """
Trigger an immediate ingest cycle (all secondary zones).
"""
@spec refresh(keyword()) :: :ok | {:error, term()}
def refresh(opts \\ []) do
GenServer.call(__MODULE__, {:refresh, opts}, Keyword.get(opts, :timeout, 60_000))
catch
:exit, {:noproc, _} -> {:error, :not_started}
:exit, {:timeout, _} -> {:error, :timeout}
end
@doc """
Transfer a single secondary zone from its masters into SecondaryStore.
Does **not** insert or update `dns_records`. Updates zone serial metadata
only when `persist?: true` (default).
"""
@spec transfer_zone(Zone.t(), keyword()) ::
{:ok, map()} | {:error, term()}
def transfer_zone(%Zone{} = zone, opts \\ []) do
if zone.kind != "secondary" do
{:error, :not_secondary}
else
masters = zone.masters || []
if masters == [] do
{:error, :no_masters}
else
do_transfer_zone(zone, masters, opts)
end
end
end
@doc """
Apply an already-parsed transfer result into SecondaryStore and optionally
persist serial metadata. Used by tests and by `transfer_zone/2`.
"""
@spec apply_transfer(Zone.t(), map(), keyword()) :: {:ok, map()} | {:error, term()}
def apply_transfer(%Zone{} = zone, %{serial: serial, records: records} = result, opts \\ [])
when is_integer(serial) and is_list(records) do
persist? = Keyword.get(opts, :persist?, true)
now = DateTime.utc_now() |> DateTime.truncate(:second)
soa = Map.get(result, :soa)
SecondaryStore.put(%{
domain: zone.domain,
serial: serial,
records: records,
soa: soa,
transferred_at: now
})
zone =
if persist? do
case persist_serial(zone, serial, soa, now) do
{:ok, updated} -> updated
{:error, _} -> zone
end
else
apply_soa_to_zone(zone, serial, soa)
end
# Publish into the authority cache without a full DB reload of records.
inject_zone_into_cache(zone, records, serial, soa)
{:ok,
%{
zone: zone,
serial: serial,
record_count: length(records),
transferred_at: now
}}
end
@impl true
def init(opts) do
poll_ms =
opts
|> Keyword.get(:poll_interval_ms, DNS.secondary_ingest_poll_interval_ms())
|> max(@min_poll_ms)
state = %{poll_interval_ms: poll_ms}
if DNS.secondary_ingest_enabled?() do
send(self(), :ingest)
end
schedule_poll(state.poll_interval_ms)
{:ok, state}
end
@impl true
def handle_call({:refresh, opts}, _from, state) do
result = run_ingest_cycle(opts)
{:reply, result, state}
end
@impl true
def handle_info(:ingest, state) do
if DNS.secondary_ingest_enabled?() do
_ = run_ingest_cycle([])
end
schedule_poll(state.poll_interval_ms)
{:noreply, state}
end
defp run_ingest_cycle(opts) do
zones = list_secondary_zones()
results =
Enum.map(zones, fn zone ->
case transfer_zone(zone, opts) do
{:ok, meta} ->
Logger.info(
"DNS secondary transfer ok domain=#{zone.domain} serial=#{meta.serial} records=#{meta.record_count}"
)
{:ok, zone.domain, meta}
{:error, :serial_unchanged} = err ->
err
{:error, reason} = err ->
Logger.warning(
"DNS secondary transfer failed domain=#{zone.domain} reason=#{inspect(reason)}"
)
_ = record_transfer_error(zone, reason)
err
end
end)
{:ok, results}
rescue
error ->
Logger.warning("DNS secondary ingest cycle failed: #{Exception.message(error)}")
{:error, error}
end
defp do_transfer_zone(zone, masters, opts) do
local_serial = SecondaryStore.serial(zone.domain) || zone.serial || 0
force? = Keyword.get(opts, :force?, false)
masters
|> Enum.reduce_while({:error, :all_masters_failed}, fn master, _acc ->
case try_master(zone, master, local_serial, force?, opts) do
{:ok, _} = ok -> {:halt, ok}
{:error, :serial_unchanged} = skip -> {:halt, skip}
{:error, _reason} = err -> {:cont, err}
end
end)
end
defp try_master(zone, master, local_serial, force?, opts) do
transport = Keyword.get(opts, :transport, &AxfrClient.transfer/3)
case transport.(zone.domain, master, Keyword.take(opts, [:timeout])) do
{:ok, %{serial: serial} = result} ->
if not force? and is_integer(local_serial) and serial == local_serial and
SecondaryStore.lookup(zone.domain) != :error do
{:error, :serial_unchanged}
else
apply_transfer(zone, result, opts)
end
{:error, reason} ->
{:error, reason}
end
end
defp list_secondary_zones do
from(z in Zone, where: z.kind == "secondary")
|> Repo.all()
rescue
# DB may be unavailable during boot; ingest will retry on the next poll.
_ -> []
end
defp persist_serial(zone, serial, soa, now) do
attrs =
%{
"serial" => serial,
"secondary_last_transferred_at" => now,
"last_error" => nil
}
|> maybe_put_soa(soa)
zone
|> Zone.changeset(attrs)
|> Repo.update()
end
defp maybe_put_soa(attrs, nil), do: attrs
defp maybe_put_soa(attrs, soa) when is_map(soa) do
attrs
|> Map.put("soa_mname", Map.get(soa, :mname) || Map.get(soa, "mname"))
|> Map.put("soa_rname", Map.get(soa, :rname) || Map.get(soa, "rname"))
|> Map.put("soa_refresh", Map.get(soa, :refresh) || Map.get(soa, "refresh"))
|> Map.put("soa_retry", Map.get(soa, :retry) || Map.get(soa, "retry"))
|> Map.put("soa_expire", Map.get(soa, :expire) || Map.get(soa, "expire"))
|> Map.put("soa_minimum", Map.get(soa, :minimum) || Map.get(soa, "minimum"))
end
defp apply_soa_to_zone(zone, serial, soa) do
zone = %{zone | serial: serial}
if is_map(soa) do
%{
zone
| soa_mname: Map.get(soa, :mname) || zone.soa_mname,
soa_rname: Map.get(soa, :rname) || zone.soa_rname,
soa_refresh: Map.get(soa, :refresh) || zone.soa_refresh,
soa_retry: Map.get(soa, :retry) || zone.soa_retry,
soa_expire: Map.get(soa, :expire) || zone.soa_expire,
soa_minimum: Map.get(soa, :minimum) || zone.soa_minimum
}
else
zone
end
end
defp inject_zone_into_cache(zone, records, serial, soa) do
zone =
zone
|> apply_soa_to_zone(serial, soa)
|> Map.put(:records, materialize_records(records))
|> Map.put(:status, zone.status || "verified")
# Only inject when the zone is answerable (verified) or secondary ingest
# is allowed to serve pre-verify (we mark verified on successful transfer).
zone = %{zone | status: "verified"}
case :ets.whereis(ZoneCache) do
:undefined ->
:ok
_ ->
:ets.insert(ZoneCache, {String.downcase(zone.domain), zone})
:ok
end
end
defp materialize_records(records) do
Enum.map(records, fn rr ->
# Keep plain maps; Query/Axfr accept maps with host/type/content.
# Attach a name field for zone_records matching.
host =
Map.get(rr, :host) || Map.get(rr, "host") || Map.get(rr, :name) || Map.get(rr, "name")
rr
|> Map.put(:host, host)
|> Map.put_new(:name, relative_name(host, Map.get(rr, :zone_domain)))
end)
end
defp relative_name(nil, _), do: "@"
defp relative_name(host, _zone_domain) do
# Query.record_name handles FQDN host; keep FQDN for fidelity.
host
end
defp record_transfer_error(zone, reason) do
message =
case reason do
atom when is_atom(atom) -> Atom.to_string(atom)
other -> inspect(other)
end
|> String.slice(0, 500)
zone
|> Zone.changeset(%{"last_error" => "secondary transfer: #{message}"})
|> Repo.update()
rescue
_ -> :ok
end
defp schedule_poll(ms) when is_integer(ms) and ms > 0 do
Process.send_after(self(), :ingest, ms)
end
defp schedule_poll(_), do: Process.send_after(self(), :ingest, @default_poll_ms)
end

View file

@ -0,0 +1,123 @@
defmodule Elektrine.DNS.SecondaryStore do
@moduledoc """
In-memory store for zone data transferred into a secondary.
Secondaries **must not** write transferred RRs into the primary
`dns_records` table. This ETS table is the answer source for
`kind=secondary` zones after a successful AXFR.
"""
use GenServer
@table __MODULE__
@type transfer_snapshot :: %{
domain: String.t(),
serial: non_neg_integer(),
records: [map()],
soa: map() | nil,
transferred_at: DateTime.t()
}
def start_link(opts \\ []) do
GenServer.start_link(__MODULE__, opts, name: __MODULE__)
end
@doc """
Look up a transferred zone snapshot by apex domain.
"""
@spec lookup(String.t()) :: {:ok, transfer_snapshot()} | :error
def lookup(domain) when is_binary(domain) do
case :ets.lookup(@table, normalize(domain)) do
[{_key, snapshot}] -> {:ok, snapshot}
[] -> :error
end
rescue
ArgumentError -> :error
end
@doc """
Store a transferred zone snapshot. Replaces any previous snapshot for the domain.
"""
@spec put(transfer_snapshot()) :: :ok
def put(%{domain: domain, serial: serial, records: records} = snapshot)
when is_binary(domain) and is_integer(serial) and is_list(records) do
entry = %{
domain: normalize(domain),
serial: serial,
records: records,
soa: Map.get(snapshot, :soa),
transferred_at: Map.get(snapshot, :transferred_at) || DateTime.utc_now()
}
insert_entry(entry)
end
defp insert_entry(entry) do
true = :ets.insert(@table, {entry.domain, entry})
:ok
rescue
ArgumentError ->
# Table may be missing when GenServer is down (tests / race on shutdown).
ensure_table()
true = :ets.insert(@table, {entry.domain, entry})
:ok
end
@doc """
Remove a transferred snapshot (e.g. zone deleted or expired).
Safe when the ETS table is missing or already torn down (test races / shutdown).
"""
@spec delete(String.t()) :: :ok
def delete(domain) when is_binary(domain) do
:ets.delete(@table, normalize(domain))
:ok
rescue
ArgumentError -> :ok
end
@doc """
List all domains currently held in the secondary store.
"""
@spec list_domains() :: [String.t()]
def list_domains do
:ets.select(@table, [{{:"$1", :_}, [], [:"$1"]}])
rescue
ArgumentError -> []
end
@doc """
Current serial for a domain, or nil when no transfer is present.
"""
@spec serial(String.t()) :: non_neg_integer() | nil
def serial(domain) when is_binary(domain) do
case lookup(domain) do
{:ok, %{serial: serial}} -> serial
:error -> nil
end
end
@impl true
def init(_opts) do
ensure_table()
{:ok, %{}}
end
defp ensure_table do
case :ets.whereis(@table) do
:undefined ->
:ets.new(@table, [:named_table, :public, read_concurrency: true])
_ ->
@table
end
end
defp normalize(domain) do
domain
|> to_string()
|> String.downcase()
|> String.trim_trailing(".")
end
end

View file

@ -6,6 +6,9 @@ defmodule Elektrine.DNS.Zone do
use Ecto.Schema
import Ecto.Changeset
@kinds ~w(native primary secondary)
@max_masters 16
schema "dns_zones" do
field :domain, :string
field :status, :string, default: "provisioning"
@ -32,6 +35,10 @@ field :axfr_enabled, :boolean, default: false
field :axfr_allow_cidrs_text, :string, virtual: true
# Populated by ZoneCache / Dnssec.Signer when DNS_DNSSEC_ENABLED and zone keys exist.
field :dnssec_bundle, :map, virtual: true, default: nil
# Master endpoints for kind=secondary (host or host:port).
field :masters, {:array, :string}, default: []
field :masters_text, :string, virtual: true
field :secondary_last_transferred_at, :utc_datetime
belongs_to :user, Elektrine.Accounts.User
belongs_to :axfr_tsig_key, Elektrine.DNS.TsigKey
@ -46,6 +53,10 @@ has_many :tunnels, Elektrine.DNS.Tunnel, foreign_key: :zone_id
end
@dnssec_statuses ~w(disabled keys_generated signed)
def kinds, do: @kinds
def secondary?(%__MODULE__{kind: "secondary"}), do: true
def secondary?(_), do: false
def changeset(zone, attrs) do
zone
@ -74,16 +85,21 @@ has_many :tunnels, Elektrine.DNS.Tunnel, foreign_key: :zone_id
:axfr_require_tsig,
:axfr_tsig_key_id,
:axfr_allow_cidrs_text,
:masters,
:masters_text,
:secondary_last_transferred_at,
:user_id
])
|> update_change(:domain, &normalize_domain/1)
|> validate_required([:domain, :status, :kind, :default_ttl, :user_id])
|> validate_inclusion(:kind, @kinds)
|> validate_inclusion(:force_https, [true, false])
|> validate_inclusion(:dnssec_enabled, [true, false])
|> validate_inclusion(:dnssec_status, @dnssec_statuses)
|> validate_inclusion(:axfr_enabled, [true, false])
|> validate_inclusion(:axfr_require_tsig, [true, false])
|> validate_number(:default_ttl, greater_than: 0, less_than_or_equal_to: 86_400)
|> validate_number(:serial, greater_than_or_equal_to: 0)
|> validate_number(:soa_refresh, greater_than: 0)
|> validate_number(:soa_retry, greater_than: 0)
|> validate_number(:soa_expire, greater_than: 0)
@ -91,6 +107,8 @@ has_many :tunnels, Elektrine.DNS.Tunnel, foreign_key: :zone_id
|> validate_number(:nameserver_set, greater_than_or_equal_to: 0, less_than_or_equal_to: 4095)
|> validate_format(:domain, ~r/^(?:[a-z0-9-]+\.)+[a-z]{2,}$/)
|> validate_axfr_allow_cidrs()
|> validate_masters()
|> validate_secondary_masters()
|> unique_constraint(:domain, name: :dns_zones_domain_ci_unique)
|> foreign_key_constraint(:user_id)
|> foreign_key_constraint(:axfr_tsig_key_id)
@ -135,6 +153,74 @@ has_many :tunnels, Elektrine.DNS.Tunnel, foreign_key: :zone_id
end
end
defp validate_masters(changeset) do
case get_change(changeset, :masters) do
nil ->
changeset
masters when is_list(masters) ->
if length(masters) > @max_masters do
message = "cannot have more than #{@max_masters} masters"
changeset
|> add_error(:masters, message)
|> add_error(:masters_text, message)
else
case normalize_masters(masters) do
{:ok, normalized} ->
put_change(changeset, :masters, normalized)
{:error, invalid} ->
message = "contains invalid master endpoints: #{Enum.join(invalid, ", ")}"
changeset
|> add_error(:masters, message)
|> add_error(:masters_text, message)
end
end
_ ->
message = "must be a list of host or host:port strings"
changeset
|> add_error(:masters, message)
|> add_error(:masters_text, message)
end
end
defp validate_secondary_masters(changeset) do
kind = get_field(changeset, :kind)
masters = get_field(changeset, :masters) || []
if kind == "secondary" and masters == [] do
message = "secondary zones require at least one master"
changeset
|> add_error(:masters, message)
|> add_error(:masters_text, message)
else
changeset
end
end
defp normalize_masters(masters) do
{ok, bad} =
Enum.reduce(masters, {[], []}, fn entry, {good, bad} ->
trimmed = entry |> to_string() |> String.trim()
case Elektrine.DNS.AxfrClient.parse_master(trimmed) do
{:ok, _} -> {[trimmed | good], bad}
{:error, _} -> {good, [trimmed | bad]}
end
end)
if bad == [] do
{:ok, Enum.reverse(ok) |> Enum.uniq()}
else
{:error, Enum.reverse(bad)}
end
end
def nameserver_records(%__MODULE__{domain: domain} = zone) when is_binary(domain) do
zone
|> Elektrine.DNS.assigned_nameservers()

View file

@ -14,6 +14,7 @@ defmodule Elektrine.DNS.ZoneCache do
alias Elektrine.DNS
alias Elektrine.DNS.Dnssec.Signer
alias Elektrine.DNS.SecondaryStore
alias Elektrine.DNS.Zone
alias Elektrine.Repo
@ -114,6 +115,13 @@ defmodule Elektrine.DNS.ZoneCache do
|> Enum.map(&maybe_prepare_dnssec(&1, sign_dnssec?))
{:ok, zones}
{:ok,
Zone
|> where([z], z.status == "verified" or z.kind == "secondary")
|> preload(:records)
|> Repo.all(repo_opts)
|> Enum.map(&prepare_zone_for_cache/1)
|> Enum.reject(&is_nil/1)}
rescue
error in [Postgrex.Error, DBConnection.OwnershipError] ->
{:error, error}
@ -138,6 +146,41 @@ defmodule Elektrine.DNS.ZoneCache do
end
defp maybe_prepare_dnssec(zone, _), do: zone
# Secondary zones answer from transferred data (SecondaryStore), never from
# primary dns_records rows. Native/primary zones keep the existing DB path.
defp prepare_zone_for_cache(%Zone{kind: "secondary"} = zone) do
case SecondaryStore.lookup(zone.domain) do
{:ok, snapshot} ->
zone
|> Map.put(:serial, snapshot.serial)
|> Map.put(:records, snapshot.records || [])
|> apply_snapshot_soa(snapshot.soa)
|> Map.put(:status, "verified")
:error ->
# No transfer yet — keep the zone out of the authority cache so we
# do not answer from empty/stale primary records.
nil
end
end
defp prepare_zone_for_cache(%Zone{} = zone) do
repair_builtin_zone_records(zone)
end
defp apply_snapshot_soa(zone, nil), do: zone
defp apply_snapshot_soa(zone, soa) when is_map(soa) do
%{
zone
| soa_mname: Map.get(soa, :mname) || zone.soa_mname,
soa_rname: Map.get(soa, :rname) || zone.soa_rname,
soa_refresh: Map.get(soa, :refresh) || zone.soa_refresh,
soa_retry: Map.get(soa, :retry) || zone.soa_retry,
soa_expire: Map.get(soa, :expire) || zone.soa_expire,
soa_minimum: Map.get(soa, :minimum) || zone.soa_minimum
}
end
defp repair_builtin_zone_records(%Zone{} = zone) do
case Elektrine.DNS.repair_builtin_user_zone_records(zone) do

View file

@ -742,6 +742,8 @@ status = DNS.zone_dnssec_status(zone)
assigned_nameservers: assigned_nameservers,
dnssec_enabled: Map.get(zone, :dnssec_enabled, false) == true,
dnssec_status: Map.get(zone, :dnssec_status) || "disabled",
masters: Map.get(zone, :masters) || [],
secondary_last_transferred_at: Map.get(zone, :secondary_last_transferred_at),
verified_at: zone.verified_at,
last_checked_at: zone.last_checked_at,
last_published_at: zone.last_published_at,

View file

@ -0,0 +1,362 @@
defmodule Elektrine.DNS.SecondaryIngestTest do
use Elektrine.DataCase, async: false
alias Elektrine.AccountsFixtures
alias Elektrine.DNS
alias Elektrine.DNS.Axfr
alias Elektrine.DNS.AxfrClient
alias Elektrine.DNS.Packet
alias Elektrine.DNS.Query
alias Elektrine.DNS.Record
alias Elektrine.DNS.SecondaryIngest
alias Elektrine.DNS.SecondaryStore
alias Elektrine.DNS.Zone
alias Elektrine.Repo
@domain "secondary-ingest.test"
@serial 2_026_080_202
setup do
previous_dns = Application.get_env(:elektrine, :dns, [])
Application.put_env(
:elektrine,
:dns,
Keyword.merge(previous_dns,
secondary_ingest_enabled: true,
secondary_axfr_enabled: true,
nameservers: ["ns1.elektrine.com", "ns2.elektrine.com"],
soa_rname: "admin.elektrine.com"
)
)
on_exit(fn ->
Application.put_env(:elektrine, :dns, previous_dns)
SecondaryStore.delete(@domain)
end)
ensure_secondary_store()
ensure_zone_cache()
SecondaryStore.delete(@domain)
user = AccountsFixtures.user_fixture()
{:ok, zone} =
DNS.create_zone(user, %{
"domain" => @domain,
"kind" => "secondary",
"masters" => ["203.0.113.50:5300"]
})
%{user: user, zone: zone}
end
test "zone kind validates native/primary/secondary and secondary requires masters", %{
user: user
} do
assert {:error, changeset} =
DNS.create_zone(user, %{
"domain" => "no-masters.example.com",
"kind" => "secondary",
"masters" => []
})
assert %{masters: _} = errors_on(changeset)
assert {:ok, primary} =
DNS.create_zone(user, %{
"domain" => "primary-kind.example.com",
"kind" => "primary"
})
assert primary.kind == "primary"
assert Zone.secondary?(primary) == false
end
test "secondary zones reject direct record writes", %{zone: zone} do
assert {:error, changeset} =
DNS.create_record(zone, %{
"name" => "@",
"type" => "A",
"content" => "203.0.113.1",
"ttl" => 300
})
assert %{zone_id: _} = errors_on(changeset)
end
test "parse_axfr_messages recovers SOA serial and RRs from primary transfer fixture" do
primary = primary_fixture_zone()
query = %{id: 42, rd: 0, qname: primary.domain, qtype: :axfr, edns: false, udp_size: 512}
messages = Packet.encode_axfr_messages(query, Axfr.zone_transfer_records(primary))
assert {:ok, result} = AxfrClient.parse_axfr_messages(messages, primary.domain)
assert result.serial == primary.serial
assert result.soa.serial == primary.serial
assert result.soa.mname == "ns1.elektrine.com"
types = Enum.map(result.records, & &1.type)
assert "A" in types
assert "MX" in types
assert "TXT" in types
assert "NS" in types
end
test "apply_transfer stores data in SecondaryStore and updates serial without dns_records", %{
zone: zone
} do
primary = primary_fixture_zone()
query = %{id: 1, rd: 0, qname: primary.domain, qtype: :axfr, edns: false, udp_size: 512}
messages = Packet.encode_axfr_messages(query, Axfr.zone_transfer_records(primary))
assert {:ok, parsed} = AxfrClient.parse_axfr_messages(messages, primary.domain)
# Remap transfer to secondary domain for the answer path.
parsed = remap_transfer(parsed, zone.domain)
assert {:ok, meta} = SecondaryIngest.apply_transfer(zone, parsed, persist?: true)
assert meta.serial == @serial
assert meta.record_count > 0
assert {:ok, snapshot} = SecondaryStore.lookup(zone.domain)
assert snapshot.serial == @serial
assert length(snapshot.records) == meta.record_count
reloaded = Repo.get!(Zone, zone.id)
assert reloaded.serial == @serial
assert reloaded.secondary_last_transferred_at != nil
assert reloaded.soa_mname == "ns1.elektrine.com"
# Critical: transferred RRs must not land in primary dns_records.
assert Repo.aggregate(
from(r in Record, where: r.zone_id == ^zone.id),
:count
) == 0
end
test "authority answers from transferred secondary data", %{zone: zone} do
primary = primary_fixture_zone()
query = %{id: 2, rd: 0, qname: primary.domain, qtype: :axfr, edns: false, udp_size: 512}
messages = Packet.encode_axfr_messages(query, Axfr.zone_transfer_records(primary))
assert {:ok, parsed} = AxfrClient.parse_axfr_messages(messages, primary.domain)
parsed = remap_transfer(parsed, zone.domain)
assert {:ok, _} = SecondaryIngest.apply_transfer(zone, parsed, persist?: false)
# ZoneCache should have the secondary zone with transferred records.
assert {:ok, cached} = Elektrine.DNS.ZoneCache.lookup(zone.domain)
assert cached.kind == "secondary"
assert cached.serial == @serial
assert cached.records != []
a_packet = build_query(zone.domain, 1)
result = Query.resolve(a_packet, transport: :udp)
assert result.rcode == :noerror
assert result.authoritative == true
assert result.zone.domain == zone.domain
answers = decode_answers(result.response)
assert Enum.any?(answers, fn rr -> rr.type == :a end)
soa_packet = build_query(zone.domain, 6)
soa_result = Query.resolve(soa_packet, transport: :udp)
assert soa_result.rcode == :noerror
soa_answers = decode_answers(soa_result.response)
assert Enum.any?(soa_answers, fn rr -> rr.type == :soa end)
end
test "serial_unchanged when store already has the same serial", %{zone: zone} do
primary = primary_fixture_zone()
query = %{id: 3, rd: 0, qname: primary.domain, qtype: :axfr, edns: false, udp_size: 512}
messages = Packet.encode_axfr_messages(query, Axfr.zone_transfer_records(primary))
assert {:ok, parsed} = AxfrClient.parse_axfr_messages(messages, primary.domain)
parsed = remap_transfer(parsed, zone.domain)
assert {:ok, _} = SecondaryIngest.apply_transfer(zone, parsed, persist?: false)
transport = fn _domain, _master, _opts -> {:ok, parsed} end
assert {:error, :serial_unchanged} =
SecondaryIngest.transfer_zone(zone,
transport: transport,
persist?: false,
force?: false
)
assert {:ok, meta} =
SecondaryIngest.transfer_zone(zone,
transport: transport,
persist?: false,
force?: true
)
assert meta.serial == @serial
end
test "AxfrClient.parse_master accepts host and host:port forms" do
assert {:ok, {"203.0.113.10", 53}} = AxfrClient.parse_master("203.0.113.10")
assert {:ok, {"203.0.113.10", 5300}} = AxfrClient.parse_master("203.0.113.10:5300")
assert {:ok, {"ns1.example.com", 53}} = AxfrClient.parse_master("ns1.example.com")
assert {:error, :empty_master} = AxfrClient.parse_master(" ")
end
test "update_zone cannot forge serial or secondary_last_transferred_at", %{zone: zone} do
original_transferred_at = zone.secondary_last_transferred_at
forged_serial = 9_999_999
assert {:ok, updated} =
DNS.update_zone(zone, %{
"serial" => forged_serial,
"secondary_last_transferred_at" => DateTime.utc_now()
})
reloaded = Repo.get!(Zone, updated.id)
# Public attrs ignore caller-supplied serial; authority touch may bump by 1.
refute reloaded.serial == forged_serial
assert reloaded.serial < forged_serial
assert reloaded.secondary_last_transferred_at == original_transferred_at
end
test "delete_zone clears SecondaryStore snapshot", %{zone: zone} do
SecondaryStore.put(%{
domain: zone.domain,
serial: @serial,
records: [%{host: zone.domain, type: "A", content: "203.0.113.9", ttl: 300}],
soa: nil,
transferred_at: DateTime.utc_now()
})
assert {:ok, _} = SecondaryStore.lookup(zone.domain)
assert {:ok, _} = DNS.delete_zone(zone)
assert :error = SecondaryStore.lookup(zone.domain)
end
test "kind change away from secondary clears SecondaryStore", %{zone: zone} do
SecondaryStore.put(%{
domain: zone.domain,
serial: @serial,
records: [%{host: zone.domain, type: "A", content: "203.0.113.9", ttl: 300}],
soa: nil,
transferred_at: DateTime.utc_now()
})
# native no longer needs masters; clear them so validation passes.
assert {:ok, _} = DNS.update_zone(zone, %{"kind" => "native", "masters" => []})
assert :error = SecondaryStore.lookup(zone.domain)
end
test "SecondaryStore.delete is safe when the ETS table is missing" do
assert :ok = SecondaryStore.delete("missing-table-safe.example")
end
# --- fixtures -------------------------------------------------------------
defp primary_fixture_zone do
%Zone{
domain: "axfr-primary.test",
status: "verified",
kind: "primary",
serial: @serial,
default_ttl: 300,
soa_mname: "ns1.elektrine.com",
soa_rname: "admin.elektrine.com",
soa_refresh: 3600,
soa_retry: 600,
soa_expire: 1_209_600,
soa_minimum: 300,
axfr_enabled: true,
axfr_allow_cidrs: ["203.0.113.0/24"],
records: [
%Record{name: "@", type: "A", content: "203.0.113.20", ttl: 300},
%Record{name: "www", type: "A", content: "203.0.113.21", ttl: 300},
%Record{name: "@", type: "MX", content: "mail.axfr-primary.test", ttl: 300, priority: 10},
%Record{name: "mail", type: "A", content: "203.0.113.22", ttl: 300},
%Record{name: "@", type: "TXT", content: "v=spf1 -all", ttl: 300}
]
}
end
defp remap_transfer(parsed, domain) do
records =
Enum.map(parsed.records, fn rr ->
host = Map.get(rr, :host) || ""
new_host =
host
|> String.replace_suffix("axfr-primary.test", domain)
|> then(fn h -> if h == "", do: domain, else: h end)
content =
case Map.get(rr, :content) || Map.get(rr, :value) do
nil -> nil
c -> String.replace(to_string(c), "axfr-primary.test", domain)
end
rr
|> Map.put(:host, new_host)
|> then(fn r ->
if content do
if Map.has_key?(r, :value),
do: Map.put(r, :value, content),
else: Map.put(r, :content, content)
else
r
end
end)
end)
soa =
parsed.soa
|> Map.put(:host, domain)
|> Map.put(:mname, "ns1.elektrine.com")
%{parsed | records: records, soa: soa}
end
defp ensure_secondary_store do
case :ets.whereis(SecondaryStore) do
:undefined ->
{:ok, _pid} = SecondaryStore.start_link([])
_ ->
:ok
end
rescue
ArgumentError ->
# Already started under a different name registration path.
:ok
end
defp ensure_zone_cache do
case :ets.whereis(Elektrine.DNS.ZoneCache) do
:undefined ->
:ets.new(Elektrine.DNS.ZoneCache, [:named_table, :public, read_concurrency: true])
_ ->
:ok
end
end
defp build_query(name, type) do
<<0xAB, 0xCD, 0x00, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00,
encode_name(name)::binary, type::16, 1::16>>
end
defp encode_name(name) do
name
|> String.split(".", trim: true)
|> Enum.map_join(fn label -> <<byte_size(label)>> <> label end)
|> Kernel.<>(<<0>>)
end
defp decode_answers(message) do
case :inet_dns.decode(message) do
{:ok, {:dns_rec, _header, _qd, answers, _ns, _ar}} ->
Enum.map(answers, fn rr when is_tuple(rr) and tuple_size(rr) >= 7 ->
%{domain: elem(rr, 1), type: elem(rr, 2), data: elem(rr, 6)}
end)
_ ->
[]
end
end
end

View file

@ -309,6 +309,9 @@ config :elektrine, :dns,
# Outbound AXFR primary (TCP zone transfer). Off by default; also requires
# per-zone axfr_enabled + allow-list CIDRs.
secondary_axfr_enabled: false,
# Inbound secondary AXFR ingest (kind=secondary zones). Off by default.
secondary_ingest_enabled: false,
secondary_ingest_poll_interval_ms: 60_000,
zone_cache_refresh_interval_ms: 300_000,
udp_port: 5300,
tcp_port: 5300,

View file

@ -196,6 +196,16 @@ secondary_axfr_enabled:
"DNS_SECONDARY_AXFR_ENABLED",
Keyword.get(dns_config, :secondary_axfr_enabled, false)
),
secondary_ingest_enabled:
parse_bool_env.(
"DNS_SECONDARY_INGEST_ENABLED",
Keyword.get(dns_config, :secondary_ingest_enabled, false)
),
secondary_ingest_poll_interval_ms:
parse_int_env.(
"DNS_SECONDARY_INGEST_POLL_INTERVAL_MS",
Keyword.get(dns_config, :secondary_ingest_poll_interval_ms, 60_000)
),
edge_proxy_enabled:
parse_bool_env.(
"DNS_EDGE_PROXY_ENABLED",

View file

@ -269,3 +269,35 @@ prefer a single sticky connector host).
- Multi-node edge site registry + heartbeat is implemented (`dns_edge_sites`,
admin `/pripyat/edge-sites`, `POST /_edge/site/v1/heartbeat`; see
edge-platform)
### Multi-NS and secondary DNS
Two distinct HA models:
1. **Multi-active authority (same Postgres)** — N× containers with
`ELEKTRINE_RUNTIME_ROLE=worker`, `DNS_AUTHORITY_ENABLED=true`, sharing one
database. Zone serial is the shared source of truth; every node answers
from `ZoneCache` reloaded from Postgres. There is no AXFR between them.
2. **RFC secondary (AXFR ingest)** — a true secondary does **not** write the
primary's `dns_records` table. Instead:
- Zone `kind` is `native` (default multi-active), `primary` (explicit
AXFR source), or `secondary` (transfer consumer).
- Primary outbound AXFR is gated by `DNS_SECONDARY_AXFR_ENABLED` plus
per-zone `axfr_enabled` and allow-list CIDRs (`Elektrine.DNS.Axfr`).
- Secondary ingest is gated by `DNS_SECONDARY_INGEST_ENABLED`. The
`SecondaryIngest` worker AXFRs each `kind=secondary` zone from its
`masters` list into `SecondaryStore` (ETS). Authority answers from
that snapshot. Only zone metadata (`serial`, SOA timers,
`secondary_last_transferred_at`) is written back to Postgres.
- TSIG MAC verification on the wire is not enabled yet (unsigned first);
keys may be stored for later use.
### Not implemented yet
- DNSSEC signing and DS rollover automation
- Signed AXFR / TSIG MAC verification
- IXFR
DNSSEC remains deferred pending operator UX and failure-mode design beyond
keys/export phases.

7
env/presets/dns.env vendored
View file

@ -33,6 +33,13 @@
# When enabled, each zone still needs axfr_enabled + allow-list CIDRs.
# DNS_SECONDARY_AXFR_ENABLED=false
# Inbound secondary AXFR ingest. Off by default. When enabled on an authority
# node, zones with kind=secondary and masters set are transferred into an
# in-memory SecondaryStore; the authority answers from that data without
# writing dns_records (primary RR table). Serial/SOA metadata is updated only.
# DNS_SECONDARY_INGEST_ENABLED=false
# DNS_SECONDARY_INGEST_POLL_INTERVAL_MS=60000
# If this host also runs NetBird or another private DNS listener, bind public
# authoritative DNS to the public interface instead of 0.0.0.0.
# PUBLIC_DNS_BIND_IP=203.0.113.10