Lotus provides a middleware pipeline that lets you hook into query execution, schema discovery and content change events. Middleware follows the familiar Plug pattern — each module implements init/1 and call/2.
How It Works
A middleware module looks like this:
defmodule MyApp.AuditMiddleware do
def init(opts), do: opts
def call(payload, _opts) do
# Inspect or transform the payload
{:cont, payload} # continue to next middleware
# or
{:halt, "reason"} # stop pipeline, the caller gets an error tuple
end
endEach middleware receives a payload map whose contents depend on the pipeline event (see below), and must return either {:cont, payload} to continue or {:halt, reason} to abort. A halt returns {:error, reason} to the caller, except on :before_content_change, where it returns {:error, {:halted, reason}}.
Pipeline Events
| Event | Triggered | Payload keys |
|---|---|---|
:before_query | Before sanitization, preflight and execution | :statement, :source, :context, :vars |
:before_execute | After sanitization and preflight pass, before execution — or, on a cache hit, before the stored result is returned | :statement, :relations, :origin, :assigns, :source, :context, :vars |
:after_query | After execution, before result returned to caller | :result, :statement, :relations, :origin, :assigns, :source, :context, :vars |
:after_list_schemas | After schema discovery and visibility filtering | :schemas, :source, :scope, :context |
:after_list_tables | After table discovery and visibility filtering | :tables, :source, :scope, :context |
:after_describe_table | After table schema introspection and column visibility | :columns, :table_name, :schema, :source, :scope, :context |
:after_list_relations | After relation discovery and visibility filtering | :relations, :source, :scope, :context |
:after_discover | After any discovery call, following the kind-specific :after_list_* event | :kind, :result, :source, :scope, :context |
:before_content_change | Before a query, visualization, dashboard, card, filter or filter mapping is created, updated, deleted, shared or reordered | :op, :resource, :record, :changeset, :context |
:after_content_change | After that change is written | :op, :resource, :record, :changes, :context |
:vars is the map of bound query variables by name, after defaults and caller-supplied values are merged. It is %{} for a raw statement run through Lotus.run_statement/3.
The contract
What a plug can rely on, release to release:
- Phase order is fixed.
:before_query, then sanitization, preflight and:before_execute, then execution, then:after_query. The order lives in one place,Lotus.Runner.run/4, whether or not the result cache is involved. - Every query event fires on a cache hit. See Caching.
:before_executegets the relations stored with the entry;:after_querygets the stored result. - Payload keys are additive. A release may add a key to a payload; it does not remove or rename one. Match on the keys you use, not on the whole map.
:relationsfollows one rule. A list is proven, an empty list means "touches nothing", a tuple means "unknown".:before_executeand:after_querycarry the same value.:assignsbelongs to one call.:before_executestarts with%{}, and:after_queryreceives what its plugs left, on a cache hit too. It is the only key Lotus reads back from:before_execute, and the result cache never stores it. See Handing a Decision to:after_query.- Halting is final. A halt returns
{:error, reason}to the caller and later events for that run do not fire. The exception is:before_content_change, whose halt returns{:error, {:halted, reason}}— see Refusing and Recording Content Changes. Observation is telemetry's job. A plug sees only its own event. To record every run, including refusals and cache-served reads, attach to
[:lotus, :run, :start | :stop | :exception]— seeLotus.Telemetry. For content changes, attach to[:lotus, :content, :change, :start | :stop | :exception].
The exact-count run
window: [count: :exact] runs a second statement, derived from the page statement, to compute meta.total_count. That run carries the caller's :context, :vars, :scope and read-only setting. It fires :before_execute (it touches the same tables as the page and is authorised the same way) but not :before_query (it is derived from the statement that hook already returned, so a rewriting plug would apply twice) and not :after_query (it has no result the caller reads). A halt on the count run leaves meta.total_count as nil and the page result intact.
Discovery event ordering
Discovery calls (Lotus.list_schemas/2, Lotus.list_tables/2, Lotus.describe_table/3, Lotus.list_relations/2) fire two events per call:
- The kind-specific event (
:after_list_schemas,:after_list_tables,:after_describe_table, or:after_list_relations). The payload uses a key that matches the returned value (e.g.:tables,:columns). Register this event when you want the full kind-specific payload. - The unified
:after_discoverevent. The payload is always%{kind:, source:, result:, scope:, context:}. Register this event when you want a single middleware module that handles every discovery kind by dispatching on:kind.
If any middleware in either phase halts, later middleware do not run and the caller receives {:error, reason}. The kind-specific event always runs before :after_discover; halting in the kind-specific phase short-circuits :after_discover.
The :kind value in the unified event is one of :list_schemas, :list_tables, :describe_table, or :list_relations. Pattern-match on it and mutate :result in-place:
defmodule MyApp.DiscoveryAuditMiddleware do
require Logger
def init(opts), do: opts
def call(%{kind: kind, source: source, result: result, scope: _scope, context: ctx} = payload, _opts) do
user = Map.get(ctx || %{}, :user_id, "anonymous")
Logger.info("[Lotus] discover kind=#{kind} source=#{source} user=#{user} count=#{length(result)}")
{:cont, payload}
end
endRewriting the Statement
A :before_query plug may replace the :statement in its payload, and the
statement it returns is the one Lotus executes:
defmodule MyApp.TenantScope do
def init(opts), do: opts
def call(%{statement: statement, context: %{tenant_id: id}} = payload, _opts) do
scoped = %{statement | body: "SELECT * FROM (#{statement.body}) t WHERE tenant_id = $#{length(statement.params) + 1}",
params: statement.params ++ [id]}
{:cont, %{payload | statement: scoped}}
end
def call(payload, _opts), do: {:cont, payload}
endBecause the rewritten statement is what runs, :before_query fires before
statement sanitization and preflight authorization — both apply to the final
statement, not the text the caller supplied. A plug cannot rewrite its way
onto a denied table, and cannot turn a read into a write when read_only is
in force.
Returning the payload unchanged leaves the original statement in place, so existing audit and access-control plugs need no changes.
Authorizing on the Tables a Statement Touches
:before_query runs before the statement is analysed, so at that point Lotus
does not yet know which tables it reads. :before_execute runs after
sanitization and preflight have passed and before the statement executes, and
its payload carries :relations — the {schema, table} pairs preflight proved
the statement touches, for the statement a :before_query plug rewrote:
defmodule MyApp.TableAuthz do
def init(opts), do: opts
def call(%{relations: relations, context: %{user: user}} = payload, _opts)
when is_list(relations) do
if Enum.all?(relations, &MyApp.Authz.may_read?(user, &1)) do
{:cont, payload}
else
{:halt, "not authorized for one of the tables this query reads"}
end
end
# A tuple — `{:unrestricted, reason}` or `{:skipped, reason}` — means Lotus
# could not name the tables. A plug that gates on the list refuses rather
# than read it as an empty set.
def call(_payload, _opts), do: {:halt, "cannot determine which tables this query reads"}
endHalting returns {:error, reason} to the caller and the statement never runs.
:relations is a list when preflight proved the set, and an empty list means
the statement touches no relation (SELECT 1 passes the plug above with
Enum.all? over nothing). It is {:unrestricted, reason} when the adapter
cannot name the relations a statement touches (Elasticsearch, for one) and the
host opted in via :allow_unrestricted_resources, and {:skipped, reason}
when the adapter does not preflight the statement at all — the SQL adapters
skip EXPLAIN, SHOW and PRAGMA. The two tuples are distinct so a plug can
log which case it hit, and both mean "unknown".
:after_query carries the same :relations, so a plug that shapes a result
by table — drop or mask columns of one relation, annotate another — does so
without a second analysis. Both events also carry :origin, :executed or
:cached.
The two query hooks answer different questions, and both can be registered:
| Hook | Question |
|---|---|
:before_query | May this caller run a statement, and what statement should run? |
:before_execute | Given what this statement provably touches, may it proceed? |
A row-level-security plug that rewrites the statement uses the first. An
authorization plug that gates on tables uses the second. A :before_execute
plug may not rewrite the statement: preflight has already authorized the one in
the payload, and that is the one that executes.
Handing a Decision to :after_query
A gate at :before_execute and a plug at :after_query often need the same
decision — the grants that let a user read a table, and the columns those grants
restrict. :assigns carries it from one to the other, so the decision is made
once per call and both plugs work from the same answer:
defmodule MyApp.GrantGate do
def init(opts), do: opts
def call(%{relations: relations, context: %{user: user}, assigns: assigns} = payload, _opts)
when is_list(relations) do
case MyApp.Grants.resolve(user, relations) do
{:ok, grants} -> {:cont, %{payload | assigns: Map.put(assigns, :grants, grants)}}
{:error, reason} -> {:halt, reason}
end
end
def call(_payload, _opts), do: {:halt, "cannot determine which tables this query reads"}
end
defmodule MyApp.ColumnRestrictions do
def init(opts), do: opts
def call(%{result: result, assigns: %{grants: grants}} = payload, _opts) do
{:cont, %{payload | result: MyApp.Grants.restrict_columns(result, grants)}}
end
def call(_payload, _opts), do: {:halt, "no grants were resolved for this query"}
endconfig :lotus,
middleware: %{
before_execute: [{MyApp.GrantGate, []}],
after_query: [{MyApp.ColumnRestrictions, []}]
}The rules:
:before_executestarts withassigns: %{}, and each plug sees what the plugs before it left. Put your own keys into the map rather than replace it, so plugs on the same event do not remove each other's keys.:after_queryreceives the:assignsof the payload the last:before_executeplug continued with. With no:before_executeplug, or when a plug leaves a value that is not a map, it receives%{}.:assignsis the only key Lotus reads back. A change to:statement,:relations,:originor:contextat:before_executeis ignored.- The assigns belong to the call that made them. On a cache hit the gate runs
again for the new caller, and
:after_queryreceives that caller's assigns. The result cache never stores them, so a grant resolved for one user never reaches another. - A halt at
:before_executeends the run, so:after_querydoes not fire.
The exact-count run fires :before_execute but not :after_query, so the
assigns its plugs leave are not used.
Refusing and Recording Content Changes
Every create, update and delete of a query, visualization, dashboard, dashboard
card, dashboard filter or filter mapping fires two events, and so do reordering
cards and enabling or disabling public sharing. The mutation functions take
opts with a :context, as Lotus.run_query/2 does, so a plug knows who made
the change:
case Lotus.update_query(query, attrs, context: %{user: current_user}) do
{:ok, query} -> ...
{:error, {:halted, reason}} -> ...
{:error, %Ecto.Changeset{} = changeset} -> ...
endThe existing arities still work, with a nil context.
| Key | Value |
|---|---|
:op | :create, :update, :delete, :enable_sharing or :disable_sharing |
:resource | :query, :visualization, :dashboard, :dashboard_card, :dashboard_filter or :filter_mapping |
:record | :before_content_change: the struct as stored, nil on a create. :after_content_change: the struct as written |
:changeset | :before_content_change only. The Ecto.Changeset Lotus is about to write; no changes on a delete |
:changes | :after_content_change only. Each changed field with its written value, embedded fields included; %{} on a delete |
:context | The caller's :context, or nil |
Refusing a change
A halt on :before_content_change makes the mutation function return
{:error, {:halted, reason}} and nothing is written. The :halted tag keeps a
refusal apart from the {:error, %Ecto.Changeset{}} and {:error, :not_found}
the same functions already return.
defmodule MyApp.ContentPermissions do
def init(opts), do: opts
def call(%{context: %{user: %{role: :viewer}}}, _opts), do: {:halt, :forbidden}
def call(payload, _opts), do: {:cont, payload}
endThe event fires whether or not the changeset is valid, so a refusal does not first tell the caller what is wrong with the input. A plug may not change the changeset: Lotus writes the changeset it built.
Public sharing
Lotus.enable_public_sharing/2 fires :enable_sharing and
Lotus.disable_public_sharing/2 fires :disable_sharing, both on :dashboard
with :public_token in the changes. On :enable_sharing, :record holds the
token as stored: nil for a first enable, the previous token for a rotation.
:public_token is still a field Lotus.update_dashboard/3 accepts, and setting
it there is a plain :update whose changes carry it. A plug that refuses public
links checks both:
defmodule MyApp.NoPublicLinks do
def init(opts), do: opts
def call(%{op: :enable_sharing}, _opts), do: {:halt, "public links are disabled"}
def call(%{op: :update, resource: :dashboard, changeset: %{changes: %{public_token: token}}}, _opts)
when is_binary(token),
do: {:halt, "public links are disabled"}
def call(payload, _opts), do: {:cont, payload}
endReordering cards
Lotus.reorder_dashboard_cards/3 fires one :update of :dashboard_card for
each card whose position changes, with :position in the changes. Every
:before_content_change runs before any position is written, so a halt on one
card leaves every position as it was. The positions are written in one
transaction, and :after_content_change fires for each card after it commits.
Recording a change
:after_content_change fires after a successful write, with the written record
and its changes — what a version history or an audit log records:
defmodule MyApp.ContentAudit do
def init(opts), do: opts
def call(%{op: op, resource: resource, record: record, context: context} = payload, _opts) do
MyApp.Audit.record(context, op, resource, record.id)
{:cont, payload}
end
endThe write has happened, so this event cannot stop it. A plug that halts, raises,
throws or exits is logged, the later plugs on the event do not run, and the
caller still receives {:ok, record}. The event fires when the write returns, or
for a reorder when its transaction commits; if the caller wraps the call in its
own transaction, that transaction has not committed yet.
A refused or failed write fires no :after_content_change, and neither does an
update whose changeset has no changes, because it writes nothing. A delete by id
that finds no record fires no event at all.
Deletes that remove other content
A delete fires one event, for the record it names. The content removed with it
fires none, so the parent's :delete is the event to gate and to record:
| Deleting | Also removes, with no event |
|---|---|
| A dashboard | Its cards, its filters and their filter mappings |
| A card | Its filter mappings |
| A filter | Its filter mappings |
| A query | Its visualizations; dashboard cards that showed it keep their place with query_id set to nil |
Observing content changes
The [:lotus, :content, :change, :start | :stop | :exception] telemetry events
bracket every content change, refused or not. :stop carries the
:after_content_change metadata, and is emitted for an update that wrote nothing
too. :exception carries the reason the caller receives: {:halted, reason} for
a refusal, the changeset for a failed validation. A consumer that records every
attempt attaches to these rather than to a plug. See Lotus.Telemetry.
Configuration
Register middleware in your Lotus config. Each entry is a {module, opts} tuple — opts is passed to init/1 at compile time:
config :lotus,
middleware: %{
before_query: [
{MyApp.AccessControlMiddleware, []},
{MyApp.QueryAuditMiddleware, [repo: MyApp.AuditRepo]}
],
before_execute: [
{MyApp.TableAuthz, []}
],
after_query: [
{MyApp.ResultRedactionMiddleware, [fields: ~w(email phone ssn)]}
],
after_list_tables: [
{MyApp.TableFilterMiddleware, []}
]
}Middleware runs in the order listed. Multiple middleware can be chained on the same event.
The config is compiled once — init/1 runs at compile time and the result is
stored in :persistent_term. Reloading the config recompiles the pipeline, and
an empty (or absent) :middleware config clears it: a reload that no longer
declares middleware actually turns it off rather than leaving the previous
pipeline in place.
Context and Scope
Two separate options reach middleware, and they are not interchangeable.
:context is opaque caller data (e.g. the current user, a request id). Lotus
never inspects it and it never affects execution. It is present on every event
payload.
:scope identifies who is asking. Lotus hashes it into cache keys and passes
it to the visibility resolver, so a resolver can hide tables or mask columns per
tenant or per role. It is present on the discovery event payloads
(:after_list_*, :after_discover), not on the query events.
# Pass both when running a query
Lotus.run_statement("SELECT * FROM orders", [],
context: %{user_id: current_user.id},
scope: %{tenant_id: current_user.tenant_id}
)The result cache key includes
:scopebut never:context. Middleware runs outside the cache callback, so a plug that masks or filters results per actor works from:contextalone — it sees every call, hit or miss. The visibility resolver does not: it runs inside the callback, so a resolver that hides tables or masks columns per actor needs the caller to pass a:scopeidentifying that actor, or one actor's stored result is served to the next. See Caching.Lotus.invalidate_scope/1clears both the discovery and result cache entries for a given scope.
defmodule MyApp.AccessControlMiddleware do
def init(opts), do: opts
def call(%{context: %{user_id: nil}} = _payload, _opts) do
{:halt, "authentication required"}
end
def call(payload, _opts) do
{:cont, payload}
end
endExamples
Audit Logging
Log each query execution with the user who ran it. :before_query runs outside
the result cache, so this records every call, cached or not — see
Caching below.
defmodule MyApp.QueryAuditMiddleware do
require Logger
def init(opts), do: opts
def call(%{statement: statement, source: source, context: context} = payload, _opts) do
user_id = Map.get(context || %{}, :user_id, "anonymous")
Logger.info("[Lotus] user=#{user_id} source=#{source} body=#{inspect(statement.body)}")
{:cont, payload}
end
endstatement.body is adapter-opaque: SQL text for an Ecto-backed source, a JSON
map or DSL term for others. Use inspect/1 rather than string interpolation so
the plug works for every source type.
Row-Level Security
Block queries that don't include a tenant filter:
defmodule MyApp.TenantMiddleware do
def init(opts), do: opts
def call(%{statement: statement, context: context} = payload, _opts) do
tenant_id = Map.get(context || %{}, :tenant_id)
cond do
is_nil(tenant_id) ->
{:halt, "tenant context required"}
not String.contains?(String.downcase(statement.body), "tenant_id") ->
{:halt, "queries must filter by tenant_id"}
true ->
{:cont, payload}
end
end
endNote: The
String.contains?check above is intentionally simplified for illustration. It can be bypassed (e.g. via SQL comments). For real row-level security, inject a parameterized filter using the:filtersoption onLotus.run_query/2instead of inspecting raw SQL text.
Limiting Variable Values
Reject a query when the caller picks a date range that is too wide. The plug reads the bound variables from :vars, so it works on every path that runs a saved query: the editor, dashboards, exports and the AI assistant.
config :lotus,
middleware: %{before_query: [{MyApp.DateRangeLimit, max_days: 5}]}
defmodule MyApp.DateRangeLimit do
def init(opts), do: opts
def call(%{vars: %{"start_date" => from, "end_date" => to}} = payload, opts) do
with {:ok, from} <- Date.from_iso8601(to_string(from)),
{:ok, to} <- Date.from_iso8601(to_string(to)),
true <- Date.diff(to, from) <= opts[:max_days] do
{:cont, payload}
else
_ -> {:halt, "Date range must be #{opts[:max_days]} days or less"}
end
end
# Queries without those variables are not affected.
def call(payload, _opts), do: {:cont, payload}
endUse payload.context in the same plug for per-user exceptions, or payload.source for per-source limits.
Redacting Sensitive Data in Results
Mask PII columns (emails, phone numbers, etc.) so non-admin users only see partial values.
Lotus.Visibility.Mask.apply/2 applies the same strategies as
column visibility policies, so the plug does
not need its own masking code:
defmodule MyApp.ResultRedactionMiddleware do
@moduledoc """
Masks sensitive columns in query results with a mask strategy per column name.
Admins (identified via context) see full values; everyone else sees masked output.
"""
alias Lotus.Visibility.Mask
def init(opts), do: Keyword.get(opts, :fields, %{})
def call(%{result: result, context: context} = payload, fields) do
if admin?(context) do
{:cont, payload}
else
strategies = Enum.map(result.columns, &Map.get(fields, &1))
redacted_rows =
Enum.map(result.rows, fn row ->
row
|> Enum.zip(strategies)
|> Enum.map(fn
{val, nil} -> val
{val, strategy} -> Mask.apply(val, strategy)
end)
end)
{:cont, put_in(payload, [:result, Access.key(:rows)], redacted_rows)}
end
end
defp admin?(%{role: :admin}), do: true
defp admin?(_), do: false
endConfigure which fields to redact and how:
config :lotus,
middleware: %{
after_query: [
{MyApp.ResultRedactionMiddleware,
[
fields: %{
"email" => {:partial, keep_domain: true, keep_first: 1, keep_last: 0},
"phone" => {:partial, keep_last: 4},
"ssn" => :sha256
}
]}
]
}A query like SELECT name, email FROM users would return:
| name | |
|---|---|
| Alice Johnson | a**@example.com |
| Bob Smith | b**@example.com |
Filtering Schema Discovery
Hide internal tables from the schema browser:
defmodule MyApp.TableFilterMiddleware do
@hidden_prefixes ["_internal_", "oban_"]
def init(opts), do: opts
def call(%{tables: tables} = payload, _opts) do
filtered = Enum.reject(tables, fn table ->
Enum.any?(@hidden_prefixes, &String.starts_with?(table.name, &1))
end)
{:cont, %{payload | tables: filtered}}
end
endCaching
Query middleware and discovery middleware both sit outside their caches. Only the raw work an adapter does is stored.
Every query event fires on a cache hit
:before_query, :before_execute and :after_query all run on a call served
from the cache:
- Side-effecting plugs see every call. An audit plug records hits and misses alike, and a plug that halts for one caller halts whether or not the cache is warm.
- Context-sensitive plugs are safe.
:contextis not part of the cache key, and for middleware it does not need to be: two callers with different:contextvalues each get their own middleware decision on the same stored rows. :before_executegets its relations either way. The relations preflight named are stored with the result, so the event carries them on a hit as well as on a miss — a plug that authorizes a statement against its tables is an access control, and one a warm cache skips would be no control at all. The payload says which case it is:origin: :cachedon a hit,:executedon a miss.:assignsstays with the call. On a miss, sanitization, preflight and:before_executerun before the entry is written, and the entry holds only the result and the relations. On a hit the gate runs again for this caller. Either way,:after_queryreceives the assigns of this call.- The stored entry keeps the raw result.
:after_queryruns on the way out, so what a plug makes of the result is returned to that caller and never written back. A halt there withholds the result from that caller; the rows themselves stay in the cache for a caller whose own plugs let them through.
A :before_query plug that rewrites the statement keys its own cache entry.
:before_query runs before pagination, so the plug is handed the query the
caller wrote rather than a LIMIT wrapper around it, and everything that keys
the entry — the body, the bound values, the window — describes the statement the
plug returned. A plug that filters through a bound parameter rather than through
the statement text is keyed just the same.
Pass cache: :bypass on a call that must never be served from the cache, or
cache: :refresh to re-run and re-seed it.
What a cache hit does skip
Everything between sanitization and the returned rows is what the entry stores,
so a hit skips it: the adapter's statement sanitization, preflight table
authorization, and column visibility — mask, omit, and the hidden-column
error. Those controls read :scope, and :scope is part of the cache key, so
two scopes never share an entry.
:scope is the only caller identity the result cache separates on. A
visibility resolver must therefore decide from (source, relations, column, scope) alone. A resolver that varies its rules by anything else — a user read
out of :context, process state, or the current request — masks the warming
caller's result and then serves it unchanged to the next caller. Carry the actor
in :scope when visibility depends on it, and keep per-actor logic that cannot
be expressed that way in middleware, which runs on every call.
Discovery middleware runs outside the schema cache
Discovery middleware (:after_list_*, :after_discover) runs outside the schema cache callback. The adapter result with visibility filtering applied is cached; middleware re-runs on every call against that cached result.
- Context-sensitive middleware is safe. Two callers with different
:contextvalues receive results filtered by their own middleware logic, not each other's cached output. - Middleware runs on every call, not only on cache misses. Side-effecting middleware (e.g. audit logging) should budget accordingly.
- Adapter calls are still cached. The schema cache short-circuits the underlying
Adapter.list_tables/3(etc.) on repeat calls — only the middleware pipeline re-runs.
Halting the Pipeline
When a middleware returns {:halt, reason}, the pipeline stops immediately and Lotus returns {:error, reason} to the caller — or {:error, {:halted, reason}} on :before_content_change, and nothing at all on :after_content_change, which fires after the write. This is useful for enforcing access control, rate limiting, or any validation that should prevent execution:
def call(%{statement: %{body: body}} = payload, _opts) when is_binary(body) do
if String.contains?(String.downcase(body), "pg_sleep") do
{:halt, "pg_sleep is not allowed"}
else
{:cont, payload}
end
end
def call(payload, _opts), do: {:cont, payload}