Skip to content
Merged
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
18 changes: 18 additions & 0 deletions config/runtime.exs
Original file line number Diff line number Diff line change
Expand Up @@ -815,6 +815,24 @@ if channels = System.get_env("GEN_RPC_SCATTER_CHANNELS") do
channels: Smolquery.RuntimeConfig.positive_integer!("GEN_RPC_SCATTER_CHANNELS", channels)
end

# T-614: how long a socket to a peer node may outlive the peer.
if ms = System.get_env("SMOLQUERY_PEER_TCP_USER_TIMEOUT_MS") do
config :smolquery, Smolquery.PeerSocket,
user_timeout_ms:
Smolquery.RuntimeConfig.positive_integer!("SMOLQUERY_PEER_TCP_USER_TIMEOUT_MS", ms)
end

if seconds = System.get_env("SMOLQUERY_PEER_TCP_KEEPALIVE_IDLE_S") do
config :smolquery, Smolquery.PeerSocket,
keepalive_idle_s:
Smolquery.RuntimeConfig.integer_in_range!(
"SMOLQUERY_PEER_TCP_KEEPALIVE_IDLE_S",
seconds,
1,
32_767
)
end

if catalog_database_url = System.get_env("CATALOG_DATABASE_URL") do
db = Smolquery.DatabaseUrl.parse!(catalog_database_url)

Expand Down
2 changes: 2 additions & 0 deletions docs/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -169,6 +169,8 @@ version's driver exists.
| `GEN_RPC_SCATTER_CHANNELS` | How many `{:scatter, _}` gen_rpc connections a node opens to each peer for distributed-query partials (`4`). A job hashes to one connection; concurrent jobs spread over the pool |
| `GEN_RPC_TLS` | `true` switches buffer/query inter-node traffic to mutual TLS (Transport Layer Security) (`false`). Verification is chain-only against the cluster CA (certificate authority); the emqx gen_rpc fork does no hostname or CN (common name) check. The CA is thus the trust boundary. Certificate files are per node (`GEN_RPC_TLS_DIR`, default `/etc/smolquery/gen-rpc-tls`; `POD_NAME` names the file). Any CA-signed certificate authenticates to any peer. A leaked node certificate thus means a rotation of the CA, not just of that node |
| `GEN_RPC_SSL_PORT` | The gen_rpc TLS port (`5870`) |
| `SMOLQUERY_PEER_TCP_USER_TIMEOUT_MS` | How long a socket to a peer node may go without an acknowledgement before the kernel closes it (`15000`): sent data in flight, or keepalive probes on an idle socket. Applies to `HotClient`'s manifest reads (T-614). Linux only |
| `SMOLQUERY_PEER_TCP_KEEPALIVE_IDLE_S` | How long a socket to a peer node may sit idle before the kernel starts keepalive probes (`5`, at most `32767`). Probes repeat every 5 s; `SMOLQUERY_PEER_TCP_USER_TIMEOUT_MS`, not a probe count, decides when an unanswered idle socket closes. Linux only for the timers; elsewhere keepalive uses the OS defaults |
| `DIST_TLS` | `true` runs Erlang distribution (cluster membership only) over TLS with the same certificates (`false`). Set it in `rel/env.sh.eex`, not `config/runtime.exs`. Distribution starts before the release's Elixir config does |
| `POD_NAME` / `POD_NAMESPACE` | When set (a Kubernetes Downward API convention), `rel/env.sh.eex` derives `RELEASE_NODE` from the pod's stable headless-service DNS name. That is the same name a peer needs to reach this node |
| `HEADLESS_SERVICE` | The headless-service name in that derived node name (`smolquery-headless`) |
Expand Down
32 changes: 32 additions & 0 deletions docs/deployment.md
Original file line number Diff line number Diff line change
Expand Up @@ -97,10 +97,42 @@ and `SMOLQUERY_DISTRIBUTED_WORKER_THREADS`. A scattered query's declared
budget on one node is the worker count `×` the worker limit, on top of the
job engine's own `job_memory_limit`.

## When a peer dies without closing its sockets

A pod that is killed hard closes none of its sockets. Its replacement has the
same name and a new IP. Linux keeps a socket to the old IP until it gives up
retransmitting: `tcp_retries2` = 15, about 15.4 minutes. On 2026-10-02 four
hard-killed pods were replaced in 13 s, and the cluster answered 5xx for 16
minutes.

Each connection between nodes now has its own bound:

| connection | what bounds it | setting |
|---|---|---|
| gen_rpc, buffer writes and replication (`:control`, `{:bulk, _}`) and scatter partials (`{:scatter, _}`) | Every client to a node is killed when distribution reports it down (T-613). A call that times out probes its channel and kills the client if the probe is unanswered (T-613). | none for the drop; the probe waits 10 s |
| gen_rpc sockets this node accepts | gen_rpc's own kernel keepalive: idle 5 s, interval 5 s, 2 probes | `config :gen_rpc, socket_keepalive_*` |
| gen_rpc sockets this node dials | `SO_KEEPALIVE` is on, with Linux's default timers: first probe after 2 h. gen_rpc 3.6.1 does not let an application set the timers or `TCP_USER_TIMEOUT` on a client socket. The first row covers it. The pod sysctls `net.ipv4.tcp_keepalive_time`, `_intvl` and `_probes` (Kubernetes safe sysctls) would shorten those timers for every socket in the pod. | none in smolquery |
| `HotClient` HTTP: the storage sealer and the query planner reading a buffer node's manifest | `TCP_USER_TIMEOUT` closes a socket to a dead peer after 15 s, and the pool drops it (T-614). Keepalive probes an idle socket every 5 s after 5 s idle, so an idle socket is checked against the same 15 s; on Linux the user timeout, not the probe count, decides when keepalive gives up. | `SMOLQUERY_PEER_TCP_USER_TIMEOUT_MS`, `SMOLQUERY_PEER_TCP_KEEPALIVE_IDLE_S`, `config :smolquery, Smolquery.PeerSocket` |
| Erlang distribution | The net tick: a silent node is down after `net_ticktime`, 60 s by default. On 2026-10-02 it reconnected in seconds. | `-kernel net_ticktime` |
| DuckDB `httpfs` reading hot segments from a buffer node | DuckDB's own `http_timeout`. Not bounded by smolquery. | none |

The keepalive timers and `TCP_USER_TIMEOUT` are Linux socket options. On
another OS only keepalive is switched on, with that OS's timers.

Neither bound cuts a slow peer short. A slow peer's kernel still
acknowledges data and keepalive probes; the bounds act only on a peer that
answers nothing at the TCP level.

## Upgrade notes

One note per release, newest first.

### 0.22.0: sockets to a replaced pod close in seconds (T-613, T-614)

- **gen_rpc:** a node's clients to a peer are killed when distribution reports the peer down, and a channel that times out is probed and redialed if it is dead (T-613). Nothing to configure.
- **Buffer manifest reads (`HotClient`)** close a socket to a dead peer after 15 s: a 15 s `TCP_USER_TIMEOUT`, with keepalive probing an idle socket every 5 s after 5 s idle. Before, a pooled connection to a replaced buffer pod stayed in the pool and failed one seal per connection, each after the full 30 s receive timeout, about once a minute. Tune with `SMOLQUERY_PEER_TCP_USER_TIMEOUT_MS` and `SMOLQUERY_PEER_TCP_KEEPALIVE_IDLE_S`.
- See [When a peer dies without closing its sockets](#when-a-peer-dies-without-closing-its-sockets) for every connection's bound.

### 0.21.0: federated DuckLake connections (T-610)

You can register another DuckLake as a connection and query it: `"kind": "ducklake"` on `POST /v1/connections`, or the kind selector on the connections page. It takes the lake's Postgres metadata database, its data path and optional S3 credentials, sealed like a password. The lake attaches read-only at its data path, and a query reads it as `name.schema.table`, alone or joined with smolquery tables. See [api.md](api.md).
Expand Down
20 changes: 19 additions & 1 deletion lib/smolquery/buffer_service/hot_client.ex
Original file line number Diff line number Diff line change
Expand Up @@ -14,8 +14,19 @@ defmodule Smolquery.BufferService.HotClient do
`HotServer` listens is the caller's configuration (`buffer_base_url` on the
storage and query runtimes) — honest for a single node; a cluster resolves
it from the ownership ring instead, which arrives with Milestone 8.

## Connections to a replaced pod

Requests go through a pooled connection with `Smolquery.PeerSocket`'s
keepalive timers and `TCP_USER_TIMEOUT` (T-614). A buffer pod killed
without closing its sockets leaves pooled connections to its old IP.
Without those bounds each one stayed in the pool until a request drew it
and waited out the whole receive timeout, one failed seal per stale
connection. With them, the kernel closes a dead socket after 15 s, idle or
in flight, and the pool drops it.
"""

alias Smolquery.PeerSocket
alias Smolquery.Segments.Store

@default_timeout_ms 30_000
Expand Down Expand Up @@ -128,7 +139,14 @@ defmodule Smolquery.BufferService.HotClient do

defp fetch(url, scope, timeout_ms) do
headers = [{Smolquery.InternalSecret.header(), Smolquery.InternalSecret.value()}]
common = [url: url, headers: headers, receive_timeout: timeout_ms, retry: false]

common = [
url: url,
headers: headers,
receive_timeout: timeout_ms,
retry: false,
connect_options: [transport_opts: PeerSocket.tcp_options()]
]

request =
if scope, do: [method: :post, json: scope] ++ common, else: [method: :get] ++ common
Expand Down
81 changes: 81 additions & 0 deletions lib/smolquery/peer_socket.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,81 @@
defmodule Smolquery.PeerSocket do
@moduledoc """
TCP options that bound how long a socket to a peer node can outlive the peer
(T-614).

A pod killed without closing its sockets leaves every connection to it
pointed at an IP that no longer answers. Linux keeps such a socket until it
gives up retransmitting, `tcp_retries2` = 15, about 15.4 minutes. Finch
turns `SO_KEEPALIVE` on for every pooled socket, but with Linux's default
timers an idle socket is first probed after two hours. A pooled HTTP
connection to a dead peer therefore stays in its pool, and each request
that draws it waits out its whole receive timeout before the pool lets it
go. On the sandbox on 2026-10-02 that was a failed seal about once a minute
after a buffer pod was replaced, one per stale connection.

`tcp_options/0` sets one bound, `user_timeout_ms` (15,000), through two
socket options:

* **`TCP_USER_TIMEOUT`** closes a socket whose sent data has gone
unacknowledged that long. A request already in flight on a dead
socket fails then, not after its receive timeout.
* **Keepalive** probes an idle socket after `keepalive_idle_s` (5) and
every `keepalive_interval_s` (5) after that. On Linux the user timeout
also governs keepalive: the kernel closes an idle socket at the first
unanswered probe at or after `user_timeout_ms`, and the probe count
plays no part. At the defaults that is 15 s. A closed idle socket
leaves its pool before a request can draw it.

So `user_timeout_ms` is the bound for both cases, and the keepalive timers
only set how often an idle socket is checked against it.

Both act only on a peer that has stopped answering at the TCP level. A slow
peer still ACKs data and probes, so neither cuts a slow response short.

The timers are Linux socket options set through `{:raw, ...}`. Elsewhere only
`keepalive: true` is set, with the operating system's own timers.

config :smolquery, Smolquery.PeerSocket,
user_timeout_ms: 15_000,
keepalive_idle_s: 5,
keepalive_interval_s: 5
"""

@ipproto_tcp 6
@tcp_keepidle 4
@tcp_keepintvl 5
@tcp_user_timeout 18

@defaults [user_timeout_ms: 15_000, keepalive_idle_s: 5, keepalive_interval_s: 5]

@doc """
The `:gen_tcp` options for a socket to a peer node, from this node's
configuration.
"""
@spec tcp_options() :: keyword()
def tcp_options, do: tcp_options(config(), :os.type())

@doc """
The `:gen_tcp` options for `config` on the operating system `os_type`, as
`:os.type/0` names it.
"""
@spec tcp_options(keyword(), {atom(), atom()}) :: keyword()
def tcp_options(config, {:unix, :linux}) do
[
keepalive: true,
raw: {@ipproto_tcp, @tcp_keepidle, int(config[:keepalive_idle_s])},
raw: {@ipproto_tcp, @tcp_keepintvl, int(config[:keepalive_interval_s])},
raw: {@ipproto_tcp, @tcp_user_timeout, int(config[:user_timeout_ms])}
]
end

def tcp_options(_config, _os_type), do: [keepalive: true]

@doc """
This node's settings, each defaulted.
"""
@spec config() :: keyword()
def config, do: Keyword.merge(@defaults, Application.get_env(:smolquery, __MODULE__, []))

defp int(value), do: <<value::32-native>>
end
33 changes: 33 additions & 0 deletions test/smolquery/buffer_service/hot_client_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ defmodule Smolquery.BufferService.HotClientTest do
@moduletag :tmp_dir

@table {"analytics", "events"}
@linux match?({:unix, :linux}, :os.type())

defp batch(range) do
%{schema: Schema.new!([{"id", :int64}]), rows: for(i <- range, do: %{"id" => i})}
Expand Down Expand Up @@ -163,4 +164,36 @@ defmodule Smolquery.BufferService.HotClientTest do
{:error, {:invalid_identifier, "../etc"}}
end
end

describe "connections (T-614)" do
@tag skip: not @linux and "Linux socket options"
test "the pooled socket to a buffer node carries PeerSocket's bounds", context do
buffer = start_buffer(context)
base_url = HotServer.base_url(buffer)
port = base_url |> URI.parse() |> Map.fetch!(:port)

assert {:ok, []} = HotClient.manifest(base_url, @table)

[socket | _rest] = sockets_to(port)

{:ok, options} =
:inet.getopts(socket, [:keepalive] ++ for(o <- [4, 5, 18], do: {:raw, 6, o, 4}))

assert {:keepalive, true} in options

assert for(
{:raw, 6, option, <<value::32-native>>} <- options,
into: %{},
do: {option, value}
) ==
%{4 => 5, 5 => 5, 18 => 15_000}
end
end

defp sockets_to(port) do
for socket <- Port.list(),
Port.info(socket, :name) == {:name, ~c"tcp_inet"},
match?({:ok, {_address, ^port}}, :inet.peername(socket)),
do: socket
end
end
55 changes: 55 additions & 0 deletions test/smolquery/peer_socket_test.exs
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
defmodule Smolquery.PeerSocketTest do
use ExUnit.Case, async: false

alias Smolquery.PeerSocket

@linux match?({:unix, :linux}, :os.type())

defp read_back(options) do
for {:raw, 6, option, <<value::32-native>>} <- options, into: %{}, do: {option, value}
end

defp decoded(options) do
for {:raw, {6, option, <<value::32-native>>}} <- options, into: %{}, do: {option, value}
end

test "on Linux sets keepalive timers and TCP_USER_TIMEOUT" do
options = PeerSocket.tcp_options(PeerSocket.config(), {:unix, :linux})

assert options[:keepalive] == true
assert decoded(options) == %{4 => 5, 5 => 5, 18 => 15_000}
end

test "elsewhere sets keepalive only" do
assert PeerSocket.tcp_options(PeerSocket.config(), {:unix, :darwin}) == [keepalive: true]
end

test "config overrides each default" do
previous = Application.fetch_env(:smolquery, PeerSocket)
Application.put_env(:smolquery, PeerSocket, user_timeout_ms: 7_000, keepalive_idle_s: 3)

on_exit(fn ->
case previous do
{:ok, value} -> Application.put_env(:smolquery, PeerSocket, value)
:error -> Application.delete_env(:smolquery, PeerSocket)
end
end)

assert PeerSocket.config()[:keepalive_interval_s] == 5
assert decoded(PeerSocket.tcp_options(PeerSocket.config(), {:unix, :linux}))[18] == 7_000
assert decoded(PeerSocket.tcp_options(PeerSocket.config(), {:unix, :linux}))[4] == 3
end

@tag skip: not @linux and "Linux socket options"
test "a socket opened with tcp_options/0 carries them" do
{:ok, listener} = :gen_tcp.listen(0, [:binary, active: false])
{:ok, port} = :inet.port(listener)
{:ok, socket} = :gen_tcp.connect(~c"127.0.0.1", port, [:binary] ++ PeerSocket.tcp_options())

{:ok, options} =
:inet.getopts(socket, [:keepalive] ++ for(o <- [4, 5, 18], do: {:raw, 6, o, 4}))

assert {:keepalive, true} in options
assert read_back(options) == %{4 => 5, 5 => 5, 18 => 15_000}
end
end
Loading