From 725d0a99b58b47f239237c520b784e0560b1fcf4 Mon Sep 17 00:00:00 2001 From: David Whittington Date: Fri, 14 Aug 2026 00:49:21 +0000 Subject: [PATCH 1/8] feat(auth): define provider-neutral authorization contracts Add validated principal and single-tenant context types, a closed capability policy with distinct unauthenticated and forbidden results, and shared Plug/LiveView attachment seams. OIDC identities derive a versioned opaque id from the exact issuer and subject without retaining raw claims or token material. Closes T-229. Refs PL-27. --- lib/smolquery/auth.ex | 62 ++++++++ lib/smolquery/auth/context.ex | 105 +++++++++++++ lib/smolquery/auth/policy.ex | 28 ++++ lib/smolquery/auth/principal.ex | 196 +++++++++++++++++++++++++ test/smolquery/auth/context_test.exs | 71 +++++++++ test/smolquery/auth/policy_test.exs | 38 +++++ test/smolquery/auth/principal_test.exs | 112 ++++++++++++++ test/smolquery/auth_test.exs | 55 +++++++ 8 files changed, 667 insertions(+) create mode 100644 lib/smolquery/auth.ex create mode 100644 lib/smolquery/auth/context.ex create mode 100644 lib/smolquery/auth/policy.ex create mode 100644 lib/smolquery/auth/principal.ex create mode 100644 test/smolquery/auth/context_test.exs create mode 100644 test/smolquery/auth/policy_test.exs create mode 100644 test/smolquery/auth/principal_test.exs create mode 100644 test/smolquery/auth_test.exs diff --git a/lib/smolquery/auth.ex b/lib/smolquery/auth.ex new file mode 100644 index 00000000..7405d7d7 --- /dev/null +++ b/lib/smolquery/auth.ex @@ -0,0 +1,62 @@ +defmodule Smolquery.Auth do + @moduledoc """ + Attachment seam for an authenticated `Smolquery.Auth.Context`. + + The same private assign key is used by Plug connections and LiveView sockets. + Assigning a value that is not a context raises `ArgumentError`; fetching a + missing or malformed assign returns `:error` and never exposes it as an + authenticated context. + """ + + alias Phoenix.LiveView.Socket + alias Smolquery.Auth.Context + + @assign_key :smolquery_auth_context + + @doc """ + Returns the stable assign key used for auth contexts. + """ + @spec assign_key() :: :smolquery_auth_context + def assign_key, do: @assign_key + + @doc """ + Assigns a context to a Plug connection or LiveView socket. + """ + @spec assign_context(Plug.Conn.t() | Socket.t(), Context.t()) :: + Plug.Conn.t() | Socket.t() + def assign_context(%Plug.Conn{} = conn, %Context{} = context), + do: assign_valid_context(conn, context) + + def assign_context(%Socket{} = socket, %Context{} = context), + do: assign_valid_context(socket, context) + + def assign_context(_target, _context), + do: raise(ArgumentError, "expected a Smolquery.Auth.Context") + + @doc """ + Fetches an assigned context from a Plug connection or LiveView socket. + """ + @spec fetch_context(Plug.Conn.t() | Socket.t()) :: {:ok, Context.t()} | :error + def fetch_context(%Plug.Conn{assigns: assigns}), do: fetch_assign(assigns) + def fetch_context(%Socket{assigns: assigns}), do: fetch_assign(assigns) + def fetch_context(_target), do: :error + + defp fetch_assign(%{@assign_key => %Context{} = context}) do + if Context.valid?(context), do: {:ok, context}, else: :error + end + + defp fetch_assign(_assigns), do: :error + + defp assign_valid_context(target, context) do + if Context.valid?(context) do + do_assign(target, context) + else + raise ArgumentError, "expected a valid Smolquery.Auth.Context" + end + end + + defp do_assign(%Plug.Conn{} = conn, context), do: Plug.Conn.assign(conn, @assign_key, context) + + defp do_assign(%Socket{} = socket, context), + do: Phoenix.Component.assign(socket, @assign_key, context) +end diff --git a/lib/smolquery/auth/context.ex b/lib/smolquery/auth/context.ex new file mode 100644 index 00000000..1070af43 --- /dev/null +++ b/lib/smolquery/auth/context.ex @@ -0,0 +1,105 @@ +defmodule Smolquery.Auth.Context do + @moduledoc """ + The authenticated principal and capabilities for one request. + + `:single_tenant` is an explicit scope sentinel for the current deployment. + It is intentionally not derived from an issuer, email domain, provider group, + or other identity-provider claim; a future tenant membership lookup can + replace it without changing the principal contract. + """ + + alias Smolquery.Auth.Principal + + @capabilities [:web_access, :query, :ingest, :catalog_manage, :platform_operate] + @single_tenant :single_tenant + + @enforce_keys [:principal, :scope, :capabilities] + defstruct [:principal, :scope, :capabilities] + + @type capability :: :web_access | :query | :ingest | :catalog_manage | :platform_operate + @type scope :: :single_tenant + @type t :: %__MODULE__{ + principal: Principal.t(), + scope: scope(), + capabilities: MapSet.t(capability()) + } + + @doc """ + Builds a context with the explicit single-tenant scope. + """ + @spec single_tenant(Principal.t(), capability() | [capability()] | MapSet.t()) :: + {:ok, t()} | {:error, term()} + def single_tenant(%Principal{} = principal, capabilities) do + if Principal.valid?(principal) do + with {:ok, capabilities} <- normalize_capabilities(capabilities) do + {:ok, + %__MODULE__{principal: principal, scope: @single_tenant, capabilities: capabilities}} + end + else + {:error, :invalid_principal} + end + end + + def single_tenant(_principal, _capabilities), do: {:error, :invalid_principal} + + @doc """ + Reports whether a context grants a known capability. + + Unknown capability checks return `false` rather than converting input into an + atom or granting access by accident. + """ + @spec granted?(t(), term()) :: boolean() + def granted?(%__MODULE__{} = context, capability) when capability in @capabilities do + valid?(context) and MapSet.member?(context.capabilities, capability) + end + + def granted?(_context, _capability), do: false + + @doc """ + Reports whether a term is a valid context produced by this module. + """ + @spec valid?(term()) :: boolean() + def valid?(%__MODULE__{principal: principal, scope: @single_tenant, capabilities: capabilities}) do + Principal.valid?(principal) and valid_capability_set?(capabilities) + end + + def valid?(_term), do: false + + @doc """ + Returns the closed set of capabilities accepted by the context contract. + """ + @spec capabilities() :: [capability()] + def capabilities, do: @capabilities + + @doc """ + Returns the explicit single-tenant scope sentinel. + """ + @spec single_tenant_scope() :: scope() + def single_tenant_scope, do: @single_tenant + + defp normalize_capabilities(capability) when capability in @capabilities do + {:ok, MapSet.new([capability])} + end + + defp normalize_capabilities(%MapSet{map: map} = capabilities) when is_map(map) do + normalize_capabilities(MapSet.to_list(capabilities)) + end + + defp normalize_capabilities(capabilities) when is_list(capabilities) do + if Enum.all?(capabilities, &(&1 in @capabilities)) do + {:ok, MapSet.new(capabilities)} + else + {:error, :invalid_capabilities} + end + end + + defp normalize_capabilities(_capabilities), do: {:error, :invalid_capabilities} + + defp valid_capability_set?(%MapSet{map: map} = capabilities) when is_map(map) do + capabilities + |> MapSet.to_list() + |> Enum.all?(&(&1 in @capabilities)) + end + + defp valid_capability_set?(_capabilities), do: false +end diff --git a/lib/smolquery/auth/policy.ex b/lib/smolquery/auth/policy.ex new file mode 100644 index 00000000..ed28d03c --- /dev/null +++ b/lib/smolquery/auth/policy.ex @@ -0,0 +1,28 @@ +defmodule Smolquery.Auth.Policy do + @moduledoc """ + Capability authorization for authenticated contexts. + + The tagged errors preserve the boundary between a missing identity (`401`) + and an identified principal without the requested capability (`403`). + """ + + alias Smolquery.Auth.Context + + @type authorization_error :: :unauthenticated | :forbidden + + @doc """ + Authorizes a capability for a context. + """ + @spec authorize(nil | Context.t(), term()) :: :ok | {:error, authorization_error()} + def authorize(nil, _capability), do: {:error, :unauthenticated} + + def authorize(%Context{} = context, capability) do + cond do + not Context.valid?(context) -> {:error, :unauthenticated} + Context.granted?(context, capability) -> :ok + true -> {:error, :forbidden} + end + end + + def authorize(_context, _capability), do: {:error, :unauthenticated} +end diff --git a/lib/smolquery/auth/principal.ex b/lib/smolquery/auth/principal.ex new file mode 100644 index 00000000..c1bc3e0f --- /dev/null +++ b/lib/smolquery/auth/principal.ex @@ -0,0 +1,196 @@ +defmodule Smolquery.Auth.Principal do + @moduledoc """ + Provider-neutral identity for an authenticated actor. + + An OIDC principal's stable identity is `{issuer, subject}`. Its `id` is the + versioned, namespaced SHA-256 digest of the exact bytes of those values, + encoded with their lengths to avoid delimiter collisions: + + "oidc:v1:" <> Base.url_encode64(:crypto.hash(:sha256, <>), padding: false) + + The digest is opaque and does not contain display metadata. `display_name` + and `client_id` are selected descriptive values only; raw tokens and claims + are deliberately not represented here. + + Local principals are for credentials managed by smolquery. Their identity is + supplied by the caller and must not use the `:oidc` authentication method. + """ + + @authn [:oidc, :api_key, :basic] + @kinds [:user, :service] + @option_keys [:display_name, :client_id, :expires_at] + @oidc_prefix "oidc:v1:" + + @enforce_keys [:id, :authn, :kind] + defstruct [ + :id, + :authn, + :kind, + :issuer, + :subject, + :display_name, + :client_id, + :expires_at + ] + + @type authn :: :oidc | :api_key | :basic + @type kind :: :user | :service + @type t :: %__MODULE__{ + id: String.t(), + authn: authn(), + kind: kind(), + issuer: String.t() | nil, + subject: String.t() | nil, + display_name: String.t() | nil, + client_id: String.t() | nil, + expires_at: non_neg_integer() | nil + } + + @doc """ + Builds an OIDC principal from an exact issuer and subject. + + Supported options are `:display_name`, `:client_id`, and `:expires_at`. + """ + @spec oidc(String.t(), String.t(), kind(), keyword()) :: {:ok, t()} | {:error, term()} + def oidc(issuer, subject, kind, opts \\ []) do + with :ok <- validate_string(issuer, :issuer), + :ok <- validate_string(subject, :subject), + :ok <- validate_kind(kind), + {:ok, attributes} <- options(opts) do + {:ok, + %__MODULE__{ + id: oidc_id(issuer, subject), + authn: :oidc, + kind: kind, + issuer: issuer, + subject: subject, + display_name: attributes.display_name, + client_id: attributes.client_id, + expires_at: attributes.expires_at + }} + end + end + + @doc """ + Builds a principal for a smolquery-managed credential. + + The `authn` value must be `:api_key` or `:basic`; OIDC identities must use + `oidc/4` so their stable identity cannot be caller-selected. + """ + @spec local(String.t(), authn(), kind(), keyword()) :: {:ok, t()} | {:error, term()} + def local(id, authn, kind, opts \\ []) do + with :ok <- validate_string(id, :id), + :ok <- validate_local_authn(authn), + :ok <- validate_kind(kind), + {:ok, attributes} <- options(opts) do + {:ok, + %__MODULE__{ + id: id, + authn: authn, + kind: kind, + issuer: nil, + subject: nil, + display_name: attributes.display_name, + client_id: attributes.client_id, + expires_at: attributes.expires_at + }} + end + end + + @doc """ + Reports whether a term is a valid principal produced by this module. + """ + @spec valid?(term()) :: boolean() + def valid?(%__MODULE__{} = principal) do + valid_string?(principal.id) and + principal.authn in @authn and + principal.kind in @kinds and + valid_identity?(principal) and + valid_optional_string?(principal.display_name) and + valid_optional_string?(principal.client_id) and + valid_expiry?(principal.expires_at) + end + + def valid?(_term), do: false + + defp valid_identity?(%__MODULE__{id: id, authn: :oidc, issuer: issuer, subject: subject}) do + valid_string?(issuer) and valid_string?(subject) and id == oidc_id(issuer, subject) + end + + defp valid_identity?(%__MODULE__{authn: authn, issuer: nil, subject: nil}) + when authn in [:api_key, :basic], + do: true + + defp valid_identity?(_principal), do: false + + defp options(opts) when is_list(opts) do + cond do + not Keyword.keyword?(opts) -> {:error, :invalid_options} + duplicate_keys?(opts) -> {:error, :invalid_options} + true -> Enum.reduce_while(opts, {:ok, defaults()}, &reduce_option/2) + end + end + + defp options(_opts), do: {:error, :invalid_options} + + defp reduce_option({key, value}, {:ok, attributes}) when key in @option_keys do + case validate_option(key, value) do + :ok -> {:cont, {:ok, Map.put(attributes, key, value)}} + {:error, _reason} = error -> {:halt, error} + end + end + + defp reduce_option({key, _value}, _acc), do: {:halt, {:error, {:unknown_option, key}}} + + defp defaults do + %{display_name: nil, client_id: nil, expires_at: nil} + end + + defp duplicate_keys?(opts) do + keys = Keyword.keys(opts) + length(keys) != length(Enum.uniq(keys)) + end + + defp validate_option(key, value) when key in [:display_name, :client_id] do + if is_nil(value) or valid_string?(value) do + :ok + else + {:error, {:invalid_option, key}} + end + end + + defp validate_option(:expires_at, value) do + if is_nil(value) or (is_integer(value) and value >= 0) do + :ok + else + {:error, {:invalid_option, :expires_at}} + end + end + + defp validate_string(value, field) do + if valid_string?(value), do: :ok, else: {:error, {:invalid, field}} + end + + defp valid_string?(value), do: is_binary(value) and byte_size(value) > 0 + + defp valid_optional_string?(nil), do: true + defp valid_optional_string?(value), do: valid_string?(value) + + defp valid_expiry?(nil), do: true + defp valid_expiry?(value), do: is_integer(value) and value >= 0 + + defp validate_kind(kind) when kind in @kinds, do: :ok + defp validate_kind(_kind), do: {:error, :invalid_kind} + + defp validate_local_authn(authn) when authn in [:api_key, :basic], do: :ok + defp validate_local_authn(:oidc), do: {:error, :oidc_requires_oidc_constructor} + defp validate_local_authn(_authn), do: {:error, :invalid_authn} + + defp oidc_id(issuer, subject) do + encoded = + <> + + @oidc_prefix <> Base.url_encode64(:crypto.hash(:sha256, encoded), padding: false) + end +end diff --git a/test/smolquery/auth/context_test.exs b/test/smolquery/auth/context_test.exs new file mode 100644 index 00000000..28fa58b7 --- /dev/null +++ b/test/smolquery/auth/context_test.exs @@ -0,0 +1,71 @@ +defmodule Smolquery.Auth.ContextTest do + use ExUnit.Case, async: true + + alias Smolquery.Auth.Context + alias Smolquery.Auth.Principal + + defp principal do + {:ok, principal} = Principal.local("static:test", :api_key, :service) + principal + end + + test "uses an explicit single-tenant sentinel" do + assert {:ok, context} = Context.single_tenant(principal(), [:query, :query, :ingest]) + + assert context.scope == :single_tenant + assert context.scope == Context.single_tenant_scope() + assert context.capabilities == MapSet.new([:query, :ingest]) + assert Context.valid?(context) + end + + test "accepts one capability or a MapSet and deduplicates" do + assert {:ok, one} = Context.single_tenant(principal(), :query) + assert {:ok, many} = Context.single_tenant(principal(), MapSet.new([:query, :query])) + + assert one.capabilities == MapSet.new([:query]) + assert many.capabilities == MapSet.new([:query]) + assert Context.granted?(one, :query) + refute Context.granted?(one, :ingest) + end + + test "rejects unknown capabilities without converting input" do + assert {:error, :invalid_capabilities} = Context.single_tenant(principal(), [:query, :admin]) + assert {:error, :invalid_capabilities} = Context.single_tenant(principal(), ["query"]) + assert {:error, :invalid_capabilities} = Context.single_tenant(principal(), :admin) + refute Context.granted?(nil, :admin) + end + + test "rejects a malformed principal" do + malformed = %Principal{id: "", authn: :api_key, kind: :service} + + assert {:error, :invalid_principal} = Context.single_tenant(malformed, :query) + end + + test "does not grant from a forged context" do + malformed = %Principal{id: "", authn: :api_key, kind: :service} + + context = %Context{ + principal: malformed, + scope: :single_tenant, + capabilities: MapSet.new([:query]) + } + + refute Context.valid?(context) + refute Context.granted?(context, :query) + end + + test "rejects a malformed MapSet without raising" do + capabilities = %MapSet{map: :forged} + + assert {:error, :invalid_capabilities} = Context.single_tenant(principal(), capabilities) + + context = %Context{ + principal: principal(), + scope: :single_tenant, + capabilities: capabilities + } + + refute Context.valid?(context) + refute Context.granted?(context, :query) + end +end diff --git a/test/smolquery/auth/policy_test.exs b/test/smolquery/auth/policy_test.exs new file mode 100644 index 00000000..9a0e0348 --- /dev/null +++ b/test/smolquery/auth/policy_test.exs @@ -0,0 +1,38 @@ +defmodule Smolquery.Auth.PolicyTest do + use ExUnit.Case, async: true + + alias Smolquery.Auth.Context + alias Smolquery.Auth.Policy + alias Smolquery.Auth.Principal + + defp context(capabilities) do + {:ok, principal} = Principal.local("static:test", :api_key, :service) + {:ok, context} = Context.single_tenant(principal, capabilities) + context + end + + test "grants a capability held by the context" do + assert :ok = Policy.authorize(context(:query), :query) + end + + test "forbids a known capability the context does not hold" do + assert {:error, :forbidden} = Policy.authorize(context(:query), :ingest) + assert {:error, :forbidden} = Policy.authorize(context(:query), :unknown) + end + + test "reports a missing or invalid context as unauthenticated" do + assert {:error, :unauthenticated} = Policy.authorize(nil, :query) + assert {:error, :unauthenticated} = Policy.authorize(:not_a_context, :query) + + malformed = %Context{ + principal: %Principal{id: "", authn: :api_key, kind: :service}, + scope: :single_tenant, + capabilities: MapSet.new([:query]) + } + + assert {:error, :unauthenticated} = Policy.authorize(malformed, :query) + + assert {:error, :unauthenticated} = + Policy.authorize(%{malformed | capabilities: %MapSet{map: :forged}}, :query) + end +end diff --git a/test/smolquery/auth/principal_test.exs b/test/smolquery/auth/principal_test.exs new file mode 100644 index 00000000..ec6b9006 --- /dev/null +++ b/test/smolquery/auth/principal_test.exs @@ -0,0 +1,112 @@ +defmodule Smolquery.Auth.PrincipalTest do + use ExUnit.Case, async: true + + alias Smolquery.Auth.Principal + + describe "oidc/4" do + test "derives a stable opaque identity from issuer and subject" do + assert {:ok, first} = Principal.oidc("https://idp.example", "user-1", :user) + assert {:ok, second} = Principal.oidc("https://idp.example", "user-1", :user) + + assert first.id == second.id + assert first.id == "oidc:v1:D69mFP3Atlfs9LemldrEtUY3onuuPpAvhUyuPWsUYsw" + assert first.id != "https://idp.example" + assert first.id != "user-1" + assert first.issuer == "https://idp.example" + assert first.subject == "user-1" + assert Principal.valid?(first) + end + + test "distinguishes the same subject at different issuers" do + assert {:ok, first} = Principal.oidc("https://one.example", "same", :user) + assert {:ok, second} = Principal.oidc("https://two.example", "same", :user) + + refute first.id == second.id + end + + test "keeps issuer and subject boundaries distinct" do + assert {:ok, first} = Principal.oidc("a\0b", "c", :user) + assert {:ok, second} = Principal.oidc("a", "b\0c", :user) + + refute first.id == second.id + end + + test "does not let display metadata affect identity" do + assert {:ok, without_metadata} = Principal.oidc("issuer", "subject", :user) + + assert {:ok, with_metadata} = + Principal.oidc("issuer", "subject", :user, + display_name: "Alice", + client_id: "web", + expires_at: 42 + ) + + assert with_metadata.id == without_metadata.id + assert with_metadata.display_name == "Alice" + assert with_metadata.client_id == "web" + assert with_metadata.expires_at == 42 + end + + test "requires exact non-empty issuer and subject strings" do + assert {:error, {:invalid, :issuer}} = Principal.oidc("", "subject", :user) + assert {:error, {:invalid, :issuer}} = Principal.oidc(:issuer, "subject", :user) + assert {:error, {:invalid, :subject}} = Principal.oidc("issuer", "", :user) + assert {:error, {:invalid, :subject}} = Principal.oidc("issuer", :subject, :user) + end + end + + describe "local/4" do + test "builds local API-key and Basic principals" do + assert {:ok, api_key} = Principal.local("static:api", :api_key, :service) + assert {:ok, basic} = Principal.local("static:web", :basic, :user) + + assert api_key.authn == :api_key + assert basic.authn == :basic + assert api_key.issuer == nil + assert api_key.subject == nil + assert Principal.valid?(api_key) + assert Principal.valid?(basic) + end + + test "does not allow local OIDC identities" do + assert {:error, :oidc_requires_oidc_constructor} = + Principal.local("id", :oidc, :user) + end + end + + describe "validation" do + test "rejects invalid kinds, options, metadata, and expiry" do + assert {:error, :invalid_kind} = Principal.oidc("issuer", "subject", :admin) + + assert {:error, {:unknown_option, :groups}} = + Principal.oidc("issuer", "subject", :user, groups: ["admins"]) + + assert {:error, {:invalid_option, :display_name}} = + Principal.local("id", :api_key, :service, display_name: "") + + assert {:error, {:invalid_option, :expires_at}} = + Principal.local("id", :api_key, :service, expires_at: -1) + + assert {:error, :invalid_options} = + Principal.local("id", :api_key, :service, %{expires_at: 1}) + + assert {:error, :invalid_options} = + Principal.local("id", :api_key, :service, expires_at: 1, expires_at: 2) + end + + test "rejects invalid authentication values and malformed principals" do + assert {:error, :invalid_authn} = Principal.local("id", :password, :service) + assert {:error, :invalid_kind} = Principal.local("id", :api_key, :admin) + + refute Principal.valid?(%Principal{ + id: "id", + authn: :api_key, + kind: :service, + issuer: "unexpected" + }) + + assert {:ok, oidc} = Principal.oidc("issuer", "subject", :user) + refute Principal.valid?(%{oidc | id: "oidc:v1:forged"}) + end + end +end diff --git a/test/smolquery/auth_test.exs b/test/smolquery/auth_test.exs new file mode 100644 index 00000000..f00eb983 --- /dev/null +++ b/test/smolquery/auth_test.exs @@ -0,0 +1,55 @@ +defmodule Smolquery.AuthTest do + use ExUnit.Case, async: true + + import Plug.Test + + alias Phoenix.LiveView.Socket + alias Smolquery.Auth + alias Smolquery.Auth.Context + alias Smolquery.Auth.Principal + + defp context do + {:ok, principal} = Principal.local("static:test", :api_key, :service) + {:ok, context} = Context.single_tenant(principal, :query) + context + end + + test "round-trips a context on a Plug connection" do + conn = Auth.assign_context(conn(:get, "/"), context()) + + assert {:ok, fetched} = Auth.fetch_context(conn) + assert fetched == context() + assert conn.assigns[Auth.assign_key()] == fetched + end + + test "round-trips a context on a LiveView socket" do + socket = Auth.assign_context(%Socket{}, context()) + + assert {:ok, fetched} = Auth.fetch_context(socket) + assert fetched == context() + assert socket.assigns[Auth.assign_key()] == fetched + end + + test "rejects a malformed assignment" do + conn = Plug.Conn.assign(conn(:get, "/"), Auth.assign_key(), %{principal: :invalid}) + + assert :error = Auth.fetch_context(conn) + assert_raise ArgumentError, fn -> Auth.assign_context(conn, %{principal: :invalid}) end + end + + test "rejects malformed context internals without raising" do + forged = %{context() | capabilities: %MapSet{map: :forged}} + conn = Plug.Conn.assign(conn(:get, "/"), Auth.assign_key(), forged) + socket = Phoenix.Component.assign(%Socket{}, Auth.assign_key(), forged) + + assert :error = Auth.fetch_context(conn) + assert :error = Auth.fetch_context(socket) + assert_raise ArgumentError, fn -> Auth.assign_context(conn, forged) end + assert_raise ArgumentError, fn -> Auth.assign_context(socket, forged) end + end + + test "rejects unsupported targets" do + assert :error = Auth.fetch_context(%{}) + assert_raise ArgumentError, fn -> Auth.assign_context(%{}, context()) end + end +end From 1ed0b82ed745706ae601f72c15e0ef83c7e645e7 Mon Sep 17 00:00:00 2001 From: David Whittington Date: Sat, 15 Aug 2026 03:10:41 +0000 Subject: [PATCH 2/8] fix(auth): tighten principal and context trust contracts Namespace local principal identities, separate assertion expiry from stable principals, enforce expiry during policy decisions, and document the structural-validation and pseudonymous-identifier boundaries. Refs T-229 and PL-27. --- lib/smolquery/auth.ex | 14 ++- lib/smolquery/auth/context.ex | 135 +++++++++++++++++----- lib/smolquery/auth/policy.ex | 36 ++++-- lib/smolquery/auth/principal.ex | 153 +++++++++++++++---------- test/smolquery/auth/context_test.exs | 76 +++++++++--- test/smolquery/auth/policy_test.exs | 31 +++-- test/smolquery/auth/principal_test.exs | 69 +++++++---- test/smolquery/auth_test.exs | 4 + 8 files changed, 361 insertions(+), 157 deletions(-) diff --git a/lib/smolquery/auth.ex b/lib/smolquery/auth.ex index 7405d7d7..2dd6326d 100644 --- a/lib/smolquery/auth.ex +++ b/lib/smolquery/auth.ex @@ -4,8 +4,9 @@ defmodule Smolquery.Auth do The same private assign key is used by Plug connections and LiveView sockets. Assigning a value that is not a context raises `ArgumentError`; fetching a - missing or malformed assign returns `:error` and never exposes it as an - authenticated context. + missing or malformed assign returns `:error`. Assignment and fetching prove + structure only; trusted adapters establish authentication provenance before + assignment, and policy authorization checks expiry and capability access. """ alias Phoenix.LiveView.Socket @@ -35,6 +36,9 @@ defmodule Smolquery.Auth do @doc """ Fetches an assigned context from a Plug connection or LiveView socket. + + Assignment checks are structural only; policy authorization checks expiry and + capability access. """ @spec fetch_context(Plug.Conn.t() | Socket.t()) :: {:ok, Context.t()} | :error def fetch_context(%Plug.Conn{assigns: assigns}), do: fetch_assign(assigns) @@ -42,16 +46,16 @@ defmodule Smolquery.Auth do def fetch_context(_target), do: :error defp fetch_assign(%{@assign_key => %Context{} = context}) do - if Context.valid?(context), do: {:ok, context}, else: :error + if Context.well_formed?(context), do: {:ok, context}, else: :error end defp fetch_assign(_assigns), do: :error defp assign_valid_context(target, context) do - if Context.valid?(context) do + if Context.well_formed?(context) do do_assign(target, context) else - raise ArgumentError, "expected a valid Smolquery.Auth.Context" + raise ArgumentError, "expected a well-formed Smolquery.Auth.Context" end end diff --git a/lib/smolquery/auth/context.ex b/lib/smolquery/auth/context.ex index 1070af43..39a66d05 100644 --- a/lib/smolquery/auth/context.ex +++ b/lib/smolquery/auth/context.ex @@ -1,69 +1,117 @@ defmodule Smolquery.Auth.Context do @moduledoc """ - The authenticated principal and capabilities for one request. + The authenticated principal, capabilities, scope, and optional expiry for + one request. `:single_tenant` is an explicit scope sentinel for the current deployment. - It is intentionally not derived from an issuer, email domain, provider group, - or other identity-provider claim; a future tenant membership lookup can - replace it without changing the principal contract. + It is intentionally not derived from an issuer, email domain, provider + group, or other identity-provider claim; a future tenant membership lookup + can replace it without changing the principal contract. + + `expires_at` is an optional Unix epoch timestamp in integer seconds. A + context is active only while `now < expires_at`; the context is expired at + exactly the boundary. Authenticator clock-skew handling occurs before this + context is constructed. + + Structural checks prove shape only, not authentication provenance. Only + trusted authenticators and mappers may construct contexts; client input must + not be decoded directly into structs or capabilities. Authentication + provenance is established before construction; policy checks expiry and + capability access. """ alias Smolquery.Auth.Principal @capabilities [:web_access, :query, :ingest, :catalog_manage, :platform_operate] @single_tenant :single_tenant - @enforce_keys [:principal, :scope, :capabilities] - defstruct [:principal, :scope, :capabilities] + defstruct [:principal, :scope, :capabilities, :expires_at] @type capability :: :web_access | :query | :ingest | :catalog_manage | :platform_operate @type scope :: :single_tenant @type t :: %__MODULE__{ principal: Principal.t(), scope: scope(), - capabilities: MapSet.t(capability()) + capabilities: MapSet.t(capability()), + expires_at: non_neg_integer() | nil } @doc """ Builds a context with the explicit single-tenant scope. + + The optional `:expires_at` option is an integer Unix epoch timestamp in + seconds and defaults to `nil` for no expiry. """ @spec single_tenant(Principal.t(), capability() | [capability()] | MapSet.t()) :: {:ok, t()} | {:error, term()} - def single_tenant(%Principal{} = principal, capabilities) do - if Principal.valid?(principal) do - with {:ok, capabilities} <- normalize_capabilities(capabilities) do + @spec single_tenant(Principal.t(), capability() | [capability()] | MapSet.t(), keyword()) :: + {:ok, t()} | {:error, term()} + def single_tenant(principal, capabilities, opts \\ []) + + def single_tenant(%Principal{} = principal, capabilities, opts) do + if Principal.well_formed?(principal) do + with {:ok, capabilities} <- normalize_capabilities(capabilities), + {:ok, expires_at} <- options(opts) do {:ok, - %__MODULE__{principal: principal, scope: @single_tenant, capabilities: capabilities}} + %__MODULE__{ + principal: principal, + scope: @single_tenant, + capabilities: capabilities, + expires_at: expires_at + }} end else {:error, :invalid_principal} end end - def single_tenant(_principal, _capabilities), do: {:error, :invalid_principal} + def single_tenant(_principal, _capabilities, _opts), do: {:error, :invalid_principal} @doc """ - Reports whether a context grants a known capability. - - Unknown capability checks return `false` rather than converting input into an - atom or granting access by accident. + Reports whether a well-formed context grants a known capability. """ @spec granted?(t(), term()) :: boolean() def granted?(%__MODULE__{} = context, capability) when capability in @capabilities do - valid?(context) and MapSet.member?(context.capabilities, capability) + well_formed?(context) and MapSet.member?(context.capabilities, capability) end def granted?(_context, _capability), do: false @doc """ - Reports whether a term is a valid context produced by this module. + Reports whether a capability belongs to the closed capability contract. + """ + @spec capability?(term()) :: boolean() + def capability?(capability), do: capability in @capabilities + + @doc """ + Reports whether a term has the structure of a context produced by this + module. This does not prove authentication provenance. """ - @spec valid?(term()) :: boolean() - def valid?(%__MODULE__{principal: principal, scope: @single_tenant, capabilities: capabilities}) do - Principal.valid?(principal) and valid_capability_set?(capabilities) + @spec well_formed?(term()) :: boolean() + def well_formed?(%__MODULE__{ + principal: principal, + scope: @single_tenant, + capabilities: capabilities, + expires_at: expires_at + }) do + Principal.well_formed?(principal) and valid_capability_set?(capabilities) and + valid_expiry?(expires_at) end - def valid?(_term), do: false + def well_formed?(_term), do: false + + @doc """ + Reports whether a context is active at a non-negative Unix epoch timestamp. + + An expiry is strict: a context is inactive when `now == expires_at`. + """ + @spec active?(term(), non_neg_integer()) :: boolean() + def active?(%__MODULE__{expires_at: expires_at} = context, now) + when is_integer(now) and now >= 0 do + well_formed?(context) and (is_nil(expires_at) or now < expires_at) + end + + def active?(_context, _now), do: false @doc """ Returns the closed set of capabilities accepted by the context contract. @@ -82,7 +130,11 @@ defmodule Smolquery.Auth.Context do end defp normalize_capabilities(%MapSet{map: map} = capabilities) when is_map(map) do - normalize_capabilities(MapSet.to_list(capabilities)) + if valid_mapset_map?(map) do + normalize_capabilities(MapSet.to_list(capabilities)) + else + {:error, :invalid_capabilities} + end end defp normalize_capabilities(capabilities) when is_list(capabilities) do @@ -95,11 +147,38 @@ defmodule Smolquery.Auth.Context do defp normalize_capabilities(_capabilities), do: {:error, :invalid_capabilities} - defp valid_capability_set?(%MapSet{map: map} = capabilities) when is_map(map) do - capabilities - |> MapSet.to_list() - |> Enum.all?(&(&1 in @capabilities)) - end + defp valid_capability_set?(%MapSet{map: map}) when is_map(map), + do: valid_mapset_map?(map) defp valid_capability_set?(_capabilities), do: false + + defp valid_mapset_map?(map) do + Enum.all?(map, fn {capability, value} -> value == [] and capability in @capabilities end) + end + + defp options(opts) when is_list(opts) do + cond do + not Keyword.keyword?(opts) -> {:error, :invalid_options} + duplicate_keys?(opts) -> {:error, :invalid_options} + true -> Enum.reduce_while(opts, {:ok, nil}, &reduce_option/2) + end + end + + defp options(_opts), do: {:error, :invalid_options} + + defp reduce_option({:expires_at, value}, _acc) do + if valid_expiry?(value), + do: {:cont, {:ok, value}}, + else: {:halt, {:error, {:invalid_option, :expires_at}}} + end + + defp reduce_option({key, _value}, _acc), do: {:halt, {:error, {:unknown_option, key}}} + + defp duplicate_keys?(opts) do + keys = Keyword.keys(opts) + length(keys) != length(Enum.uniq(keys)) + end + + defp valid_expiry?(nil), do: true + defp valid_expiry?(value), do: is_integer(value) and value >= 0 end diff --git a/lib/smolquery/auth/policy.ex b/lib/smolquery/auth/policy.ex index ed28d03c..f757b139 100644 --- a/lib/smolquery/auth/policy.ex +++ b/lib/smolquery/auth/policy.ex @@ -2,27 +2,45 @@ defmodule Smolquery.Auth.Policy do @moduledoc """ Capability authorization for authenticated contexts. - The tagged errors preserve the boundary between a missing identity (`401`) - and an identified principal without the requested capability (`403`). + Structural context checks prove shape only. Trusted authenticators and + mappers construct contexts, and client input must not be decoded directly + into structs or capabilities. Expiry is checked here using the system Unix + epoch clock; authenticator clock-skew handling occurs before construction. + + The tagged errors preserve the boundary between a missing, malformed, or + expired identity (`:unauthenticated`), an unknown requested capability + (`:invalid_capability`), and an identified principal without a known + requested capability (`:forbidden`). """ alias Smolquery.Auth.Context - @type authorization_error :: :unauthenticated | :forbidden + @type authorization_error :: :unauthenticated | :invalid_capability | :forbidden + + @doc """ + Authorizes a capability for a context using the current Unix epoch time. + """ + @spec authorize(term(), term()) :: :ok | {:error, authorization_error()} + def authorize(context, capability), + do: authorize(context, capability, System.system_time(:second)) @doc """ - Authorizes a capability for a context. + Authorizes a capability for a context at a supplied non-negative Unix epoch + timestamp. This arity supports deterministic authorization tests. """ - @spec authorize(nil | Context.t(), term()) :: :ok | {:error, authorization_error()} - def authorize(nil, _capability), do: {:error, :unauthenticated} + @spec authorize(term(), term(), non_neg_integer()) :: + :ok | {:error, authorization_error()} + def authorize(nil, _capability, _now), do: {:error, :unauthenticated} - def authorize(%Context{} = context, capability) do + def authorize(%Context{} = context, capability, now) do cond do - not Context.valid?(context) -> {:error, :unauthenticated} + not Context.well_formed?(context) -> {:error, :unauthenticated} + not Context.active?(context, now) -> {:error, :unauthenticated} + not Context.capability?(capability) -> {:error, :invalid_capability} Context.granted?(context, capability) -> :ok true -> {:error, :forbidden} end end - def authorize(_context, _capability), do: {:error, :unauthenticated} + def authorize(_context, _capability, _now), do: {:error, :unauthenticated} end diff --git a/lib/smolquery/auth/principal.ex b/lib/smolquery/auth/principal.ex index c1bc3e0f..e86bb9ce 100644 --- a/lib/smolquery/auth/principal.ex +++ b/lib/smolquery/auth/principal.ex @@ -6,34 +6,42 @@ defmodule Smolquery.Auth.Principal do versioned, namespaced SHA-256 digest of the exact bytes of those values, encoded with their lengths to avoid delimiter collisions: - "oidc:v1:" <> Base.url_encode64(:crypto.hash(:sha256, <>), padding: false) - - The digest is opaque and does not contain display metadata. `display_name` - and `client_id` are selected descriptive values only; raw tokens and claims - are deliberately not represented here. - - Local principals are for credentials managed by smolquery. Their identity is - supplied by the caller and must not use the `:oidc` authentication method. + "oidc:v1:" <> Base.url_encode64(:crypto.hash(:sha256, <>), padding: false) + + Local principal source keys are stable, non-secret credential identifiers. + The constructor derives their final IDs in authn-specific, versioned + namespaces. Derived IDs are opaque stable pseudonymous identifiers, not + anonymous or non-PII values. The source key is not retained in the struct. + Local namespaces cannot collide with OIDC or with one another. + + Every identity string is framed by its unsigned big-endian 32-bit byte + length. Digests are SHA-256 values encoded with URL-safe Base64 without + padding. Smolquery treats exact differing `{issuer, subject}` pairs as + distinct identities, including pairwise subjects across clients or sectors; + the derived IDs rely on SHA-256 collision resistance, and mutable display + claims must never merge identities. Display metadata is descriptive only, and raw + tokens and claims are deliberately not represented here. + + Structural checks prove shape only, not authentication provenance. Only + trusted authenticators and mappers may construct principals; client input + must not be decoded directly into structs or capabilities. """ @authn [:oidc, :api_key, :basic] + @local_authn [:api_key, :basic] @kinds [:user, :service] - @option_keys [:display_name, :client_id, :expires_at] + @option_keys [:display_name, :client_id] @oidc_prefix "oidc:v1:" + @local_prefixes %{api_key: "api_key:v1:", basic: "basic:v1:"} + @max_frame_size 4_294_967_295 + @digest_size 32 + @encoded_digest_size 43 @enforce_keys [:id, :authn, :kind] - defstruct [ - :id, - :authn, - :kind, - :issuer, - :subject, - :display_name, - :client_id, - :expires_at - ] - - @type authn :: :oidc | :api_key | :basic + defstruct [:id, :authn, :kind, :issuer, :subject, :display_name, :client_id] + + @type authn :: :oidc | local_authn() + @type local_authn :: :api_key | :basic @type kind :: :user | :service @type t :: %__MODULE__{ id: String.t(), @@ -42,14 +50,13 @@ defmodule Smolquery.Auth.Principal do issuer: String.t() | nil, subject: String.t() | nil, display_name: String.t() | nil, - client_id: String.t() | nil, - expires_at: non_neg_integer() | nil + client_id: String.t() | nil } @doc """ Builds an OIDC principal from an exact issuer and subject. - Supported options are `:display_name`, `:client_id`, and `:expires_at`. + Supported options are `:display_name` and `:client_id`. """ @spec oidc(String.t(), String.t(), kind(), keyword()) :: {:ok, t()} | {:error, term()} def oidc(issuer, subject, kind, opts \\ []) do @@ -65,8 +72,7 @@ defmodule Smolquery.Auth.Principal do issuer: issuer, subject: subject, display_name: attributes.display_name, - client_id: attributes.client_id, - expires_at: attributes.expires_at + client_id: attributes.client_id }} end end @@ -74,52 +80,53 @@ defmodule Smolquery.Auth.Principal do @doc """ Builds a principal for a smolquery-managed credential. - The `authn` value must be `:api_key` or `:basic`; OIDC identities must use + `source_key` is a stable, non-secret source key. The constructor derives an + opaque ID from it in an authn-specific namespace. OIDC identities must use `oidc/4` so their stable identity cannot be caller-selected. """ - @spec local(String.t(), authn(), kind(), keyword()) :: {:ok, t()} | {:error, term()} - def local(id, authn, kind, opts \\ []) do - with :ok <- validate_string(id, :id), + @spec local(String.t(), local_authn(), kind(), keyword()) :: {:ok, t()} | {:error, term()} + def local(source_key, authn, kind, opts \\ []) do + with :ok <- validate_string(source_key, :source_key), :ok <- validate_local_authn(authn), :ok <- validate_kind(kind), {:ok, attributes} <- options(opts) do {:ok, %__MODULE__{ - id: id, + id: local_id(source_key, authn), authn: authn, kind: kind, issuer: nil, subject: nil, display_name: attributes.display_name, - client_id: attributes.client_id, - expires_at: attributes.expires_at + client_id: attributes.client_id }} end end @doc """ - Reports whether a term is a valid principal produced by this module. + Reports whether a term has the structure of a principal produced by this + module. This does not prove authentication provenance. """ - @spec valid?(term()) :: boolean() - def valid?(%__MODULE__{} = principal) do + @spec well_formed?(term()) :: boolean() + def well_formed?(%__MODULE__{} = principal) do valid_string?(principal.id) and principal.authn in @authn and principal.kind in @kinds and valid_identity?(principal) and valid_optional_string?(principal.display_name) and - valid_optional_string?(principal.client_id) and - valid_expiry?(principal.expires_at) + valid_optional_string?(principal.client_id) end - def valid?(_term), do: false + def well_formed?(_term), do: false defp valid_identity?(%__MODULE__{id: id, authn: :oidc, issuer: issuer, subject: subject}) do valid_string?(issuer) and valid_string?(subject) and id == oidc_id(issuer, subject) end - defp valid_identity?(%__MODULE__{authn: authn, issuer: nil, subject: nil}) - when authn in [:api_key, :basic], - do: true + defp valid_identity?(%__MODULE__{id: id, authn: authn, issuer: nil, subject: nil}) + when authn in @local_authn do + local_id_shape?(id, authn) + end defp valid_identity?(_principal), do: false @@ -142,9 +149,7 @@ defmodule Smolquery.Auth.Principal do defp reduce_option({key, _value}, _acc), do: {:halt, {:error, {:unknown_option, key}}} - defp defaults do - %{display_name: nil, client_id: nil, expires_at: nil} - end + defp defaults, do: %{display_name: nil, client_id: nil} defp duplicate_keys?(opts) do keys = Keyword.keys(opts) @@ -159,38 +164,62 @@ defmodule Smolquery.Auth.Principal do end end - defp validate_option(:expires_at, value) do - if is_nil(value) or (is_integer(value) and value >= 0) do - :ok - else - {:error, {:invalid_option, :expires_at}} - end - end - defp validate_string(value, field) do if valid_string?(value), do: :ok, else: {:error, {:invalid, field}} end - defp valid_string?(value), do: is_binary(value) and byte_size(value) > 0 + defp valid_string?(value), + do: is_binary(value) and byte_size(value) > 0 and byte_size(value) <= @max_frame_size defp valid_optional_string?(nil), do: true defp valid_optional_string?(value), do: valid_string?(value) - defp valid_expiry?(nil), do: true - defp valid_expiry?(value), do: is_integer(value) and value >= 0 - defp validate_kind(kind) when kind in @kinds, do: :ok defp validate_kind(_kind), do: {:error, :invalid_kind} - defp validate_local_authn(authn) when authn in [:api_key, :basic], do: :ok + defp validate_local_authn(authn) when authn in @local_authn, do: :ok defp validate_local_authn(:oidc), do: {:error, :oidc_requires_oidc_constructor} defp validate_local_authn(_authn), do: {:error, :invalid_authn} defp oidc_id(issuer, subject) do - encoded = - <> + encoded = frame(issuer, subject) + @oidc_prefix <> digest(encoded) + end - @oidc_prefix <> Base.url_encode64(:crypto.hash(:sha256, encoded), padding: false) + defp local_id(source_key, authn) do + namespace = Map.fetch!(@local_prefixes, authn) + namespace <> digest(frame(namespace, source_key)) end + + defp frame(first, second) do + <> + end + + defp digest(bytes), do: Base.url_encode64(:crypto.hash(:sha256, bytes), padding: false) + + defp local_id_shape?(id, authn) do + prefix = Map.fetch!(@local_prefixes, authn) + prefix_size = byte_size(prefix) + + case id do + <> when candidate == prefix -> + canonical_digest?(digest_part) + + _other -> + false + end + end + + defp canonical_digest?(digest_part) when byte_size(digest_part) == @encoded_digest_size do + case Base.url_decode64(digest_part, padding: false) do + {:ok, <<_::binary-size(@digest_size)>> = decoded} -> + Base.url_encode64(decoded, padding: false) == digest_part + + :error -> + false + end + end + + defp canonical_digest?(_digest_part), do: false end diff --git a/test/smolquery/auth/context_test.exs b/test/smolquery/auth/context_test.exs index 28fa58b7..26ae0ab2 100644 --- a/test/smolquery/auth/context_test.exs +++ b/test/smolquery/auth/context_test.exs @@ -9,39 +9,75 @@ defmodule Smolquery.Auth.ContextTest do principal end - test "uses an explicit single-tenant sentinel" do - assert {:ok, context} = Context.single_tenant(principal(), [:query, :query, :ingest]) + test "uses an explicit single-tenant sentinel and optional expiry" do + assert {:ok, context} = + Context.single_tenant(principal(), [:query, :query, :ingest], expires_at: 100) assert context.scope == :single_tenant assert context.scope == Context.single_tenant_scope() assert context.capabilities == MapSet.new([:query, :ingest]) - assert Context.valid?(context) + assert context.expires_at == 100 + assert Context.well_formed?(context) end - test "accepts one capability or a MapSet and deduplicates" do + test "defaults expiry to nil and accepts one capability or a MapSet" do assert {:ok, one} = Context.single_tenant(principal(), :query) assert {:ok, many} = Context.single_tenant(principal(), MapSet.new([:query, :query])) + assert one.expires_at == nil assert one.capabilities == MapSet.new([:query]) assert many.capabilities == MapSet.new([:query]) assert Context.granted?(one, :query) refute Context.granted?(one, :ingest) end - test "rejects unknown capabilities without converting input" do + test "identifies known capabilities without converting input" do + assert Context.capability?(:query) + refute Context.capability?(:admin) + refute Context.capability?("query") + assert {:error, :invalid_capabilities} = Context.single_tenant(principal(), [:query, :admin]) assert {:error, :invalid_capabilities} = Context.single_tenant(principal(), ["query"]) assert {:error, :invalid_capabilities} = Context.single_tenant(principal(), :admin) refute Context.granted?(nil, :admin) end - test "rejects a malformed principal" do + test "rejects malformed expiry options" do + assert {:error, {:invalid_option, :expires_at}} = + Context.single_tenant(principal(), :query, expires_at: -1) + + assert {:error, {:invalid_option, :expires_at}} = + Context.single_tenant(principal(), :query, expires_at: "later") + + assert {:error, :invalid_options} = + Context.single_tenant(principal(), :query, %{expires_at: 1}) + end + + test "expires before and at the boundary" do + assert {:ok, context} = Context.single_tenant(principal(), :query, expires_at: 100) + + assert Context.active?(context, 99) + refute Context.active?(context, 100) + refute Context.active?(context, 101) + refute Context.active?(context, -1) + refute Context.active?(context, "100") + + assert {:ok, never_expires} = Context.single_tenant(principal(), :query) + assert Context.active?(never_expires, 0) + end + + test "rejects malformed principal and expiry structure" do malformed = %Principal{id: "", authn: :api_key, kind: :service} assert {:error, :invalid_principal} = Context.single_tenant(malformed, :query) + + assert {:ok, context} = Context.single_tenant(principal(), :query) + malformed_expiry = %{context | expires_at: -1} + refute Context.well_formed?(malformed_expiry) + refute Context.active?(malformed_expiry, 0) end - test "does not grant from a forged context" do + test "structural checks do not grant from a forged context" do malformed = %Principal{id: "", authn: :api_key, kind: :service} context = %Context{ @@ -50,22 +86,26 @@ defmodule Smolquery.Auth.ContextTest do capabilities: MapSet.new([:query]) } - refute Context.valid?(context) + refute Context.well_formed?(context) refute Context.granted?(context, :query) end - test "rejects a malformed MapSet without raising" do - capabilities = %MapSet{map: :forged} + test "rejects malformed MapSet internals without raising" do + non_map = %MapSet{map: :forged} + invalid_entry = %MapSet{map: %{query: :forged}} - assert {:error, :invalid_capabilities} = Context.single_tenant(principal(), capabilities) + assert {:error, :invalid_capabilities} = Context.single_tenant(principal(), non_map) + assert {:error, :invalid_capabilities} = Context.single_tenant(principal(), invalid_entry) - context = %Context{ - principal: principal(), - scope: :single_tenant, - capabilities: capabilities - } + for capabilities <- [non_map, invalid_entry] do + context = %Context{ + principal: principal(), + scope: :single_tenant, + capabilities: capabilities + } - refute Context.valid?(context) - refute Context.granted?(context, :query) + refute Context.well_formed?(context) + refute Context.granted?(context, :query) + end end end diff --git a/test/smolquery/auth/policy_test.exs b/test/smolquery/auth/policy_test.exs index 9a0e0348..ae503bb0 100644 --- a/test/smolquery/auth/policy_test.exs +++ b/test/smolquery/auth/policy_test.exs @@ -5,24 +5,30 @@ defmodule Smolquery.Auth.PolicyTest do alias Smolquery.Auth.Policy alias Smolquery.Auth.Principal - defp context(capabilities) do + defp context(capabilities, opts \\ []) do {:ok, principal} = Principal.local("static:test", :api_key, :service) - {:ok, context} = Context.single_tenant(principal, capabilities) + {:ok, context} = Context.single_tenant(principal, capabilities, opts) context end - test "grants a capability held by the context" do - assert :ok = Policy.authorize(context(:query), :query) + test "grants a capability held by an active context" do + assert :ok = Policy.authorize(context(:query), :query, 100) end test "forbids a known capability the context does not hold" do - assert {:error, :forbidden} = Policy.authorize(context(:query), :ingest) - assert {:error, :forbidden} = Policy.authorize(context(:query), :unknown) + assert {:error, :forbidden} = Policy.authorize(context(:query), :ingest, 100) end - test "reports a missing or invalid context as unauthenticated" do - assert {:error, :unauthenticated} = Policy.authorize(nil, :query) - assert {:error, :unauthenticated} = Policy.authorize(:not_a_context, :query) + test "rejects an unknown capability after structural and expiry checks" do + assert {:error, :invalid_capability} = Policy.authorize(context(:query), :unknown, 100) + + assert {:error, :unauthenticated} = + Policy.authorize(context(:query, expires_at: 100), :unknown, 100) + end + + test "reports a missing, malformed, or expired context as unauthenticated" do + assert {:error, :unauthenticated} = Policy.authorize(nil, :query, 100) + assert {:error, :unauthenticated} = Policy.authorize(:not_a_context, :query, 100) malformed = %Context{ principal: %Principal{id: "", authn: :api_key, kind: :service}, @@ -30,9 +36,12 @@ defmodule Smolquery.Auth.PolicyTest do capabilities: MapSet.new([:query]) } - assert {:error, :unauthenticated} = Policy.authorize(malformed, :query) + assert {:error, :unauthenticated} = Policy.authorize(malformed, :query, 100) + + assert {:error, :unauthenticated} = + Policy.authorize(%{malformed | capabilities: %MapSet{map: :forged}}, :query, 100) assert {:error, :unauthenticated} = - Policy.authorize(%{malformed | capabilities: %MapSet{map: :forged}}, :query) + Policy.authorize(context(:query, expires_at: 100), :query, 100) end end diff --git a/test/smolquery/auth/principal_test.exs b/test/smolquery/auth/principal_test.exs index ec6b9006..47e47785 100644 --- a/test/smolquery/auth/principal_test.exs +++ b/test/smolquery/auth/principal_test.exs @@ -14,14 +14,16 @@ defmodule Smolquery.Auth.PrincipalTest do assert first.id != "user-1" assert first.issuer == "https://idp.example" assert first.subject == "user-1" - assert Principal.valid?(first) + assert Principal.well_formed?(first) end - test "distinguishes the same subject at different issuers" do + test "distinguishes exact issuer and subject pairs" do assert {:ok, first} = Principal.oidc("https://one.example", "same", :user) assert {:ok, second} = Principal.oidc("https://two.example", "same", :user) + assert {:ok, pairwise} = Principal.oidc("https://one.example", "same-client-subject", :user) refute first.id == second.id + refute first.id == pairwise.id end test "keeps issuer and subject boundaries distinct" do @@ -37,14 +39,12 @@ defmodule Smolquery.Auth.PrincipalTest do assert {:ok, with_metadata} = Principal.oidc("issuer", "subject", :user, display_name: "Alice", - client_id: "web", - expires_at: 42 + client_id: "web" ) assert with_metadata.id == without_metadata.id assert with_metadata.display_name == "Alice" assert with_metadata.client_id == "web" - assert with_metadata.expires_at == 42 end test "requires exact non-empty issuer and subject strings" do @@ -56,26 +56,50 @@ defmodule Smolquery.Auth.PrincipalTest do end describe "local/4" do - test "builds local API-key and Basic principals" do - assert {:ok, api_key} = Principal.local("static:api", :api_key, :service) - assert {:ok, basic} = Principal.local("static:web", :basic, :user) - - assert api_key.authn == :api_key - assert basic.authn == :basic - assert api_key.issuer == nil - assert api_key.subject == nil - assert Principal.valid?(api_key) - assert Principal.valid?(basic) + test "derives stable authn-specific opaque IDs without retaining the source key" do + assert {:ok, first} = Principal.local("stable-source", :api_key, :service) + assert {:ok, second} = Principal.local("stable-source", :api_key, :service) + + assert first.id == second.id + assert first.id == "api_key:v1:YYMKF704-BVc8TMLshzdSjww0wilh6UT20EHWm27Q20" + assert first.authn == :api_key + assert first.issuer == nil + assert first.subject == nil + refute first.id =~ "stable-source" + refute Map.has_key?(first, :source_key) + assert Principal.well_formed?(first) + end + + test "separates local authn namespaces and OIDC" do + assert {:ok, api_key} = Principal.local("same-source", :api_key, :service) + assert {:ok, basic} = Principal.local("same-source", :basic, :service) + assert {:ok, oidc} = Principal.oidc("same-source", "same-source", :service) + + refute api_key.id == basic.id + refute api_key.id == oidc.id + refute basic.id == oidc.id end - test "does not allow local OIDC identities" do + test "rejects non-canonical or repeated local ID encodings" do + assert {:ok, principal} = Principal.local("stable-source", :api_key, :service) + prefix = "api_key:v1:" + body_size = byte_size(principal.id) - 1 + <> = principal.id + + refute Principal.well_formed?(%{principal | id: body <> "1"}) + refute Principal.well_formed?(%{principal | id: prefix <> principal.id}) + end + + test "does not allow local OIDC identities or empty source keys" do assert {:error, :oidc_requires_oidc_constructor} = Principal.local("id", :oidc, :user) + + assert {:error, {:invalid, :source_key}} = Principal.local("", :api_key, :service) end end describe "validation" do - test "rejects invalid kinds, options, metadata, and expiry" do + test "rejects invalid kinds, options, and metadata" do assert {:error, :invalid_kind} = Principal.oidc("issuer", "subject", :admin) assert {:error, {:unknown_option, :groups}} = @@ -84,21 +108,18 @@ defmodule Smolquery.Auth.PrincipalTest do assert {:error, {:invalid_option, :display_name}} = Principal.local("id", :api_key, :service, display_name: "") - assert {:error, {:invalid_option, :expires_at}} = - Principal.local("id", :api_key, :service, expires_at: -1) - assert {:error, :invalid_options} = - Principal.local("id", :api_key, :service, %{expires_at: 1}) + Principal.local("id", :api_key, :service, %{display_name: "name"}) assert {:error, :invalid_options} = - Principal.local("id", :api_key, :service, expires_at: 1, expires_at: 2) + Principal.local("id", :api_key, :service, display_name: "one", display_name: "two") end test "rejects invalid authentication values and malformed principals" do assert {:error, :invalid_authn} = Principal.local("id", :password, :service) assert {:error, :invalid_kind} = Principal.local("id", :api_key, :admin) - refute Principal.valid?(%Principal{ + refute Principal.well_formed?(%Principal{ id: "id", authn: :api_key, kind: :service, @@ -106,7 +127,7 @@ defmodule Smolquery.Auth.PrincipalTest do }) assert {:ok, oidc} = Principal.oidc("issuer", "subject", :user) - refute Principal.valid?(%{oidc | id: "oidc:v1:forged"}) + refute Principal.well_formed?(%{oidc | id: "oidc:v1:forged"}) end end end diff --git a/test/smolquery/auth_test.exs b/test/smolquery/auth_test.exs index f00eb983..a16bd494 100644 --- a/test/smolquery/auth_test.exs +++ b/test/smolquery/auth_test.exs @@ -14,6 +14,10 @@ defmodule Smolquery.AuthTest do context end + test "uses the stable direct assign key" do + assert Auth.assign_key() == :smolquery_auth_context + end + test "round-trips a context on a Plug connection" do conn = Auth.assign_context(conn(:get, "/"), context()) From 4d395eddf1f6121004aad155593aaeace708e53c Mon Sep 17 00:00:00 2001 From: David Whittington Date: Sat, 15 Aug 2026 14:22:51 +0000 Subject: [PATCH 3/8] fix(auth): require expiry for OIDC contexts Fail closed when an OIDC-derived authorization context omits assertion expiry, and strengthen the identity framing regression coverage. Refs T-229 and PL-27. --- lib/smolquery/auth/context.ex | 29 ++++++++++++++++++-------- test/smolquery/auth/context_test.exs | 20 ++++++++++++++++++ test/smolquery/auth/policy_test.exs | 13 ++++++++++++ test/smolquery/auth/principal_test.exs | 4 ++-- 4 files changed, 55 insertions(+), 11 deletions(-) diff --git a/lib/smolquery/auth/context.ex b/lib/smolquery/auth/context.ex index 39a66d05..2aa31958 100644 --- a/lib/smolquery/auth/context.ex +++ b/lib/smolquery/auth/context.ex @@ -1,6 +1,6 @@ defmodule Smolquery.Auth.Context do @moduledoc """ - The authenticated principal, capabilities, scope, and optional expiry for + The authenticated principal, capabilities, scope, and expiry contract for one request. `:single_tenant` is an explicit scope sentinel for the current deployment. @@ -8,10 +8,11 @@ defmodule Smolquery.Auth.Context do group, or other identity-provider claim; a future tenant membership lookup can replace it without changing the principal contract. - `expires_at` is an optional Unix epoch timestamp in integer seconds. A - context is active only while `now < expires_at`; the context is expired at - exactly the boundary. Authenticator clock-skew handling occurs before this - context is constructed. + `expires_at` is a Unix epoch timestamp in integer seconds. It is required + for OIDC principals and optional for local static principals. A context is + active only while `now < expires_at`; the context is expired at exactly the + boundary. Authenticator clock-skew handling occurs before this context is + constructed. Structural checks prove shape only, not authentication provenance. Only trusted authenticators and mappers may construct contexts; client input must @@ -39,8 +40,8 @@ defmodule Smolquery.Auth.Context do @doc """ Builds a context with the explicit single-tenant scope. - The optional `:expires_at` option is an integer Unix epoch timestamp in - seconds and defaults to `nil` for no expiry. + The `:expires_at` option is an integer Unix epoch timestamp in seconds. It + is required for OIDC principals and defaults to `nil` for local principals. """ @spec single_tenant(Principal.t(), capability() | [capability()] | MapSet.t()) :: {:ok, t()} | {:error, term()} @@ -51,7 +52,8 @@ defmodule Smolquery.Auth.Context do def single_tenant(%Principal{} = principal, capabilities, opts) do if Principal.well_formed?(principal) do with {:ok, capabilities} <- normalize_capabilities(capabilities), - {:ok, expires_at} <- options(opts) do + {:ok, expires_at} <- options(opts), + :ok <- validate_context_expiry(principal, expires_at) do {:ok, %__MODULE__{ principal: principal, @@ -95,7 +97,7 @@ defmodule Smolquery.Auth.Context do expires_at: expires_at }) do Principal.well_formed?(principal) and valid_capability_set?(capabilities) and - valid_expiry?(expires_at) + validate_context_expiry(principal, expires_at) == :ok end def well_formed?(_term), do: false @@ -179,6 +181,15 @@ defmodule Smolquery.Auth.Context do length(keys) != length(Enum.uniq(keys)) end + defp validate_context_expiry(%Principal{authn: :oidc}, nil), + do: {:error, :oidc_requires_expiry} + + defp validate_context_expiry(_principal, expires_at) do + if valid_expiry?(expires_at), + do: :ok, + else: {:error, {:invalid_option, :expires_at}} + end + defp valid_expiry?(nil), do: true defp valid_expiry?(value), do: is_integer(value) and value >= 0 end diff --git a/test/smolquery/auth/context_test.exs b/test/smolquery/auth/context_test.exs index 26ae0ab2..3ed96e99 100644 --- a/test/smolquery/auth/context_test.exs +++ b/test/smolquery/auth/context_test.exs @@ -9,6 +9,11 @@ defmodule Smolquery.Auth.ContextTest do principal end + defp oidc_principal do + {:ok, principal} = Principal.oidc("https://idp.example", "user-1", :user) + principal + end + test "uses an explicit single-tenant sentinel and optional expiry" do assert {:ok, context} = Context.single_tenant(principal(), [:query, :query, :ingest], expires_at: 100) @@ -31,6 +36,21 @@ defmodule Smolquery.Auth.ContextTest do refute Context.granted?(one, :ingest) end + test "requires expiry for OIDC contexts" do + assert {:error, :oidc_requires_expiry} = + Context.single_tenant(oidc_principal(), :query) + + assert {:error, :oidc_requires_expiry} = + Context.single_tenant(oidc_principal(), :query, expires_at: nil) + + assert {:ok, context} = + Context.single_tenant(oidc_principal(), :query, expires_at: 100) + + assert Context.well_formed?(context) + refute Context.well_formed?(%{context | expires_at: nil}) + refute Context.active?(%{context | expires_at: nil}, 0) + end + test "identifies known capabilities without converting input" do assert Context.capability?(:query) refute Context.capability?(:admin) diff --git a/test/smolquery/auth/policy_test.exs b/test/smolquery/auth/policy_test.exs index ae503bb0..b5dc7995 100644 --- a/test/smolquery/auth/policy_test.exs +++ b/test/smolquery/auth/policy_test.exs @@ -26,6 +26,19 @@ defmodule Smolquery.Auth.PolicyTest do Policy.authorize(context(:query, expires_at: 100), :unknown, 100) end + test "fails closed for an OIDC context without expiry" do + assert {:ok, principal} = Principal.oidc("https://idp.example", "user-1", :user) + + context = %Context{ + principal: principal, + scope: :single_tenant, + capabilities: MapSet.new([:platform_operate]) + } + + assert {:error, :unauthenticated} = + Policy.authorize(context, :platform_operate, 4_102_444_800) + end + test "reports a missing, malformed, or expired context as unauthenticated" do assert {:error, :unauthenticated} = Policy.authorize(nil, :query, 100) assert {:error, :unauthenticated} = Policy.authorize(:not_a_context, :query, 100) diff --git a/test/smolquery/auth/principal_test.exs b/test/smolquery/auth/principal_test.exs index 47e47785..f08ef8f7 100644 --- a/test/smolquery/auth/principal_test.exs +++ b/test/smolquery/auth/principal_test.exs @@ -27,8 +27,8 @@ defmodule Smolquery.Auth.PrincipalTest do end test "keeps issuer and subject boundaries distinct" do - assert {:ok, first} = Principal.oidc("a\0b", "c", :user) - assert {:ok, second} = Principal.oidc("a", "b\0c", :user) + assert {:ok, first} = Principal.oidc("ab", "c", :user) + assert {:ok, second} = Principal.oidc("a", "bc", :user) refute first.id == second.id end From 87530f97ec96bc626106f52ca1caf7a86c8ed485 Mon Sep 17 00:00:00 2001 From: David Whittington Date: Sat, 15 Aug 2026 14:47:58 +0000 Subject: [PATCH 4/8] fix(auth): make principal validation total Reject partial struct-tagged maps without raising through context attachment and authorization paths, and pin the durable Basic-auth principal ID derivation with a compatibility vector. Refs T-229 and PL-27. --- lib/smolquery/auth/principal.ex | 29 ++++++++++++++++---------- test/smolquery/auth/principal_test.exs | 5 +++++ test/smolquery/auth_test.exs | 20 +++++++++++------- 3 files changed, 36 insertions(+), 18 deletions(-) diff --git a/lib/smolquery/auth/principal.ex b/lib/smolquery/auth/principal.ex index e86bb9ce..0b6b55a5 100644 --- a/lib/smolquery/auth/principal.ex +++ b/lib/smolquery/auth/principal.ex @@ -108,27 +108,34 @@ defmodule Smolquery.Auth.Principal do module. This does not prove authentication provenance. """ @spec well_formed?(term()) :: boolean() - def well_formed?(%__MODULE__{} = principal) do - valid_string?(principal.id) and - principal.authn in @authn and - principal.kind in @kinds and - valid_identity?(principal) and - valid_optional_string?(principal.display_name) and - valid_optional_string?(principal.client_id) + def well_formed?(%__MODULE__{ + id: id, + authn: authn, + kind: kind, + issuer: issuer, + subject: subject, + display_name: display_name, + client_id: client_id + }) do + valid_string?(id) and + authn in @authn and + kind in @kinds and + valid_identity?(id, authn, issuer, subject) and + valid_optional_string?(display_name) and + valid_optional_string?(client_id) end def well_formed?(_term), do: false - defp valid_identity?(%__MODULE__{id: id, authn: :oidc, issuer: issuer, subject: subject}) do + defp valid_identity?(id, :oidc, issuer, subject) do valid_string?(issuer) and valid_string?(subject) and id == oidc_id(issuer, subject) end - defp valid_identity?(%__MODULE__{id: id, authn: authn, issuer: nil, subject: nil}) - when authn in @local_authn do + defp valid_identity?(id, authn, nil, nil) when authn in @local_authn do local_id_shape?(id, authn) end - defp valid_identity?(_principal), do: false + defp valid_identity?(_id, _authn, _issuer, _subject), do: false defp options(opts) when is_list(opts) do cond do diff --git a/test/smolquery/auth/principal_test.exs b/test/smolquery/auth/principal_test.exs index f08ef8f7..3690aabf 100644 --- a/test/smolquery/auth/principal_test.exs +++ b/test/smolquery/auth/principal_test.exs @@ -59,9 +59,11 @@ defmodule Smolquery.Auth.PrincipalTest do test "derives stable authn-specific opaque IDs without retaining the source key" do assert {:ok, first} = Principal.local("stable-source", :api_key, :service) assert {:ok, second} = Principal.local("stable-source", :api_key, :service) + assert {:ok, basic} = Principal.local("stable-source", :basic, :service) assert first.id == second.id assert first.id == "api_key:v1:YYMKF704-BVc8TMLshzdSjww0wilh6UT20EHWm27Q20" + assert basic.id == "basic:v1:WwzWc2-Vju_Dsy7TNFUrFyXBToNxLan-lXElQkwzypA" assert first.authn == :api_key assert first.issuer == nil assert first.subject == nil @@ -126,6 +128,9 @@ defmodule Smolquery.Auth.PrincipalTest do issuer: "unexpected" }) + refute Principal.well_formed?(%{__struct__: Principal}) + refute Principal.well_formed?(%{__struct__: Principal, id: "partial"}) + assert {:ok, oidc} = Principal.oidc("issuer", "subject", :user) refute Principal.well_formed?(%{oidc | id: "oidc:v1:forged"}) end diff --git a/test/smolquery/auth_test.exs b/test/smolquery/auth_test.exs index a16bd494..f8ff6758 100644 --- a/test/smolquery/auth_test.exs +++ b/test/smolquery/auth_test.exs @@ -42,14 +42,20 @@ defmodule Smolquery.AuthTest do end test "rejects malformed context internals without raising" do - forged = %{context() | capabilities: %MapSet{map: :forged}} - conn = Plug.Conn.assign(conn(:get, "/"), Auth.assign_key(), forged) - socket = Phoenix.Component.assign(%Socket{}, Auth.assign_key(), forged) + malformed_contexts = [ + %{context() | capabilities: %MapSet{map: :forged}}, + %{context() | principal: %{__struct__: Principal}} + ] - assert :error = Auth.fetch_context(conn) - assert :error = Auth.fetch_context(socket) - assert_raise ArgumentError, fn -> Auth.assign_context(conn, forged) end - assert_raise ArgumentError, fn -> Auth.assign_context(socket, forged) end + for malformed <- malformed_contexts do + conn = Plug.Conn.assign(conn(:get, "/"), Auth.assign_key(), malformed) + socket = Phoenix.Component.assign(%Socket{}, Auth.assign_key(), malformed) + + assert :error = Auth.fetch_context(conn) + assert :error = Auth.fetch_context(socket) + assert_raise ArgumentError, fn -> Auth.assign_context(conn, malformed) end + assert_raise ArgumentError, fn -> Auth.assign_context(socket, malformed) end + end end test "rejects unsupported targets" do From ba9858249e4e67e735920389cc2854b0942c4546 Mon Sep 17 00:00:00 2001 From: David Whittington Date: Sat, 15 Aug 2026 15:10:25 +0000 Subject: [PATCH 5/8] docs(auth): clarify context assign key Describe the shared Plug and LiveView storage location as a stable assign key rather than implying that it uses conn.private. Refs T-229 and PL-27. --- lib/smolquery/auth.ex | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/smolquery/auth.ex b/lib/smolquery/auth.ex index 2dd6326d..ffa120c2 100644 --- a/lib/smolquery/auth.ex +++ b/lib/smolquery/auth.ex @@ -2,7 +2,7 @@ defmodule Smolquery.Auth do @moduledoc """ Attachment seam for an authenticated `Smolquery.Auth.Context`. - The same private assign key is used by Plug connections and LiveView sockets. + The same stable assign key is used by Plug connections and LiveView sockets. Assigning a value that is not a context raises `ArgumentError`; fetching a missing or malformed assign returns `:error`. Assignment and fetching prove structure only; trusted adapters establish authentication provenance before From 1eb77559b46be1c6c90836b1b0b7c26ec7938e13 Mon Sep 17 00:00:00 2001 From: Chase Granberry Date: Sun, 16 Aug 2026 07:38:26 -0700 Subject: [PATCH 6/8] Derive compact_max_rows from the compaction engine's memory budget (T-262) (#166) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * Derive compact_max_rows from the compaction engine's memory budget (T-262) A fixed 4Mi-row cap OOMs a 1 GiB compaction engine: the sandbox measured a 4Mi-row group pinning the whole budget during the merge, with preserve_insertion_order already off. A group the engine cannot merge re-plans identically every sweep, so one group blocks its table forever. Splitting the group does not converge — the halves' outputs re-enter the next group up to the same cap. The cap now derives from the budget at 512 bytes per row, half the measured pin rate. The sandbox's 1 GiB budget yields 2Mi rows. An explicit SMOLQUERY_COMPACT_MAX_ROWS still wins, a tiny budget floors at 64Ki rows, and no resolvable budget keeps the old 4Mi default. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01GFEnyub9GiAqwDWRRtMKBs * The row cap adapts per table: halve on merge OOM, recover on success (T-262) A budget-derived cap fixes the budget dimension and not the workload one. The 512 B/row constant is calibrated on this bench's data; wide repetitive text fields pin kilobytes per row while compressing well enough to pass both static caps. No footer statistic predicts pin cost across workloads, so the cap answers the workload instead: a merge OOM halves the table's cap (floor 64Ki rows), a successful compaction doubles it back, and the override is shed at the resolved cap. Caps live in the compactor's state; a restart re-learns them on the next OOM. The sweeper contract grows an optional three-tuple so a sweep can carry state; retention's two-tuple form is unchanged. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01GFEnyub9GiAqwDWRRtMKBs * Hoist split_result out of the Sweeper quote for dialyzer Injected into each using module, the helper's unused clauses are dead code per expansion — the compactor never returns a two-tuple error, the retention sweeper never returns three-tuples — and dialyzer rightly flags patterns that can never match. As one shared function the clauses are all reachable. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01GFEnyub9GiAqwDWRRtMKBs * Alias the Sweeper in its own quote for credo Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01GFEnyub9GiAqwDWRRtMKBs --------- Co-authored-by: Claude Fable 5 --- lib/smolquery/storage_service/compactor.ex | 88 +++++++++++++++++-- lib/smolquery/storage_service/runtime.ex | 85 ++++++++++++++++-- lib/smolquery/storage_service/sweeper.ex | 19 +++- .../storage_service/compactor_test.exs | 60 +++++++++++++ .../storage_service/runtime_test.exs | 51 ++++++++++- 5 files changed, 282 insertions(+), 21 deletions(-) diff --git a/lib/smolquery/storage_service/compactor.ex b/lib/smolquery/storage_service/compactor.ex index 6c08a51f..0e3c42ae 100644 --- a/lib/smolquery/storage_service/compactor.ex +++ b/lib/smolquery/storage_service/compactor.ex @@ -69,6 +69,10 @@ defmodule Smolquery.StorageService.Compactor do cap never needs this — two files under `compact_below_bytes` always fit — but a single small file can carry more rows than the cap admits twice. + The row cap also adapts per table: a merge OOM halves the table's cap and + a success doubles it back, because no static bytes-per-row constant fits + every workload — see `adjusted_row_caps/3` (T-262). + The sizing query chunks the same way, oldest first, and reads each file's `num_rows` in the same call. Sizing stops once the undersized bytes found reach `compact_max_bytes`, the rows found reach `compact_max_rows`, or the @@ -130,7 +134,7 @@ defmodule Smolquery.StorageService.Compactor do alias Smolquery.StorageService.Runtime @enforce_keys [:runtime] - defstruct [:runtime] + defstruct [:runtime, row_caps: %{}] @group_max_staging_chunks 64 @stage_chunk_target_bytes 67_108_864 @@ -158,17 +162,83 @@ defmodule Smolquery.StorageService.Compactor do def sweep(name, timeout \\ 60_000), do: GenServer.call(Runtime.compactor(name), :sweep, timeout) defp run(state) do - with {:ok, tables} <- Catalog.tables(state.runtime.catalog) do - outcomes = Enum.map(tables, &compact_table(state.runtime, &1)) - - {:ok, - %{ - compacted: for({:ok, swap} <- outcomes, do: swap), - failed: for({:failed, failure} <- outcomes, do: failure) - }} + runtime = Runtime.with_compact_max_rows(state.runtime) + + with {:ok, tables} <- Catalog.tables(runtime.catalog) do + outcomes = Enum.map(tables, &compact_table(runtime, state.row_caps, &1)) + + report = %{ + compacted: for({:ok, swap} <- outcomes, do: swap), + failed: for({:failed, failure} <- outcomes, do: failure) + } + + row_caps = adjusted_row_caps(state.row_caps, outcomes, runtime.compact_max_rows) + + {:ok, report, %{state | row_caps: row_caps}} end end + @row_cap_floor 65_536 + + @doc """ + The per-table row caps after a sweep — the workload-adaptive half of the + row bound (T-262). + + The derived cap assumes 512 bytes of pinned memory per row, and a workload + decides its own pin rate: wide, repetitive text fields pin kilobytes per + row while compressing well enough to pass every static cap. So the caps + answer the workload instead of predicting it. A merge that OOMs halves the + table's cap, never below #{@row_cap_floor} rows; a table that compacts + doubles its cap back toward the runtime's, and sheds the override when it + gets there. Only tables that OOMed carry an entry, so the map stays empty + on healthy deployments. The caps live in the compactor's state: a restart + forgets them, and the first OOM after the restart re-learns them. + """ + @spec adjusted_row_caps(%{Catalog.table_ref() => pos_integer()}, [term()], pos_integer()) :: + %{Catalog.table_ref() => pos_integer()} + def adjusted_row_caps(row_caps, outcomes, resolved) do + Enum.reduce(outcomes, row_caps, fn + {:ok, %{table: table}}, caps -> + relax(caps, table, resolved) + + {:failed, %{table: table, reason: reason}}, caps -> + if merge_oom?(reason), do: tighten(caps, table, resolved), else: caps + + _skip, caps -> + caps + end) + end + + defp tighten(caps, table, resolved) do + attempted = caps |> Map.get(table, resolved) |> min(resolved) + Map.put(caps, table, max(div(attempted, 2), @row_cap_floor)) + end + + defp relax(caps, table, resolved) do + case caps do + %{^table => cap} when cap * 2 >= resolved -> Map.delete(caps, table) + %{^table => cap} -> Map.put(caps, table, cap * 2) + _no_override -> caps + end + end + + defp merge_oom?({:put_failed, _key, {:merge_failed, %Adbc.Error{message: message}}}), + do: message =~ "Out of Memory" + + defp merge_oom?(_reason), do: false + + defp compact_table(runtime, row_caps, table_ref) do + runtime = table_capped(runtime, row_caps, table_ref) + compact_table(runtime, table_ref) + end + + defp table_capped(runtime, row_caps, table_ref) do + cap = + row_caps |> Map.get(table_ref, runtime.compact_max_rows) |> min(runtime.compact_max_rows) + + %{runtime | compact_max_rows: cap} + end + defp compact_table(runtime, table_ref) do started_at = System.monotonic_time(:microsecond) diff --git a/lib/smolquery/storage_service/runtime.ex b/lib/smolquery/storage_service/runtime.ex index d4c3d5a5..750c993a 100644 --- a/lib/smolquery/storage_service/runtime.ex +++ b/lib/smolquery/storage_service/runtime.ex @@ -86,10 +86,20 @@ defmodule Smolquery.StorageService.Runtime do ~25M rows, and merge cost — the staging inserts and the clustered `ORDER BY` on the final `COPY` — scales with rows, so the byte-bounded group blew the merge's five-minute `COPY` budget and re-planned identically every - sweep. The default of 4Mi rows keeps a group of that shape near a minute. - Normal ~3x-compressible data hits the byte cap first, so the row cap only - bites the pathological case it exists for. Sizing already reads each - footer's `num_rows`, so the cap costs no new I/O. + sweep. Left `nil`, the cap derives from the compaction engine's memory + budget at 512 bytes per row, resolved by `with_compact_max_rows/2` at each + sweep (T-262): the sandbox measured a 4Mi-row group pinning ~1 GiB during + the merge — the whole of a 1 GiB budget — and a group the engine cannot + merge re-plans identically forever, because splitting it only produces + outputs the next group re-ingests up to the same cap. Half the observed + pin rate leaves 2x headroom. Without a resolvable budget the cap falls + back to 4Mi rows. The 512 B/row constant is a starting point, not a + promise — the compactor halves a table's cap after a merge OOM and + recovers it on success (`Smolquery.StorageService.Compactor.adjusted_row_caps/3`), + so a workload that pins more per row converges on its own cap. Normal + ~3x-compressible data hits the byte cap first, so the row cap only bites + the pathological case it exists for. Sizing already reads each footer's + `num_rows`, so the cap costs no new I/O. Compaction runs on its own engine, `compact_engine/1` (T-259). `compact_engine_memory_limit` sizes it the way `engine_memory_limit` sizes @@ -184,7 +194,7 @@ defmodule Smolquery.StorageService.Runtime do compact_below_bytes: 33_554_432, compact_min_inputs: 2, compact_max_bytes: 134_217_728, - compact_max_rows: 4_194_304, + compact_max_rows: nil, compact_engine_memory_limit: nil, merge_engine: nil, merge_inputs_per_call: 12, @@ -215,7 +225,7 @@ defmodule Smolquery.StorageService.Runtime do compact_below_bytes: pos_integer(), compact_min_inputs: pos_integer(), compact_max_bytes: pos_integer(), - compact_max_rows: pos_integer(), + compact_max_rows: pos_integer() | nil, compact_engine_memory_limit: String.t() | nil, merge_engine: atom() | nil, merge_inputs_per_call: pos_integer(), @@ -339,6 +349,63 @@ defmodule Smolquery.StorageService.Runtime do defp derived_memory_limit(nil, :none, _divisor), do: nil + @compact_bytes_per_row 512 + @compact_max_rows_fallback 4_194_304 + @compact_max_rows_floor 65_536 + + @doc """ + The runtime with `compact_max_rows` resolved to an integer. + + An explicit cap survives untouched. A `nil` cap derives from the + compaction engine's memory budget at #{@compact_bytes_per_row} bytes per + row — half the pin rate the sandbox measured, so a group at the cap keeps + 2x headroom (T-262). A budget too small for the floor still yields + #{@compact_max_rows_floor} rows, and no resolvable budget — no explicit + limit, no cgroup — falls back to #{@compact_max_rows_fallback} rows, the + fixed default this derivation replaced. + """ + @spec with_compact_max_rows(t(), {:ok, pos_integer()} | :none) :: t() + def with_compact_max_rows(runtime, cgroup \\ Smolquery.CgroupMemory.limit_bytes()) + + def with_compact_max_rows(%__MODULE__{compact_max_rows: rows} = runtime, _cgroup) + when is_integer(rows), + do: runtime + + def with_compact_max_rows(%__MODULE__{} = runtime, cgroup) do + rows = + runtime + |> compact_engine_memory_limit(cgroup) + |> derived_compact_max_rows() + + %{runtime | compact_max_rows: rows} + end + + defp derived_compact_max_rows(nil), do: @compact_max_rows_fallback + + defp derived_compact_max_rows(limit) do + case memory_limit_bytes(limit) do + nil -> @compact_max_rows_fallback + bytes -> max(div(bytes, @compact_bytes_per_row), @compact_max_rows_floor) + end + end + + defp memory_limit_bytes(limit) do + case Regex.run(~r/^(\d+)\s*(B|KB|MB|GB|TB|KiB|MiB|GiB|TiB)$/i, String.trim(limit)) do + [_, digits, unit] -> String.to_integer(digits) * unit_bytes(String.downcase(unit)) + _ -> nil + end + end + + defp unit_bytes("b"), do: 1 + defp unit_bytes("kb"), do: 1_000 + defp unit_bytes("mb"), do: 1_000_000 + defp unit_bytes("gb"), do: 1_000_000_000 + defp unit_bytes("tb"), do: 1_000_000_000_000 + defp unit_bytes("kib"), do: 1_024 + defp unit_bytes("mib"), do: 1_048_576 + defp unit_bytes("gib"), do: 1_073_741_824 + defp unit_bytes("tib"), do: 1_099_511_627_776 + @doc """ The engine a merge's calls run on. @@ -456,12 +523,14 @@ defmodule Smolquery.StorageService.Runtime do end defp validate_compact_max_rows(%__MODULE__{compact_max_rows: rows} = runtime) - when is_integer(rows) and rows > 0, + when (is_integer(rows) and rows > 0) or is_nil(rows), do: runtime defp validate_compact_max_rows(%__MODULE__{compact_max_rows: rows}) do raise ArgumentError, - "unsupported compact_max_rows: #{inspect(rows)} (expected a positive integer)" + "unsupported compact_max_rows: #{inspect(rows)} " <> + "(expected a positive integer, or nil to derive it from the " <> + "compaction engine's memory budget)" end defp validate_compact_engine_memory_limit( diff --git a/lib/smolquery/storage_service/sweeper.ex b/lib/smolquery/storage_service/sweeper.ex index f83fbf5f..30cc7009 100644 --- a/lib/smolquery/storage_service/sweeper.ex +++ b/lib/smolquery/storage_service/sweeper.ex @@ -12,9 +12,17 @@ defmodule Smolquery.StorageService.Sweeper do The using module supplies `run/1` — state in, `{:ok, report}` or `{:error, reason}` out — and defines its struct (with a `runtime` field) - before the `use`. + before the `use`. A sweeper that carries state across sweeps returns + `{:ok, report, state}` or `{:error, reason, state}` instead; the two-tuple + forms keep the state unchanged. """ + @doc false + @spec split_result(term(), state) :: {term(), state} when state: var + def split_result({:ok, report, state}, _state), do: {{:ok, report}, state} + def split_result({:error, reason, state}, _state), do: {{:error, reason}, state} + def split_result(result, state), do: {result, state} + defmacro __using__(opts) do interval = Keyword.fetch!(opts, :interval) @@ -23,6 +31,8 @@ defmodule Smolquery.StorageService.Sweeper do require Logger + alias Smolquery.StorageService.Sweeper + @impl GenServer def init(%Smolquery.StorageService.Runtime{} = runtime) do {:ok, schedule(%__MODULE__{runtime: runtime})} @@ -30,12 +40,15 @@ defmodule Smolquery.StorageService.Sweeper do @impl GenServer def handle_call(:sweep, _from, state) do - {:reply, run(state), state} + {reply, state} = Sweeper.split_result(run(state), state) + {:reply, reply, state} end @impl GenServer def handle_info(:sweep, state) do - case run(state) do + {result, state} = Sweeper.split_result(run(state), state) + + case result do {:ok, _report} -> :ok diff --git a/test/smolquery/storage_service/compactor_test.exs b/test/smolquery/storage_service/compactor_test.exs index e464a09f..502c8fd4 100644 --- a/test/smolquery/storage_service/compactor_test.exs +++ b/test/smolquery/storage_service/compactor_test.exs @@ -288,4 +288,64 @@ defmodule Smolquery.StorageService.CompactorTest do assert Compactor.sweep(context.storage) == {:ok, %{compacted: [], failed: []}} assert {:ok, [_a, _b]} = Catalog.segments(context.catalog, @table, :current) end + + describe "adjusted_row_caps/3 (T-262)" do + @resolved 4_194_304 + + defp oom_failure(table) do + {:failed, + %{ + table: table, + reason: + {:put_failed, "analytics/events/x.parquet", + {:merge_failed, %Adbc.Error{message: "Out of Memory Error: failed to pin block"}}} + }} + end + + test "a merge OOM halves the table's cap" do + caps = Compactor.adjusted_row_caps(%{}, [oom_failure(@table)], @resolved) + + assert caps == %{@table => div(@resolved, 2)} + end + + test "repeated OOMs keep halving, never below the floor" do + caps = + Enum.reduce(1..30, %{}, fn _sweep, caps -> + Compactor.adjusted_row_caps(caps, [oom_failure(@table)], @resolved) + end) + + assert caps == %{@table => 65_536} + end + + test "a success doubles the cap and sheds the override at the resolved cap" do + caps = %{@table => div(@resolved, 4)} + success = {:ok, %{table: @table, key: "k", replaced: 2, snapshot: 1}} + + doubled = Compactor.adjusted_row_caps(caps, [success], @resolved) + assert doubled == %{@table => div(@resolved, 2)} + + assert Compactor.adjusted_row_caps(doubled, [success], @resolved) == %{} + end + + test "a success without an override changes nothing" do + success = {:ok, %{table: @table, key: "k", replaced: 2, snapshot: 1}} + + assert Compactor.adjusted_row_caps(%{}, [success], @resolved) == %{} + end + + test "a failure that is not a merge OOM changes nothing" do + timeout = + {:failed, + %{ + table: @table, + reason: + {:put_failed, "analytics/events/x.parquet", + {:merge_failed, %Smolquery.Engine.CallExited{reason: :timeout}}} + }} + + sizing = {:failed, %{table: @table, reason: {:sizing_failed, %Adbc.Error{message: "x"}}}} + + assert Compactor.adjusted_row_caps(%{}, [timeout, sizing, :skip], @resolved) == %{} + end + end end diff --git a/test/smolquery/storage_service/runtime_test.exs b/test/smolquery/storage_service/runtime_test.exs index 736aae45..612d0b62 100644 --- a/test/smolquery/storage_service/runtime_test.exs +++ b/test/smolquery/storage_service/runtime_test.exs @@ -113,7 +113,7 @@ defmodule Smolquery.StorageService.RuntimeTest do end test "refuses an unusable compact_max_rows at boot, not at first sweep" do - for rows <- [0, -1, "4194304", nil] do + for rows <- [0, -1, "4194304"] do assert_raise ArgumentError, ~r/unsupported compact_max_rows/, fn -> Runtime.new(name: __MODULE__.BadRowCap, compact_max_rows: rows) end @@ -174,6 +174,55 @@ defmodule Smolquery.StorageService.RuntimeTest do end end + describe "with_compact_max_rows/2" do + test "an explicit cap survives untouched" do + runtime = Runtime.new(name: __MODULE__.ExplicitRowCap, compact_max_rows: 20) + + assert Runtime.with_compact_max_rows(runtime, {:ok, 4_294_967_296}).compact_max_rows == 20 + end + + test "derives 512 bytes per row from the cgroup-quartered budget" do + runtime = Runtime.new(name: __MODULE__.DerivedRowCap) + + assert Runtime.with_compact_max_rows(runtime, {:ok, 4_294_967_296}).compact_max_rows == + 2_097_152 + end + + test "derives from an explicit engine budget when there is no cgroup" do + runtime = + Runtime.new(name: __MODULE__.BudgetRowCap, compact_engine_memory_limit: "1GiB") + + assert Runtime.with_compact_max_rows(runtime, :none).compact_max_rows == 2_097_152 + end + + test "reads decimal units the way DuckDB does" do + runtime = + Runtime.new(name: __MODULE__.DecimalRowCap, compact_engine_memory_limit: "1GB") + + assert Runtime.with_compact_max_rows(runtime, :none).compact_max_rows == 1_953_125 + end + + test "a tiny budget still yields the floor" do + runtime = + Runtime.new(name: __MODULE__.TinyRowCap, compact_engine_memory_limit: "1MiB") + + assert Runtime.with_compact_max_rows(runtime, :none).compact_max_rows == 65_536 + end + + test "no resolvable budget falls back to the fixed default" do + runtime = Runtime.new(name: __MODULE__.FallbackRowCap) + + assert Runtime.with_compact_max_rows(runtime, :none).compact_max_rows == 4_194_304 + end + + test "a budget string DuckDB would accept but the parser does not falls back" do + runtime = + Runtime.new(name: __MODULE__.OpaqueRowCap, compact_engine_memory_limit: "1.5GB") + + assert Runtime.with_compact_max_rows(runtime, :none).compact_max_rows == 4_194_304 + end + end + describe "merge_engine/1" do test "resolves to the seal merge engine unless overridden" do runtime = Runtime.new(name: Storage) From 5ae463942ec0bcadddb95c4bac5f4afeaacad790 Mon Sep 17 00:00:00 2001 From: Chase Granberry Date: Sun, 16 Aug 2026 07:38:27 -0700 Subject: [PATCH 7/8] Bound ingest admission before the body is read (T-245) (#167) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The write path's 429s were correct and too late: buffer_full and {:overloaded, _} fire after the request body is in RAM, so nothing bounded resident bodies and memory grew with client concurrency until the kernel OOMKilled the pod — three kills in one in-region sweep at up to 128 VUs x 6.87 MiB. SmolqueryApi.Admission counts POST insert/load bodies against an in-flight byte limit between auth and parsing, while the request is still a header. The reservation is the declared content-length capped at the route's own body limit; no declared length reserves the limit outright. Over the limit answers 429 with retry-after; an idle counter always admits one request, so a limit below one body cannot brick ingest. The counter releases when the response sends, and a monitor releases on request crash. The limit derives as a quarter of the cgroup memory limit, floored at one NDJSON body; SMOLQUERY_INSERT_MAX_IN_FLIGHT_BYTES overrides it. Claude-Session: https://claude.ai/code/session_01GFEnyub9GiAqwDWRRtMKBs Co-authored-by: Claude Fable 5 --- config/runtime.exs | 8 + lib/smolquery_api/admission.ex | 179 ++++++++++++++++++ .../controllers/insert_controller.ex | 7 + lib/smolquery_api/router.ex | 11 +- lib/smolquery_api/runtime.ex | 40 +++- lib/smolquery_api/supervisor.ex | 2 +- test/smolquery_api/admission_test.exs | 159 ++++++++++++++++ test/smolquery_api/runtime_test.exs | 27 +++ 8 files changed, 426 insertions(+), 7 deletions(-) create mode 100644 lib/smolquery_api/admission.ex create mode 100644 test/smolquery_api/admission_test.exs diff --git a/config/runtime.exs b/config/runtime.exs index 7b427aa7..c76fd3c6 100644 --- a/config/runtime.exs +++ b/config/runtime.exs @@ -8,6 +8,14 @@ if api_key = System.get_env("SMOLQUERY_API_KEY") do config :smolquery, SmolqueryApi, api_key: api_key end +# T-245: the 429 must fire while the request is still a header. Unset, the +# in-flight limit derives as a quarter of the cgroup memory limit. +if bytes = System.get_env("SMOLQUERY_INSERT_MAX_IN_FLIGHT_BYTES") do + config :smolquery, SmolqueryApi, + insert_max_in_flight_bytes: + Smolquery.RuntimeConfig.positive_integer!("SMOLQUERY_INSERT_MAX_IN_FLIGHT_BYTES", bytes) +end + if internal_secret = System.get_env("SMOLQUERY_INTERNAL_SECRET") do config :smolquery, :internal_secret, internal_secret end diff --git a/lib/smolquery_api/admission.ex b/lib/smolquery_api/admission.ex new file mode 100644 index 00000000..7c13caa2 --- /dev/null +++ b/lib/smolquery_api/admission.ex @@ -0,0 +1,179 @@ +defmodule SmolqueryApi.Admission do + @moduledoc """ + Bounds the bytes of ingest bodies in flight, before any body is read (T-245). + + The write path's 429s were correct and too late: `buffer_full` and + `{:overloaded, _}` fire after the request body is in RAM, so nothing bounded + how many bodies were resident at once, and memory grew with client + concurrency until the kernel OOMKilled the pod — measured twice on the + sandbox, at 128 VUs x 6.87 MiB from an in-region load generator. The right + refusal costs a header read, not a body. + + `SmolqueryApi.Router` calls `admit_conn/1` between auth and parsing, so an + ingest body is counted before the first byte is read and an unauthenticated + request never reaches the counter. Only `POST .../insert` and `POST .../load` + are counted — every other route carries small bodies the parsers already + bound. The reservation is the request's `content-length`, capped at the + route's own body limit; a request that does not declare a length reserves + that limit outright, so a chunked body cannot slip under the counter. + + Admission over the limit answers 429 with `retry-after`, the same contract + the buffer's own refusals speak. An idle counter always admits one request, + whatever its size: the route's body cap decides what is too large (413), + and a limit misconfigured below one body must not brick ingest. + + The counter releases when the response is sent, and a monitor on the + request process releases on crash, so an abandoned request cannot leak its + reservation. The server is one process per API instance; two calls per + request is noise next to a multi-megabyte body. + + A request dispatched for an instance with no admission server passes + uncounted. In production the server starts under `SmolqueryApi.Supervisor` + ahead of the endpoint, so that window is a supervisor restart; in tests it + is the default, and only admission's own tests start the server. + """ + + use GenServer + + require Logger + + alias SmolqueryApi.Errors + alias SmolqueryApi.Runtime + + @doc """ + Starts the admission server for an API runtime. + """ + @spec start_link(Runtime.t()) :: GenServer.on_start() + def start_link(%Runtime{} = runtime) do + GenServer.start_link(__MODULE__, runtime, name: server(runtime.name)) + end + + @doc """ + The registered name of an instance's admission server. + """ + @spec server(atom()) :: atom() + def server(name), do: Module.concat(name, "Admission") + + @doc """ + Admits or refuses `conn` — the plug seam `SmolqueryApi.Router` dispatches + through. + + A refused conn is halted with a 429 and `retry-after`. An admitted ingest + conn releases its reservation when the response is sent. + """ + @spec admit_conn(Plug.Conn.t()) :: Plug.Conn.t() + def admit_conn( + %Plug.Conn{method: "POST", path_info: ["v1", "datasets", _, "tables", _, action]} = conn + ) + when action in ["insert", "load"] do + instance = conn.private.smolquery_api + + case Process.whereis(server(instance)) do + nil -> conn + server -> admit_conn(conn, server, reservation(conn, action, instance)) + end + end + + def admit_conn(%Plug.Conn{} = conn), do: conn + + defp admit_conn(conn, server, bytes) do + case GenServer.call(server, {:admit, bytes, self()}) do + {:ok, reservation} -> + Plug.Conn.register_before_send(conn, fn conn -> + release(server, reservation) + conn + end) + + {:error, :admission_full} -> + conn + |> Plug.Conn.put_resp_header("retry-after", "1") + |> Errors.send_error( + 429, + "RESOURCE_EXHAUSTED", + "too many ingest bytes in flight, retry later" + ) + |> Plug.Conn.halt() + end + end + + @doc """ + Releases an admitted reservation. + """ + @spec release(GenServer.server(), reference()) :: :ok + def release(server, reservation), do: GenServer.cast(server, {:release, reservation}) + + @doc """ + The bytes currently admitted on an instance — what tests assert on. + """ + @spec in_flight(atom()) :: non_neg_integer() + def in_flight(instance), do: GenServer.call(server(instance), :in_flight) + + defp reservation(conn, action, instance) do + ceiling = ceiling(action, instance) + + case Plug.Conn.get_req_header(conn, "content-length") do + [length | _] -> + case Integer.parse(length) do + {bytes, ""} when bytes >= 0 -> min(bytes, ceiling) + _not_a_length -> ceiling + end + + [] -> + ceiling + end + end + + defp ceiling("insert", _instance), do: SmolqueryApi.InsertController.max_ndjson_bytes() + + defp ceiling("load", instance) do + {:ok, runtime} = Runtime.fetch(instance) + runtime.load_max_bytes + end + + @impl GenServer + def init(%Runtime{} = runtime) do + limit = Runtime.insert_max_in_flight_bytes(runtime) + Logger.info("ingest admission limit=#{limit} bytes in flight") + + {:ok, %{limit: limit, in_flight: 0, reservations: %{}}} + end + + @impl GenServer + def handle_call({:admit, bytes, pid}, _from, state) do + if state.in_flight > 0 and state.in_flight + bytes > state.limit do + {:reply, {:error, :admission_full}, state} + else + reservation = Process.monitor(pid) + + {:reply, {:ok, reservation}, + %{ + state + | in_flight: state.in_flight + bytes, + reservations: Map.put(state.reservations, reservation, bytes) + }} + end + end + + def handle_call(:in_flight, _from, state), do: {:reply, state.in_flight, state} + + @impl GenServer + def handle_cast({:release, reservation}, state) do + Process.demonitor(reservation, [:flush]) + {:noreply, drop(state, reservation)} + end + + @impl GenServer + def handle_info({:DOWN, reservation, :process, _pid, _reason}, state) do + {:noreply, drop(state, reservation)} + end + + defp drop(state, reservation) do + case Map.pop(state.reservations, reservation) do + {nil, _reservations} -> + state + + {bytes, reservations} -> + %{state | in_flight: state.in_flight - bytes, reservations: reservations} + end + end +end diff --git a/lib/smolquery_api/controllers/insert_controller.ex b/lib/smolquery_api/controllers/insert_controller.ex index 4e905c78..e0ea60b8 100644 --- a/lib/smolquery_api/controllers/insert_controller.ex +++ b/lib/smolquery_api/controllers/insert_controller.ex @@ -18,6 +18,13 @@ defmodule SmolqueryApi.InsertController do @max_ndjson_bytes 8_000_000 + @doc """ + The NDJSON body ceiling — what `SmolqueryApi.Admission` reserves for an + insert that declares no `content-length`. + """ + @spec max_ndjson_bytes() :: pos_integer() + def max_ndjson_bytes, do: @max_ndjson_bytes + @doc """ Inserts the body's rows into a table. diff --git a/lib/smolquery_api/router.ex b/lib/smolquery_api/router.ex index 2146f126..c04a3a2b 100644 --- a/lib/smolquery_api/router.ex +++ b/lib/smolquery_api/router.ex @@ -7,9 +7,11 @@ defmodule SmolqueryApi.Router do client modules and `Smolquery.Catalog`, never into a service's internals, so the split-out rules hold by construction. Started by the `:api` role. - One pipeline, and its order is the contract: instance → auth → parsers. - `SmolqueryApi.Auth` runs before `SmolqueryApi.Parsers`, so an - unauthenticated body is never read; the catch-all route at the bottom runs + One pipeline, and its order is the contract: instance → auth → admission → + parsers. `SmolqueryApi.Auth` runs before `SmolqueryApi.Parsers`, so an + unauthenticated body is never read; `SmolqueryApi.Admission` sits between + them, so an ingest body is counted against the in-flight limit before its + first byte is read (T-245); the catch-all route at the bottom runs *inside* the authed pipeline, so an unauthenticated request is answered 401 whether its path exists or not — 404 never reveals which routes are real. @@ -45,6 +47,7 @@ defmodule SmolqueryApi.Router do pipeline :api do plug :put_instance plug SmolqueryApi.Auth + plug :admit_ingest plug SmolqueryApi.Parsers end @@ -77,4 +80,6 @@ defmodule SmolqueryApi.Router do _absent -> put_private(conn, :smolquery_api, SmolqueryApi) end end + + defp admit_ingest(conn, _opts), do: SmolqueryApi.Admission.admit_conn(conn) end diff --git a/lib/smolquery_api/runtime.ex b/lib/smolquery_api/runtime.ex index 52733b43..a520fec9 100644 --- a/lib/smolquery_api/runtime.ex +++ b/lib/smolquery_api/runtime.ex @@ -40,7 +40,8 @@ defmodule SmolqueryApi.Runtime do :catalog_opts, ingest_name: Smolquery.IngestService, query_name: Smolquery.QueryService, - load_max_bytes: 268_435_456 + load_max_bytes: 268_435_456, + insert_max_in_flight_bytes: nil ] @type t :: %__MODULE__{ @@ -50,7 +51,8 @@ defmodule SmolqueryApi.Runtime do catalog_opts: keyword() | nil, ingest_name: atom(), query_name: atom(), - load_max_bytes: pos_integer() + load_max_bytes: pos_integer(), + insert_max_in_flight_bytes: pos_integer() | nil } @doc """ @@ -80,9 +82,41 @@ defmodule SmolqueryApi.Runtime do catalog: catalog, catalog_opts: catalog_opts } - |> struct!(Keyword.take(config, [:ingest_name, :query_name, :load_max_bytes])) + |> struct!( + Keyword.take(config, [ + :ingest_name, + :query_name, + :load_max_bytes, + :insert_max_in_flight_bytes + ]) + ) end + @in_flight_floor 8_000_000 + @in_flight_fallback 268_435_456 + + @doc """ + The most ingest-body bytes `SmolqueryApi.Admission` admits at once (T-245). + + An explicit `insert_max_in_flight_bytes` wins. Left `nil`, the limit + derives as a quarter of the container's cgroup memory limit — in-flight + bodies are resident heap, and the write path needs the rest of the budget + for encode buffers and the accumulators — floored at one NDJSON body so a + small container still ingests. Without a cgroup limit the fallback is + #{@in_flight_fallback} bytes. + """ + @spec insert_max_in_flight_bytes(t(), {:ok, pos_integer()} | :none) :: pos_integer() + def insert_max_in_flight_bytes(runtime, cgroup \\ Smolquery.CgroupMemory.limit_bytes()) + + def insert_max_in_flight_bytes(%__MODULE__{insert_max_in_flight_bytes: bytes}, _cgroup) + when is_integer(bytes), + do: bytes + + def insert_max_in_flight_bytes(%__MODULE__{}, {:ok, bytes}), + do: max(div(bytes, 4), @in_flight_floor) + + def insert_max_in_flight_bytes(%__MODULE__{}, :none), do: @in_flight_fallback + use Smolquery.Runtime @doc """ diff --git a/lib/smolquery_api/supervisor.ex b/lib/smolquery_api/supervisor.ex index 31e3078d..bd24f9be 100644 --- a/lib/smolquery_api/supervisor.ex +++ b/lib/smolquery_api/supervisor.ex @@ -38,7 +38,7 @@ defmodule SmolqueryApi.Supervisor do children = DuckLake.children(runtime.catalog_opts, Runtime.catalog_engine(runtime.name)) ++ - [SmolqueryApi.Endpoint] + [{SmolqueryApi.Admission, runtime}, SmolqueryApi.Endpoint] Supervisor.init(children, strategy: :rest_for_one) end diff --git a/test/smolquery_api/admission_test.exs b/test/smolquery_api/admission_test.exs new file mode 100644 index 00000000..ba281c31 --- /dev/null +++ b/test/smolquery_api/admission_test.exs @@ -0,0 +1,159 @@ +defmodule SmolqueryApi.AdmissionTest do + use ExUnit.Case, async: true + + import Plug.Conn, only: [put_req_header: 3] + import Plug.Test + + alias Smolquery.BufferService + alias Smolquery.Catalog + alias Smolquery.IngestService + alias Smolquery.Schema + alias Smolquery.Test.ApiEndpoint + alias Smolquery.Test.MapCatalog + alias SmolqueryApi.Admission + alias SmolqueryApi.Runtime + + @moduletag :tmp_dir + + @key "admission-test-key" + @path "/v1/datasets/analytics/tables/events/insert" + + setup context do + buffer = :"admission_buffer_#{:erlang.unique_integer([:positive])}" + + start_supervised!( + {BufferService.Supervisor, + name: buffer, dir: Path.join(context.tmp_dir, "buffer"), flush_max_rows: 1}, + id: buffer + ) + + on_exit(fn -> BufferService.Runtime.delete(buffer) end) + + catalog = MapCatalog.new() + :ok = Catalog.create_dataset(catalog, "analytics") + + :ok = + Catalog.create_table( + catalog, + {"analytics", "events"}, + Schema.new!([{"id", :int64, nullable: false}]) + ) + + ingest = :"admission_ingest_#{:erlang.unique_integer([:positive])}" + + start_supervised!( + {IngestService.Supervisor, name: ingest, catalog: catalog, buffer_name: buffer}, + id: ingest + ) + + on_exit(fn -> IngestService.Runtime.delete(ingest) end) + + name = :"api_admission_#{:erlang.unique_integer([:positive])}" + + runtime = + Runtime.new( + name: name, + api_key: @key, + catalog: catalog, + ingest_name: ingest, + insert_max_in_flight_bytes: 100 + ) + + Runtime.put(runtime) + on_exit(fn -> Runtime.delete(name) end) + + start_supervised!({Admission, runtime}, id: {:admission, name}) + + %{name: name} + end + + defp post_ndjson(name, body) do + conn(:post, @path, body) + |> put_req_header("content-type", "application/x-ndjson") + |> put_req_header("content-length", Integer.to_string(byte_size(body))) + |> put_req_header("authorization", "Bearer #{@key}") + |> then(&ApiEndpoint.request(name, &1)) + end + + defp hold(name, bytes) do + server = Admission.server(name) + {:ok, reservation} = GenServer.call(server, {:admit, bytes, self()}) + {server, reservation} + end + + test "an admitted insert lands and releases its reservation", %{name: name} do + response = post_ndjson(name, ~s({"id": 1}\n)) + + assert response.status == 200 + assert Admission.in_flight(name) == 0 + end + + test "a full counter refuses before the body, with the 429 contract", %{name: name} do + hold(name, 95) + + response = post_ndjson(name, ~s({"id": 1}\n)) + + assert response.status == 429 + assert Plug.Conn.get_resp_header(response, "retry-after") == ["1"] + assert %{"error" => %{"status" => "RESOURCE_EXHAUSTED"}} = JSON.decode!(response.resp_body) + end + + test "a released reservation readmits the next insert", %{name: name} do + {server, reservation} = hold(name, 95) + assert post_ndjson(name, ~s({"id": 1}\n)).status == 429 + + Admission.release(server, reservation) + + assert post_ndjson(name, ~s({"id": 1}\n)).status == 200 + end + + test "an idle counter admits a body over the limit", %{name: name} do + body = ~s({"id": 1, "pad": "#{String.duplicate("x", 200)}"}\n) + assert byte_size(body) > 100 + + assert post_ndjson(name, body).status == 200 + assert Admission.in_flight(name) == 0 + end + + test "a crashed holder frees its bytes", %{name: name} do + server = Admission.server(name) + + holder = + spawn(fn -> + GenServer.call(server, {:admit, 95, self()}) + + receive do + :never -> :ok + end + end) + + await(fn -> Admission.in_flight(name) == 95 end) + Process.exit(holder, :kill) + await(fn -> Admission.in_flight(name) == 0 end) + end + + test "routes without ingest bodies pass a full counter", %{name: name} do + hold(name, 100) + + response = + conn(:get, "/v1/datasets/analytics/tables") + |> put_req_header("authorization", "Bearer #{@key}") + |> then(&ApiEndpoint.request(name, &1)) + + assert response.status == 200 + end + + test "an instance without an admission server passes uncounted", %{name: name} do + stop_supervised!({:admission, name}) + + assert post_ndjson(name, ~s({"id": 1}\n)).status == 200 + end + + defp await(check, attempts \\ 50) do + cond do + check.() -> :ok + attempts == 0 -> flunk("condition never held") + true -> Process.sleep(10) && await(check, attempts - 1) + end + end +end diff --git a/test/smolquery_api/runtime_test.exs b/test/smolquery_api/runtime_test.exs index 38560650..94d3bb02 100644 --- a/test/smolquery_api/runtime_test.exs +++ b/test/smolquery_api/runtime_test.exs @@ -40,4 +40,31 @@ defmodule SmolqueryApi.RuntimeTest do assert Runtime.fetch(:api_runtime_roundtrip) == :error end end + + describe "insert_max_in_flight_bytes/2 (T-245)" do + test "an explicit limit wins over the cgroup" do + runtime = + Runtime.new(name: :api_admission_explicit, api_key: "k", insert_max_in_flight_bytes: 42) + + assert Runtime.insert_max_in_flight_bytes(runtime, {:ok, 4_294_967_296}) == 42 + end + + test "derives a quarter of the cgroup limit" do + runtime = Runtime.new(name: :api_admission_derived, api_key: "k") + + assert Runtime.insert_max_in_flight_bytes(runtime, {:ok, 4_294_967_296}) == 1_073_741_824 + end + + test "a tiny cgroup limit still admits one NDJSON body" do + runtime = Runtime.new(name: :api_admission_floor, api_key: "k") + + assert Runtime.insert_max_in_flight_bytes(runtime, {:ok, 1_000_000}) == 8_000_000 + end + + test "without a cgroup limit the fallback applies" do + runtime = Runtime.new(name: :api_admission_fallback, api_key: "k") + + assert Runtime.insert_max_in_flight_bytes(runtime, :none) == 268_435_456 + end + end end From 10b019b99f1fb43dcd2f3f59d329e642910f69d0 Mon Sep 17 00:00:00 2001 From: Chase Granberry Date: Sun, 16 Aug 2026 07:38:27 -0700 Subject: [PATCH 8/8] Review fixes for the admission stack (T-264, T-265, T-266) (#169) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * Review fixes for the admission stack (T-264, T-265, T-266) - Compactor: a staging-phase OOM now tightens the row cap; only the store-wrapped shape matched before, so the same group re-OOMed every sweep. - Compactor: a cap raise now needs evidence — a cap-filling success or a cap-forced skip — and `patience` such sweeps; each OOM doubles the patience, so a table at its true limit probes it rarely instead of every second sweep, and a cap wedged at the floor unwedges. - Compactor: a timed-out final COPY (store-wrapped CallExited) now recycles the wedged compaction engine. - Compactor: `compact_max_rows` resolves once at start, not on every sweep. - Storage runtime: the row cap derives from the cgroup bytes directly, with no string round-trip; fractional budgets like "1.5GB" parse; an unreadable explicit budget logs a warning instead of a silent fallback. - Admission: a stray message no longer crashes the server (and, under rest_for_one, the Endpoint with it). - Admission: a server that dies between the lookup and the call passes the request uncounted instead of a 500. - API runtime: an explicit `insert_max_in_flight_bytes` is validated at boot; the derivation floor reuses `InsertController.max_ndjson_bytes/0`. - Whole-request load reservations stay by design; the moduledoc now states the tradeoff. - Docs: the `SMOLQUERY_COMPACT_MAX_ROWS` derivation, the new `SMOLQUERY_INSERT_MAX_IN_FLIGHT_BYTES` row, and the admission 429 in api.md. Refs T-264, T-265, T-266. Co-Authored-By: Claude Fable 5 * The row cap starts at the default, not at a bench-derived pin rate (T-266) - The 512 B/row figure came from a bench with very large rows. A cap derived from it starves typical workloads. - An unset `compact_max_rows` now starts at `4194304` rows. The per-table adaptation — halve on OOM, earn it back with evidence — fits the cap to each workload instead. - This removes the budget-string parser and the per-sweep cgroup reads with it; `with_compact_max_rows/1` no longer takes a cgroup argument. Co-Authored-By: Claude Fable 5 * One helper owns the 429 shed-load contract (T-265) - `Errors.send_resource_exhausted/3` pairs the 429 envelope with its retry-after. - The four hand-built copies — admission, buffer full, overload, job ceiling — now call it, so the contract cannot drift. Co-Authored-By: Claude Fable 5 --------- Co-authored-by: Claude Fable 5 --- docs/api.md | 4 +- docs/configuration.md | 3 +- lib/smolquery/storage_service/compactor.ex | 153 +++++++++++++----- lib/smolquery/storage_service/runtime.ex | 85 +++------- lib/smolquery_api/admission.ex | 32 +++- .../controllers/insert_controller.ex | 12 +- .../controllers/job_controller.ex | 4 +- lib/smolquery_api/errors.ex | 14 ++ lib/smolquery_api/runtime.ex | 17 +- .../storage_service/compactor_test.exs | 86 ++++++++-- .../storage_service/runtime_test.exs | 42 +---- test/smolquery_api/admission_test.exs | 17 ++ test/smolquery_api/errors_test.exs | 15 ++ test/smolquery_api/runtime_test.exs | 16 ++ 14 files changed, 327 insertions(+), 173 deletions(-) diff --git a/docs/api.md b/docs/api.md index e5f3900a..2eba7eb8 100644 --- a/docs/api.md +++ b/docs/api.md @@ -41,8 +41,8 @@ curl -H "$auth" -H "$json" -d '{"query": "SELECT count(*) AS n FROM analytics.ev | `POST /v1/datasets/:ds/tables` | create a table — re-creating with the same schema is a 200, with a different one a 409, never a silent no-op | | `GET /v1/datasets/:ds/tables/:t` | a table's schema, retention policy, and clustering key | | `PATCH /v1/datasets/:ds/tables/:t` | set or clear retention and/or clustering. Retention: `{"retention": {"column": "ts", "ttlMs": 2592000000}}` ages rows out of `ts` after 30 days, segment-grained and conservative (a segment is dropped only once *every* row in it has aged out); `{"retention": null}` keeps rows forever again. Clustering: `{"clustering": ["project_id", "ts"]}` sorts future writes by those columns (ClickHouse `ORDER BY` analog); `{"clustering": []}` clears it. Columns must exist on the schema (unknown names are 422). Changing clustering does not rewrite existing segments. A body carrying both fields applies them atomically — an error response means neither changed | -| `POST /v1/datasets/:ds/tables/:t/insert` | streaming insert — **`application/x-ndjson` only**, one JSON object per line, the same bytes ClickHouse takes as `JSONEachRow`; `insertId` is a query parameter. A JSON-array body is a 415: two content types were two ingest paths and the array one measured 3-4x slower with nothing announcing which you got. A 200 means the buffer service has every accepted row durable and queryable; rejected rows come back per-index in `insertErrors` (partial failure is a 200, BigQuery-style); a full or overloaded buffer is a 429 whose `retry-after` says how far behind the write path is. An optional `insertId` makes the request idempotent: retrying after a timeout or dropped response with the same id (and the same rows) cannot double-count — without one, retries are at-least-once | -| `POST /v1/datasets/:ds/tables/:t/load` | batch load — the body is the file (`application/x-ndjson`, `text/csv`, or `application/vnd.apache.parquet`), pushed through the same insert path in chunks; capped by `load_max_bytes` (413 past it), synchronous, and not atomic — a mid-load failure reports what was already durable, and unlike `/insert` it takes no `insertId`, so a retry re-inserts. Two measured caveats ([benchmarks](benchmarks.md)): the body spools to disk but the parser materializes every row, so a load peaks at **~10× the file in memory**; and the cap is in *bytes*, which at 61 columns is ~120k rows of NDJSON but ~254k of CSV. It is also **not** the fast path — concurrent `/insert` is 2.4× quicker | +| `POST /v1/datasets/:ds/tables/:t/insert` | streaming insert — **`application/x-ndjson` only**, one JSON object per line, the same bytes ClickHouse takes as `JSONEachRow`; `insertId` is a query parameter. A JSON-array body is a 415: two content types were two ingest paths and the array one measured 3-4x slower with nothing announcing which you got. A 200 means the buffer service has every accepted row durable and queryable; rejected rows come back per-index in `insertErrors` (partial failure is a 200, BigQuery-style); a full or overloaded buffer is a 429 whose `retry-after` says how far behind the write path is; a node with too many ingest-body bytes already in flight also answers 429, with `retry-after: 1`, before it reads the body (`SMOLQUERY_INSERT_MAX_IN_FLIGHT_BYTES`, T-245). An optional `insertId` makes the request idempotent: retrying after a timeout or dropped response with the same id (and the same rows) cannot double-count — without one, retries are at-least-once | +| `POST /v1/datasets/:ds/tables/:t/load` | batch load — the body is the file (`application/x-ndjson`, `text/csv`, or `application/vnd.apache.parquet`), pushed through the same insert path in chunks; capped by `load_max_bytes` (413 past it), counted whole against the same in-flight admission limit as `/insert` for the request's full duration, synchronous, and not atomic — a mid-load failure reports what was already durable, and unlike `/insert` it takes no `insertId`, so a retry re-inserts. Two measured caveats ([benchmarks](benchmarks.md)): the body spools to disk but the parser materializes every row, so a load peaks at **~10× the file in memory**; and the cap is in *bytes*, which at 61 columns is ~120k rows of NDJSON but ~254k of CSV. It is also **not** the fast path — concurrent `/insert` is 2.4× quicker | | `POST /v1/queries` | sync query — the finished job plus its first page of rows (`maxResults`, default 1000); a query that outlives `timeoutMs` is cancelled and answered 504 | | `POST /v1/jobs` | the same query as an async job — returns it pending | | `GET /v1/jobs/:id` | status and stats; once the result TTL expires, answered from durable job history | diff --git a/docs/configuration.md b/docs/configuration.md index 27ed9fd1..ce9ad888 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -20,6 +20,7 @@ error. | `SMOLQUERY_ROLES` | which service subtrees start — `all`, or a comma-separated subset of `api,ingest,buffer,storage,query,web` (all). An unknown name fails the boot | | `SMOLQUERY_API_KEY` | the Bearer key every `/v1` route requires; a node with the `:api` role and no key refuses to boot | | `SMOLQUERY_API_IP` / `SMOLQUERY_API_PORT` | API bind (`0.0.0.0` in the prod image / `4000`) | +| `SMOLQUERY_INSERT_MAX_IN_FLIGHT_BYTES` | the most ingest-body bytes an API node admits at once (T-245). `SmolqueryApi.Admission` counts `POST .../insert` and `.../load` bodies by declared `content-length` before any body is read, and refuses past the limit with a 429 and `retry-after: 1` — the refusal costs a header read, never a body, so a client burst sheds load instead of OOMKilling the pod. Unset, the limit is **a quarter of the container's cgroup memory limit**, floored at one NDJSON body (`8000000`); without a cgroup limit, `268435456`. An idle counter always admits one request — the route's own body cap decides what is too large | | `SMOLQUERY_WEB_IP` / `SMOLQUERY_WEB_PORT` | web UI bind — expose the listener only on purpose (`127.0.0.1` / `4002`) | | `SMOLQUERY_WEB_USERNAME` / `SMOLQUERY_WEB_PASSWORD` | the basic-auth credential every UI route requires; a node with the `:web` role and no credential refuses to boot | | `SMOLQUERY_WEB_HOST` | the public host of the UI; also the default `check_origin` source (`localhost`) | @@ -39,7 +40,7 @@ error. | `SMOLQUERY_MERGE_INPUTS_PER_CALL` | cap on `read_parquet` inputs any one of the merge's engine calls carries (`12`, T-246/T-247). Per-input cost over `httpfs` is what outruns the engine's 30 s call timeout. The merge reads a larger input list in capped chunks into a temp table, so a seal claim of any size merges. The default's derivation is in `Smolquery.StorageService.Runtime`'s docs | | `SMOLQUERY_STORAGE_MEMORY_LIMIT` | DuckDB memory limit for the storage merge engine (T-250). Unset, the limit is **half the container's cgroup memory limit**, so the merge scales with the pod; only without a cgroup limit does the engine fall back to `SMOLQUERY_MEMORY_LIMIT` — one size for every engine on every role, which is what left a 4 Gi pod merging inside 954 MiB. The resolved value and its source are logged at boot | | `SMOLQUERY_STORAGE_COMPACT_MEMORY_LIMIT` | DuckDB memory limit for the compaction engine (T-259). Compaction runs on its own engine so a timed-out merge cannot starve seals, and the compactor recycles it after a call exit. Unset, the limit is **a quarter of the container's cgroup memory limit**; only without a cgroup limit does the engine fall back to `SMOLQUERY_MEMORY_LIMIT`. The resolved value and its source are logged at boot | -| `SMOLQUERY_COMPACT_MAX_ROWS` | cap on a compaction group's summed rows (`4194304`, T-260). `compact_max_bytes` bounds compressed bytes, and on ~100x-compressible data a 47 MiB group held ~25M rows — merge cost scales with rows, so the group blew the merge's five-minute budget and re-planned identically every sweep. Sizing already reads each footer's `num_rows`, so the cap costs no new I/O. A head file no neighbor fits beside under the cap is skipped, so a row-heavy file cannot wedge the table's backlog | +| `SMOLQUERY_COMPACT_MAX_ROWS` | cap on a compaction group's summed rows (`4194304`, T-260). `compact_max_bytes` bounds compressed bytes, and on ~100x-compressible data a 47 MiB group held ~25M rows — merge cost scales with rows, so the group blew the merge's five-minute budget and re-planned identically every sweep. Sizing already reads each footer's `num_rows`, so the cap costs no new I/O. A head file no neighbor fits beside under the cap is skipped, so a row-heavy file cannot wedge the table's backlog. The default is a start, not a pin-rate prediction: the compactor adapts the cap per table — a merge OOM halves a table's cap, never below `65536` rows, and sustained evidence at the tightened cap earns it back (`Smolquery.StorageService.Compactor.adjusted_row_caps/3`, T-262) | | `SMOLQUERY_S3_BUCKET` | puts the sealed tier on an S3-compatible store: points both the storage service's and the query service's `store:` at `Segments.Store.S3` | | `SMOLQUERY_S3_ACCESS_KEY_ID` / `SMOLQUERY_S3_SECRET_ACCESS_KEY` | static S3 credentials. Set both, or neither — leaving both out uses the [AWS credential chain](#s3-credentials) instead. One without the other is rejected at startup | | `SMOLQUERY_S3_ENDPOINT` | S3-compatible endpoint (unset targets AWS S3) | diff --git a/lib/smolquery/storage_service/compactor.ex b/lib/smolquery/storage_service/compactor.ex index 0e3c42ae..49048747 100644 --- a/lib/smolquery/storage_service/compactor.ex +++ b/lib/smolquery/storage_service/compactor.ex @@ -69,9 +69,10 @@ defmodule Smolquery.StorageService.Compactor do cap never needs this — two files under `compact_below_bytes` always fit — but a single small file can carry more rows than the cap admits twice. - The row cap also adapts per table: a merge OOM halves the table's cap and - a success doubles it back, because no static bytes-per-row constant fits - every workload — see `adjusted_row_caps/3` (T-262). + The row cap also adapts per table: a merge OOM halves the table's cap, and + evidence at the tightened cap earns it back, because no static + bytes-per-row constant fits every workload — see `adjusted_row_caps/3` + (T-262). The sizing query chunks the same way, oldest first, and reads each file's `num_rows` in the same call. Sizing stops once the undersized bytes found @@ -149,6 +150,7 @@ defmodule Smolquery.StorageService.Compactor do """ @spec start_link(Runtime.t()) :: GenServer.on_start() def start_link(%Runtime{} = runtime) do + runtime = Runtime.with_compact_max_rows(runtime) GenServer.start_link(__MODULE__, runtime, name: Runtime.compactor(runtime.name)) end @@ -162,7 +164,7 @@ defmodule Smolquery.StorageService.Compactor do def sweep(name, timeout \\ 60_000), do: GenServer.call(Runtime.compactor(name), :sweep, timeout) defp run(state) do - runtime = Runtime.with_compact_max_rows(state.runtime) + runtime = state.runtime with {:ok, tables} <- Catalog.tables(runtime.catalog) do outcomes = Enum.map(tables, &compact_table(runtime, state.row_caps, &1)) @@ -179,50 +181,100 @@ defmodule Smolquery.StorageService.Compactor do end @row_cap_floor 65_536 + @relax_patience_start 2 + @relax_patience_max 64 @doc """ The per-table row caps after a sweep — the workload-adaptive half of the row bound (T-262). - The derived cap assumes 512 bytes of pinned memory per row, and a workload - decides its own pin rate: wide, repetitive text fields pin kilobytes per - row while compressing well enough to pass every static cap. So the caps - answer the workload instead of predicting it. A merge that OOMs halves the - table's cap, never below #{@row_cap_floor} rows; a table that compacts - doubles its cap back toward the runtime's, and sheds the override when it - gets there. Only tables that OOMed carry an entry, so the map stays empty - on healthy deployments. The caps live in the compactor's state: a restart - forgets them, and the first OOM after the restart re-learns them. + No static bytes-per-row constant predicts a workload's pin rate: wide, + repetitive text fields pin kilobytes per row while compressing well enough + to pass every static cap, and narrow rows pin almost nothing. So the + runtime's cap is only a start, and the caps here answer the workload + instead of predicting it. A merge that OOMs halves the + table's cap, never below #{@row_cap_floor} rows. Only tables that OOMed + carry an entry, so the map stays empty on healthy deployments. The caps + live in the compactor's state: a restart forgets them, and the first OOM + after the restart re-learns them. + + Raising a cap needs evidence, because a raise re-probes the level that + OOMed and a failed probe burns a multi-minute merge. A sweep counts toward + the raise when the table's group compacts at more than half its cap — a + two-file success proves nothing about a cap-sized group — or when the cap + itself makes every plan skip, since a cap wedged at the floor produces no + successes to learn from. The cap doubles after `patience` such sweeps, and + sheds the override when it reaches the runtime's. Each OOM doubles the + table's patience, up to #{@relax_patience_max} sweeps, so a table sitting + on its true limit probes it rarely instead of every other sweep. """ - @spec adjusted_row_caps(%{Catalog.table_ref() => pos_integer()}, [term()], pos_integer()) :: - %{Catalog.table_ref() => pos_integer()} + @spec adjusted_row_caps( + %{Catalog.table_ref() => map()}, + [term()], + pos_integer() + ) :: %{Catalog.table_ref() => map()} def adjusted_row_caps(row_caps, outcomes, resolved) do Enum.reduce(outcomes, row_caps, fn - {:ok, %{table: table}}, caps -> - relax(caps, table, resolved) + {:ok, %{table: table, rows: rows}}, caps -> + relax(caps, table, rows, resolved) + + {:skip, table}, caps -> + case caps do + %{^table => entry} -> advance(caps, table, entry, resolved) + _no_override -> caps + end {:failed, %{table: table, reason: reason}}, caps -> if merge_oom?(reason), do: tighten(caps, table, resolved), else: caps - _skip, caps -> + _not_owned, caps -> caps end) end defp tighten(caps, table, resolved) do - attempted = caps |> Map.get(table, resolved) |> min(resolved) - Map.put(caps, table, max(div(attempted, 2), @row_cap_floor)) + {attempted, patience} = + case caps do + %{^table => %{cap: cap, patience: patience}} -> + {min(cap, resolved), min(patience * 2, @relax_patience_max)} + + _no_override -> + {resolved, @relax_patience_start} + end + + Map.put(caps, table, %{ + cap: max(div(attempted, 2), @row_cap_floor), + streak: 0, + patience: patience + }) end - defp relax(caps, table, resolved) do + defp relax(caps, table, rows, resolved) do case caps do - %{^table => cap} when cap * 2 >= resolved -> Map.delete(caps, table) - %{^table => cap} -> Map.put(caps, table, cap * 2) - _no_override -> caps + %{^table => %{cap: cap} = entry} when rows * 2 >= cap -> + advance(caps, table, entry, resolved) + + _small_group_or_no_override -> + caps + end + end + + defp advance(caps, table, entry, resolved) do + cond do + entry.streak + 1 < entry.patience -> + Map.put(caps, table, %{entry | streak: entry.streak + 1}) + + entry.cap * 2 >= resolved -> + Map.delete(caps, table) + + true -> + Map.put(caps, table, %{entry | cap: entry.cap * 2, streak: 0}) end end - defp merge_oom?({:put_failed, _key, {:merge_failed, %Adbc.Error{message: message}}}), + defp merge_oom?({:put_failed, _key, reason}), do: merge_oom?(reason) + + defp merge_oom?({:merge_failed, %Adbc.Error{message: message}}), do: message =~ "Out of Memory" defp merge_oom?(_reason), do: false @@ -234,7 +286,10 @@ defmodule Smolquery.StorageService.Compactor do defp table_capped(runtime, row_caps, table_ref) do cap = - row_caps |> Map.get(table_ref, runtime.compact_max_rows) |> min(runtime.compact_max_rows) + case row_caps do + %{^table_ref => %{cap: cap}} -> min(cap, runtime.compact_max_rows) + _no_override -> runtime.compact_max_rows + end %{runtime | compact_max_rows: cap} end @@ -248,7 +303,7 @@ defmodule Smolquery.StorageService.Compactor do swap(runtime, table_ref, group, started_at) else false -> :skip - :skip -> :skip + :skip -> {:skip, table_ref} {:error, reason} -> failed(runtime, table_ref, reason, started_at) end end @@ -374,7 +429,8 @@ defmodule Smolquery.StorageService.Compactor do %{result: :ok} ) - {:ok, %{table: table_ref, key: key, replaced: length(paths), snapshot: snapshot}} + {:ok, + %{table: table_ref, key: key, replaced: length(paths), rows: row_count, snapshot: snapshot}} else {:error, reason} -> failed(runtime, table_ref, reason, started_at) end @@ -427,24 +483,39 @@ defmodule Smolquery.StorageService.Compactor do {:failed, %{table: table_ref, reason: reason}} end - defp recycle_on_exit(runtime, {step, %CallExited{}}) - when step in [:sizing_failed, :merge_failed] do - engine = Runtime.compact_engine(runtime.name) - stale_connection = Process.whereis(Engine.connection_name(engine)) + @doc """ + Whether a compaction failure carries a `Smolquery.Engine.CallExited` — what + `failed/4` recycles the compaction engine on. A final `COPY`'s exit arrives + wrapped by the store as `{:put_failed, key, {:merge_failed, exit}}`; sizing + and staging exits arrive bare. + """ + @spec engine_call_exited?(term()) :: boolean() + def engine_call_exited?({:put_failed, _key, reason}), do: engine_call_exited?(reason) + + def engine_call_exited?({step, %CallExited{}}) when step in [:sizing_failed, :merge_failed], + do: true + + def engine_call_exited?(_reason), do: false - case Process.whereis(Engine.database_name(engine)) do - nil -> - await_engine(engine, stale_connection, @engine_recycle_wait_ms) + defp recycle_on_exit(runtime, reason) do + if engine_call_exited?(reason) do + engine = Runtime.compact_engine(runtime.name) + stale_connection = Process.whereis(Engine.connection_name(engine)) - database -> - Logger.warning("recycling the compaction engine after an engine call exit") - Process.exit(database, :kill) - await_engine(engine, stale_connection, @engine_recycle_wait_ms) + case Process.whereis(Engine.database_name(engine)) do + nil -> + await_engine(engine, stale_connection, @engine_recycle_wait_ms) + + database -> + Logger.warning("recycling the compaction engine after an engine call exit") + Process.exit(database, :kill) + await_engine(engine, stale_connection, @engine_recycle_wait_ms) + end + else + :ok end end - defp recycle_on_exit(_runtime, _reason), do: :ok - defp await_engine(_engine, _stale_connection, remaining_ms) when remaining_ms <= 0, do: :ok defp await_engine(engine, stale_connection, remaining_ms) do diff --git a/lib/smolquery/storage_service/runtime.ex b/lib/smolquery/storage_service/runtime.ex index 750c993a..9b0d6291 100644 --- a/lib/smolquery/storage_service/runtime.ex +++ b/lib/smolquery/storage_service/runtime.ex @@ -86,17 +86,16 @@ defmodule Smolquery.StorageService.Runtime do ~25M rows, and merge cost — the staging inserts and the clustered `ORDER BY` on the final `COPY` — scales with rows, so the byte-bounded group blew the merge's five-minute `COPY` budget and re-planned identically every - sweep. Left `nil`, the cap derives from the compaction engine's memory - budget at 512 bytes per row, resolved by `with_compact_max_rows/2` at each - sweep (T-262): the sandbox measured a 4Mi-row group pinning ~1 GiB during - the merge — the whole of a 1 GiB budget — and a group the engine cannot - merge re-plans identically forever, because splitting it only produces - outputs the next group re-ingests up to the same cap. Half the observed - pin rate leaves 2x headroom. Without a resolvable budget the cap falls - back to 4Mi rows. The 512 B/row constant is a starting point, not a - promise — the compactor halves a table's cap after a merge OOM and - recovers it on success (`Smolquery.StorageService.Compactor.adjusted_row_caps/3`), - so a workload that pins more per row converges on its own cap. Normal + sweep. Left `nil`, the cap starts at 4Mi rows, resolved by + `with_compact_max_rows/1` once when the compactor starts. The start + assumes no bytes-per-row pin rate on purpose: the sandbox did measure a + 4Mi-row group pinning ~1 GiB during a merge — the whole of a 1 GiB budget + — but those were very large rows, and a cap derived from their pin rate + would starve typical workloads. The compactor fits the cap to the + workload instead: a merge OOM halves a table's cap and sustained evidence + earns it back (`Smolquery.StorageService.Compactor.adjusted_row_caps/3`, + T-262), so a workload that pins more per row converges on its own cap + after one bad merge, and every other workload keeps the full cap. Normal ~3x-compressible data hits the byte cap first, so the row cap only bites the pathological case it exists for. Sizing already reads each footer's `num_rows`, so the cap costs no new I/O. @@ -349,62 +348,27 @@ defmodule Smolquery.StorageService.Runtime do defp derived_memory_limit(nil, :none, _divisor), do: nil - @compact_bytes_per_row 512 - @compact_max_rows_fallback 4_194_304 - @compact_max_rows_floor 65_536 + @compact_max_rows_default 4_194_304 @doc """ The runtime with `compact_max_rows` resolved to an integer. - An explicit cap survives untouched. A `nil` cap derives from the - compaction engine's memory budget at #{@compact_bytes_per_row} bytes per - row — half the pin rate the sandbox measured, so a group at the cap keeps - 2x headroom (T-262). A budget too small for the floor still yields - #{@compact_max_rows_floor} rows, and no resolvable budget — no explicit - limit, no cgroup — falls back to #{@compact_max_rows_fallback} rows, the - fixed default this derivation replaced. + An explicit cap survives untouched. A `nil` cap starts at + #{@compact_max_rows_default} rows. The start deliberately assumes no + bytes-per-row pin rate: the one figure the sandbox measured came from a + bench with very large rows, and deriving every deployment's cap from it + would starve typical workloads. The compactor owns fitting the cap to the + workload instead — a merge OOM halves a table's cap and sustained evidence + earns it back (`Smolquery.StorageService.Compactor.adjusted_row_caps/3`, + T-262). """ - @spec with_compact_max_rows(t(), {:ok, pos_integer()} | :none) :: t() - def with_compact_max_rows(runtime, cgroup \\ Smolquery.CgroupMemory.limit_bytes()) - - def with_compact_max_rows(%__MODULE__{compact_max_rows: rows} = runtime, _cgroup) + @spec with_compact_max_rows(t()) :: t() + def with_compact_max_rows(%__MODULE__{compact_max_rows: rows} = runtime) when is_integer(rows), do: runtime - def with_compact_max_rows(%__MODULE__{} = runtime, cgroup) do - rows = - runtime - |> compact_engine_memory_limit(cgroup) - |> derived_compact_max_rows() - - %{runtime | compact_max_rows: rows} - end - - defp derived_compact_max_rows(nil), do: @compact_max_rows_fallback - - defp derived_compact_max_rows(limit) do - case memory_limit_bytes(limit) do - nil -> @compact_max_rows_fallback - bytes -> max(div(bytes, @compact_bytes_per_row), @compact_max_rows_floor) - end - end - - defp memory_limit_bytes(limit) do - case Regex.run(~r/^(\d+)\s*(B|KB|MB|GB|TB|KiB|MiB|GiB|TiB)$/i, String.trim(limit)) do - [_, digits, unit] -> String.to_integer(digits) * unit_bytes(String.downcase(unit)) - _ -> nil - end - end - - defp unit_bytes("b"), do: 1 - defp unit_bytes("kb"), do: 1_000 - defp unit_bytes("mb"), do: 1_000_000 - defp unit_bytes("gb"), do: 1_000_000_000 - defp unit_bytes("tb"), do: 1_000_000_000_000 - defp unit_bytes("kib"), do: 1_024 - defp unit_bytes("mib"), do: 1_048_576 - defp unit_bytes("gib"), do: 1_073_741_824 - defp unit_bytes("tib"), do: 1_099_511_627_776 + def with_compact_max_rows(%__MODULE__{} = runtime), + do: %{runtime | compact_max_rows: @compact_max_rows_default} @doc """ The engine a merge's calls run on. @@ -529,8 +493,7 @@ defmodule Smolquery.StorageService.Runtime do defp validate_compact_max_rows(%__MODULE__{compact_max_rows: rows}) do raise ArgumentError, "unsupported compact_max_rows: #{inspect(rows)} " <> - "(expected a positive integer, or nil to derive it from the " <> - "compaction engine's memory budget)" + "(expected a positive integer, or nil for the default)" end defp validate_compact_engine_memory_limit( diff --git a/lib/smolquery_api/admission.ex b/lib/smolquery_api/admission.ex index 7c13caa2..5144ac98 100644 --- a/lib/smolquery_api/admission.ex +++ b/lib/smolquery_api/admission.ex @@ -27,10 +27,20 @@ defmodule SmolqueryApi.Admission do reservation. The server is one process per API instance; two calls per request is noise next to a multi-megabyte body. + A load reserves its declared size for the whole request, upload phase + included, even though the body spools to disk in small chunks and only the + parse phase materializes it. That is deliberate: the counter must cover the + request's resident peak — a load's parse holds many times the file — and a + reservation that shrank during the upload would admit inserts the parse + phase then competes with. The cost is that one large load can hold most of + a small limit for its full duration. + A request dispatched for an instance with no admission server passes uncounted. In production the server starts under `SmolqueryApi.Supervisor` ahead of the endpoint, so that window is a supervisor restart; in tests it - is the default, and only admission's own tests start the server. + is the default, and only admission's own tests start the server. The same + contract covers a server that dies between the lookup and the call: the + request passes uncounted rather than crash. """ use GenServer @@ -77,7 +87,10 @@ defmodule SmolqueryApi.Admission do def admit_conn(%Plug.Conn{} = conn), do: conn defp admit_conn(conn, server, bytes) do - case GenServer.call(server, {:admit, bytes, self()}) do + case admit(server, bytes) do + :no_server -> + conn + {:ok, reservation} -> Plug.Conn.register_before_send(conn, fn conn -> release(server, reservation) @@ -86,16 +99,17 @@ defmodule SmolqueryApi.Admission do {:error, :admission_full} -> conn - |> Plug.Conn.put_resp_header("retry-after", "1") - |> Errors.send_error( - 429, - "RESOURCE_EXHAUSTED", - "too many ingest bytes in flight, retry later" - ) + |> Errors.send_resource_exhausted(1, "too many ingest bytes in flight, retry later") |> Plug.Conn.halt() end end + defp admit(server, bytes) do + GenServer.call(server, {:admit, bytes, self()}) + catch + :exit, _reason -> :no_server + end + @doc """ Releases an admitted reservation. """ @@ -167,6 +181,8 @@ defmodule SmolqueryApi.Admission do {:noreply, drop(state, reservation)} end + def handle_info(_message, state), do: {:noreply, state} + defp drop(state, reservation) do case Map.pop(state.reservations, reservation) do {nil, _reservations} -> diff --git a/lib/smolquery_api/controllers/insert_controller.ex b/lib/smolquery_api/controllers/insert_controller.ex index e0ea60b8..e9be5fdb 100644 --- a/lib/smolquery_api/controllers/insert_controller.ex +++ b/lib/smolquery_api/controllers/insert_controller.ex @@ -124,17 +124,13 @@ defmodule SmolqueryApi.InsertController do """ @spec insert_error(Plug.Conn.t(), term()) :: Plug.Conn.t() def insert_error(conn, :buffer_full) do - conn - |> put_resp_header("retry-after", "1") - |> Errors.send_error(429, "RESOURCE_EXHAUSTED", "buffer full, retry later") + Errors.send_resource_exhausted(conn, 1, "buffer full, retry later") end def insert_error(conn, {:overloaded, predicted_ms}) do - conn - |> put_resp_header("retry-after", Integer.to_string(max(ceil(predicted_ms / 1000), 1))) - |> Errors.send_error( - 429, - "RESOURCE_EXHAUSTED", + Errors.send_resource_exhausted( + conn, + max(ceil(predicted_ms / 1000), 1), "write path overloaded, ~#{predicted_ms} ms behind; retry later" ) end diff --git a/lib/smolquery_api/controllers/job_controller.ex b/lib/smolquery_api/controllers/job_controller.ex index 89b3bcad..666578cb 100644 --- a/lib/smolquery_api/controllers/job_controller.ex +++ b/lib/smolquery_api/controllers/job_controller.ex @@ -143,9 +143,7 @@ defmodule SmolqueryApi.JobController do """ @spec query_error(Plug.Conn.t(), term()) :: Plug.Conn.t() def query_error(conn, :too_many_jobs) do - conn - |> put_resp_header("retry-after", "1") - |> Errors.send_error(429, "RESOURCE_EXHAUSTED", "too many jobs in flight, retry later") + Errors.send_resource_exhausted(conn, 1, "too many jobs in flight, retry later") end def query_error(conn, :query_service_unavailable) do diff --git a/lib/smolquery_api/errors.ex b/lib/smolquery_api/errors.ex index e723c055..af5710b0 100644 --- a/lib/smolquery_api/errors.ex +++ b/lib/smolquery_api/errors.ex @@ -23,6 +23,20 @@ defmodule SmolqueryApi.Errors do |> send_resp(code, body) end + @doc """ + Sends the shed-load refusal every over-capacity route speaks: 429, + `RESOURCE_EXHAUSTED`, and a `retry-after` of `seconds`. + + One function so the contract cannot drift between the buffer's refusals, + the job ceiling, and ingest admission. + """ + @spec send_resource_exhausted(Plug.Conn.t(), pos_integer(), String.t()) :: Plug.Conn.t() + def send_resource_exhausted(conn, seconds, message) do + conn + |> put_resp_header("retry-after", Integer.to_string(seconds)) + |> send_error(429, "RESOURCE_EXHAUSTED", message) + end + @doc """ Sends the envelope a failure reason maps to. diff --git a/lib/smolquery_api/runtime.ex b/lib/smolquery_api/runtime.ex index a520fec9..beb80d12 100644 --- a/lib/smolquery_api/runtime.ex +++ b/lib/smolquery_api/runtime.ex @@ -90,9 +90,22 @@ defmodule SmolqueryApi.Runtime do :insert_max_in_flight_bytes ]) ) + |> validate_insert_max_in_flight_bytes() + end + + defp validate_insert_max_in_flight_bytes( + %__MODULE__{insert_max_in_flight_bytes: bytes} = runtime + ) + when (is_integer(bytes) and bytes > 0) or is_nil(bytes), + do: runtime + + defp validate_insert_max_in_flight_bytes(%__MODULE__{insert_max_in_flight_bytes: bytes}) do + raise ArgumentError, + "unsupported insert_max_in_flight_bytes: #{inspect(bytes)} " <> + "(expected a positive integer, or nil to derive it from the " <> + "container's cgroup memory limit)" end - @in_flight_floor 8_000_000 @in_flight_fallback 268_435_456 @doc """ @@ -113,7 +126,7 @@ defmodule SmolqueryApi.Runtime do do: bytes def insert_max_in_flight_bytes(%__MODULE__{}, {:ok, bytes}), - do: max(div(bytes, 4), @in_flight_floor) + do: max(div(bytes, 4), SmolqueryApi.InsertController.max_ndjson_bytes()) def insert_max_in_flight_bytes(%__MODULE__{}, :none), do: @in_flight_fallback diff --git a/test/smolquery/storage_service/compactor_test.exs b/test/smolquery/storage_service/compactor_test.exs index 502c8fd4..ec1d1fc5 100644 --- a/test/smolquery/storage_service/compactor_test.exs +++ b/test/smolquery/storage_service/compactor_test.exs @@ -302,35 +302,72 @@ defmodule Smolquery.StorageService.CompactorTest do }} end + defp staging_oom_failure(table) do + {:failed, + %{ + table: table, + reason: {:merge_failed, %Adbc.Error{message: "Out of Memory Error: failed to pin block"}} + }} + end + + defp compacted(table, rows) do + {:ok, %{table: table, key: "k", replaced: 2, rows: rows, snapshot: 1}} + end + test "a merge OOM halves the table's cap" do caps = Compactor.adjusted_row_caps(%{}, [oom_failure(@table)], @resolved) - assert caps == %{@table => div(@resolved, 2)} + assert caps == %{@table => %{cap: div(@resolved, 2), streak: 0, patience: 2}} + end + + test "a staging-phase OOM tightens the cap the same way" do + caps = Compactor.adjusted_row_caps(%{}, [staging_oom_failure(@table)], @resolved) + + assert caps == %{@table => %{cap: div(@resolved, 2), streak: 0, patience: 2}} end - test "repeated OOMs keep halving, never below the floor" do + test "repeated OOMs keep halving, never below the floor, and grow the patience" do caps = Enum.reduce(1..30, %{}, fn _sweep, caps -> Compactor.adjusted_row_caps(caps, [oom_failure(@table)], @resolved) end) - assert caps == %{@table => 65_536} + assert caps == %{@table => %{cap: 65_536, streak: 0, patience: 64}} end - test "a success doubles the cap and sheds the override at the resolved cap" do - caps = %{@table => div(@resolved, 4)} - success = {:ok, %{table: @table, key: "k", replaced: 2, snapshot: 1}} + test "cap-filling successes raise the cap after the patience and shed at the resolved cap" do + caps = %{@table => %{cap: div(@resolved, 4), streak: 0, patience: 2}} + quarter = compacted(@table, div(@resolved, 4)) + half = compacted(@table, div(@resolved, 2)) + + counted = Compactor.adjusted_row_caps(caps, [quarter], @resolved) + assert counted == %{@table => %{cap: div(@resolved, 4), streak: 1, patience: 2}} + + doubled = Compactor.adjusted_row_caps(counted, [quarter], @resolved) + assert doubled == %{@table => %{cap: div(@resolved, 2), streak: 0, patience: 2}} + + counted = Compactor.adjusted_row_caps(doubled, [half], @resolved) + assert Compactor.adjusted_row_caps(counted, [half], @resolved) == %{} + end - doubled = Compactor.adjusted_row_caps(caps, [success], @resolved) - assert doubled == %{@table => div(@resolved, 2)} + test "a small group's success is not evidence for a raise" do + caps = %{@table => %{cap: div(@resolved, 4), streak: 1, patience: 2}} - assert Compactor.adjusted_row_caps(doubled, [success], @resolved) == %{} + assert Compactor.adjusted_row_caps(caps, [compacted(@table, 100)], @resolved) == caps end - test "a success without an override changes nothing" do - success = {:ok, %{table: @table, key: "k", replaced: 2, snapshot: 1}} + test "a cap that makes every plan skip still earns its raise, so the floor unwedges" do + caps = %{@table => %{cap: 65_536, streak: 0, patience: 2}} + + counted = Compactor.adjusted_row_caps(caps, [{:skip, @table}], @resolved) + doubled = Compactor.adjusted_row_caps(counted, [{:skip, @table}], @resolved) + + assert doubled == %{@table => %{cap: 131_072, streak: 0, patience: 2}} + end - assert Compactor.adjusted_row_caps(%{}, [success], @resolved) == %{} + test "a success or skip without an override changes nothing" do + assert Compactor.adjusted_row_caps(%{}, [compacted(@table, 100)], @resolved) == %{} + assert Compactor.adjusted_row_caps(%{}, [{:skip, @table}], @resolved) == %{} end test "a failure that is not a merge OOM changes nothing" do @@ -348,4 +385,29 @@ defmodule Smolquery.StorageService.CompactorTest do assert Compactor.adjusted_row_caps(%{}, [timeout, sizing, :skip], @resolved) == %{} end end + + describe "engine_call_exited?/1" do + test "matches bare sizing and merge exits" do + exit = %CallExited{reason: :timeout} + + assert Compactor.engine_call_exited?({:sizing_failed, exit}) + assert Compactor.engine_call_exited?({:merge_failed, exit}) + end + + test "matches a final COPY's exit, which the store wraps" do + wrapped = + {:put_failed, "analytics/events/x.parquet", + {:merge_failed, %CallExited{reason: :timeout}}} + + assert Compactor.engine_call_exited?(wrapped) + end + + test "does not match plain errors" do + refute Compactor.engine_call_exited?({:merge_failed, %Adbc.Error{message: "x"}}) + + refute Compactor.engine_call_exited?( + {:put_failed, "x.parquet", {:merge_failed, %Adbc.Error{message: "x"}}} + ) + end + end end diff --git a/test/smolquery/storage_service/runtime_test.exs b/test/smolquery/storage_service/runtime_test.exs index 612d0b62..9f7fc5e3 100644 --- a/test/smolquery/storage_service/runtime_test.exs +++ b/test/smolquery/storage_service/runtime_test.exs @@ -174,52 +174,24 @@ defmodule Smolquery.StorageService.RuntimeTest do end end - describe "with_compact_max_rows/2" do + describe "with_compact_max_rows/1" do test "an explicit cap survives untouched" do runtime = Runtime.new(name: __MODULE__.ExplicitRowCap, compact_max_rows: 20) - assert Runtime.with_compact_max_rows(runtime, {:ok, 4_294_967_296}).compact_max_rows == 20 + assert Runtime.with_compact_max_rows(runtime).compact_max_rows == 20 end - test "derives 512 bytes per row from the cgroup-quartered budget" do - runtime = Runtime.new(name: __MODULE__.DerivedRowCap) + test "an unset cap starts at the default, for the compactor to adapt per table" do + runtime = Runtime.new(name: __MODULE__.DefaultRowCap) - assert Runtime.with_compact_max_rows(runtime, {:ok, 4_294_967_296}).compact_max_rows == - 2_097_152 + assert Runtime.with_compact_max_rows(runtime).compact_max_rows == 4_194_304 end - test "derives from an explicit engine budget when there is no cgroup" do + test "an explicit engine budget does not change the start" do runtime = Runtime.new(name: __MODULE__.BudgetRowCap, compact_engine_memory_limit: "1GiB") - assert Runtime.with_compact_max_rows(runtime, :none).compact_max_rows == 2_097_152 - end - - test "reads decimal units the way DuckDB does" do - runtime = - Runtime.new(name: __MODULE__.DecimalRowCap, compact_engine_memory_limit: "1GB") - - assert Runtime.with_compact_max_rows(runtime, :none).compact_max_rows == 1_953_125 - end - - test "a tiny budget still yields the floor" do - runtime = - Runtime.new(name: __MODULE__.TinyRowCap, compact_engine_memory_limit: "1MiB") - - assert Runtime.with_compact_max_rows(runtime, :none).compact_max_rows == 65_536 - end - - test "no resolvable budget falls back to the fixed default" do - runtime = Runtime.new(name: __MODULE__.FallbackRowCap) - - assert Runtime.with_compact_max_rows(runtime, :none).compact_max_rows == 4_194_304 - end - - test "a budget string DuckDB would accept but the parser does not falls back" do - runtime = - Runtime.new(name: __MODULE__.OpaqueRowCap, compact_engine_memory_limit: "1.5GB") - - assert Runtime.with_compact_max_rows(runtime, :none).compact_max_rows == 4_194_304 + assert Runtime.with_compact_max_rows(runtime).compact_max_rows == 4_194_304 end end diff --git a/test/smolquery_api/admission_test.exs b/test/smolquery_api/admission_test.exs index ba281c31..88ae36ba 100644 --- a/test/smolquery_api/admission_test.exs +++ b/test/smolquery_api/admission_test.exs @@ -149,6 +149,23 @@ defmodule SmolqueryApi.AdmissionTest do assert post_ndjson(name, ~s({"id": 1}\n)).status == 200 end + test "a server that dies between the lookup and the call passes uncounted", %{name: name} do + stop_supervised!({:admission, name}) + + stub = spawn(fn -> receive do: (_call -> :ok) end) + Process.register(stub, Admission.server(name)) + + assert post_ndjson(name, ~s({"id": 1}\n)).status == 200 + end + + test "a stray message does not crash the server", %{name: name} do + server = Process.whereis(Admission.server(name)) + send(server, :stray) + + assert Admission.in_flight(name) == 0 + assert Process.alive?(server) + end + defp await(check, attempts \\ 50) do cond do check.() -> :ok diff --git a/test/smolquery_api/errors_test.exs b/test/smolquery_api/errors_test.exs index 3e7f0cb4..f51f5208 100644 --- a/test/smolquery_api/errors_test.exs +++ b/test/smolquery_api/errors_test.exs @@ -16,4 +16,19 @@ defmodule SmolqueryApi.ErrorsTest do assert {"content-type", "application/json; charset=utf-8"} in response.resp_headers end + + test "sends the shed-load refusal with its retry-after" do + response = Errors.send_resource_exhausted(conn(:post, "/"), 3, "buffer full, retry later") + + assert response.status == 429 + assert {"retry-after", "3"} in response.resp_headers + + assert response.resp_body |> JSON.decode!() == %{ + "error" => %{ + "code" => 429, + "status" => "RESOURCE_EXHAUSTED", + "message" => "buffer full, retry later" + } + } + end end diff --git a/test/smolquery_api/runtime_test.exs b/test/smolquery_api/runtime_test.exs index 94d3bb02..b1d807ed 100644 --- a/test/smolquery_api/runtime_test.exs +++ b/test/smolquery_api/runtime_test.exs @@ -66,5 +66,21 @@ defmodule SmolqueryApi.RuntimeTest do assert Runtime.insert_max_in_flight_bytes(runtime, :none) == 268_435_456 end + + test "an explicit non-positive limit refuses to boot" do + assert_raise ArgumentError, ~r/insert_max_in_flight_bytes/, fn -> + Runtime.new(name: :api_admission_zero, api_key: "k", insert_max_in_flight_bytes: 0) + end + end + + test "an explicit non-integer limit refuses to boot" do + assert_raise ArgumentError, ~r/insert_max_in_flight_bytes/, fn -> + Runtime.new( + name: :api_admission_string, + api_key: "k", + insert_max_in_flight_bytes: "256MB" + ) + end + end end end