Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions config/runtime.exs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 2 additions & 2 deletions docs/api.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down
3 changes: 2 additions & 1 deletion docs/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`) |
Expand All @@ -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) |
Expand Down
66 changes: 66 additions & 0 deletions lib/smolquery/auth.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,66 @@
defmodule Smolquery.Auth do
@moduledoc """
Attachment seam for an authenticated `Smolquery.Auth.Context`.

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
assignment, and policy authorization checks expiry and capability access.
"""

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.

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)
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.well_formed?(context), do: {:ok, context}, else: :error
end

defp fetch_assign(_assigns), do: :error

defp assign_valid_context(target, context) do
if Context.well_formed?(context) do
do_assign(target, context)
else
raise ArgumentError, "expected a well-formed 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
195 changes: 195 additions & 0 deletions lib/smolquery/auth/context.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,195 @@
defmodule Smolquery.Auth.Context do
@moduledoc """
The authenticated principal, capabilities, scope, and expiry contract 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.

`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
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, :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()),
expires_at: non_neg_integer() | nil
}

@doc """
Builds a context with the explicit single-tenant scope.

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()}
@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),
:ok <- validate_context_expiry(principal, expires_at) do
{:ok,
%__MODULE__{
principal: principal,
scope: @single_tenant,
capabilities: capabilities,
expires_at: expires_at
}}
end
else
{:error, :invalid_principal}
end
end

def single_tenant(_principal, _capabilities, _opts), do: {:error, :invalid_principal}

@doc """
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
well_formed?(context) and MapSet.member?(context.capabilities, capability)
end

def granted?(_context, _capability), do: false

@doc """
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 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
validate_context_expiry(principal, expires_at) == :ok
end

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.
"""
@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
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
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}) 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 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
Loading
Loading