In-memory InfluxDB client for fast, isolated testing.
Stores data in ETS tables, enabling safe async: true tests with full
isolation between test instances. Each call to start/1 creates an
independent ETS table.
Parses real line protocol on write, stores points as maps, and responds
with realistic InfluxDB response formats on query. Parsing is split out:
InfluxElixir.Client.Local.LineProtocolParser handles writes and
InfluxElixir.Client.Local.SQLParser handles the SQL subset; this module
owns storage, profiles and the InfluxQL and Flux paths; SQL execution is
InfluxElixir.Client.Local.SQLExecutor.
Profiles
LocalClient enforces an InfluxDB version profile that determines which operations are available. This prevents tests from accidentally using operations that the real InfluxDB backend doesn't support.
| Profile | Write | SQL | InfluxQL | Flux | DB CRUD | Bucket CRUD | Tokens |
|---|---|---|---|---|---|---|---|
:v3_core | yes | yes | yes | no | yes | no | no |
:v3_enterprise | yes | yes | yes | no | yes | no | yes |
:v2 | yes | no | no | yes | no | yes | no |
Operations outside the configured profile return
{:error, :unsupported_operation}.
Usage
# Match your production InfluxDB version
setup do
{:ok, conn} = InfluxElixir.Client.Local.start(
databases: ["test_db"],
profile: :v3_core
)
on_exit(fn -> InfluxElixir.Client.Local.stop(conn) end)
{:ok, conn: conn}
endChecking Profile Support
Use supports?/2 to check if an operation is available:
if Local.supports?(conn, :query_sql) do
Local.query_sql(conn, "SELECT * FROM cpu", database: "test_db")
endETS Key Layout
One :ordered_set per instance. Every mutation is a single insert or
delete of its own key, so concurrent writers — async: true tests sharing
one database, BatchWriter flushes racing direct writes — never
read-modify-write a shared value and no write is ever lost:
{:database, name}=>true{:bucket, name}=>true{:token, id}=>map()— the token map{:point, database, measurement, seq}=>point_map()—seqis a monotonic integer, so points scan in insertion order{:column, database, measurement, column}=> the column's kind (iox::column_type::tagoriox::column_type::field::<type>), fixed by the first write that names the column
Write Rules
What a write accepts is what InfluxDB 3 accepts, verified against the engine:
- A payload is applied line by line. A line with a syntax error, or a
column whose kind conflicts with the measurement's schema, is dropped
and reported; the other lines are stored. The result is then
{:error, %{status: 400, body: json}}with the engine's body —"partial write of line protocol occurred"and onedataentry per bad line (error_message,line_number,original_line). - A column's kind is fixed by the first write that names it, per
database and measurement: a tag stays a tag, an integer field stays
an integer (
v=1ithenv=2.0is "invalid column type for column 'v', expected iox::column_type::field::integer, got iox::column_type::field::float"). Deleting the database drops the schema with the data. timeis a reserved column; a key cannot be both a tag and a field on one line; an integer must fit in 64 bits (7uis unsigned); a newline inside a quoted string value is part of the value; an empty payload is "incoming write was empty".
SQL Query Support
query_sql/3 understands a subset of SQL:
SELECT * FROM measurementSELECT col1, col2 [, ...] FROM measurementwith optionalAS alias(projects fields and tags;timeis selectable).timeandDATE_BINbuckets areDateTimevalues with microsecond precision, the same as the HTTP and Flight transports return; compare them withDateTime.compare/2or a six-digit sigil (~U[... .000000Z]). A projected column may be an arithmetic expression with an alias ((bid + ask) / 2 AS mid); a null operand makes the column null (omitted).ORDER BYmay name a projected alias.WITH name AS (<select>)[, name AS (<select>)] <select>— non-recursive CTEs. Each body is a query in this subset, run in order over the store or an earlier CTE; the finalSELECTmay read from any of them (FROM w). A CTE's output columns are its fields (timestaystime).FROM a CROSS JOIN b— every row ofapaired with every row ofb(the usual use is broadcasting a one-row CTE such as a median across the rows it screens). A column present on both sides is refused as ambiguous, because qualifiers are dropped and the two could not be told apart; the engine refuses the unqualified reference too. Other joins, set operations,HAVINGand window functions are rejected by name rather than silently ignored.- Table qualifiers and aliases:
FROM q AS w/FROM q w, andw.time,q.bidin any clause — one table per query, so the prefix is dropped. WHEREwith=,!=/<>,<,<=,>,>=, combined withAND,OR,NOTand parentheses (ANDbinds tighter thanOR, as in SQL). A quoted literal is always a string, exactly as in InfluxDB v3:'08338636'keeps its leading zero and matches a string tag, and comparing it against a numeric field compares the field's text rendering (soamount >= '1000.00'is a lexical comparison — DataFusion casts the numeric side to Utf8). The other way round, a string column against a bare number compares the number's text rendering, also lexically (rack = 2matches the tag"2";rack > 3does not match"10"). Bare literals (42,1.5,true) are typed and compare numerically against numeric fields. Either side may be an arithmetic expression over columns (price <= med * 3,2 * price > volume); a bare word is a column reference, as in SQL. A column that no row has — named anywhere:SELECT, an aggregate,WHERE,GROUP BY,ORDER BY,DISTINCT— is the engine's schema error ("No field named prod", HTTP 500), which is what a typo or a forgotten pair of quotes produces in production. With no rows the schema is unknown and nothing is checked.col = NULL(anilparam) is never true.WHERE col IN (v1, v2, ...)andWHERE col NOT IN (v1, v2, ...)— each item a literal, a column or an expression, as in SQL (a bare word is a column reference, never a string)- A constant with an alias in any select list (
0.0 AS volume,'x' AS label); an unaliased constant is refused because DataFusion names it after its own rendering WHERE col IS NULLandWHERE col IS NOT NULLWHERE col [NOT] BETWEEN low AND high(inclusive;timetoo)WHERE col [NOT] LIKE 'pattern'andILIKE(%any run,_one character;LIKEis case-sensitive,ILIKEis not).LIKEover a numeric column is the engine's planning error, reproduced.WHERE time <op> <comparand>— exactly what InfluxDB 3 accepts against a Timestamp: a quoted ISO-8601 datetime ('2026-03-31T12:00:00Z', zone-less or fractional forms too), a quoted date ('2026-03-31', midnight UTC), ornow()offset by+/-INTERVAL 'N unit'terms (now() - INTERVAL '5 minutes'). A bare integer (time > 1700000000) and an integer-as-string are rejected, as DataFusion rejects them ("Cannot infer common argument type for comparison operation Timestamp(ns) > Int64"), rather than silently matching nothing.SELECT DISTINCT col[, col ...] FROM measurement(sorted combinations;ORDER BYmust name a selected column, as in DataFusion)ORDER BY a [ASC|DESC][, b [ASC|DESC] ...]— each termtime, a column, an output alias, or (on raw and projected rows) an expression such asCAST(level AS INTEGER) DESC; every term applies, each with its own directionCAST(expr AS INTEGER | INT | BIGINT | DOUBLE | FLOAT | VARCHAR | STRING)and DataFusion'scol::TYPEshorthand, wherever an expression is allowed:WHERE(CAST(level AS INTEGER) <= 20compares a numeric tag numerically),BETWEEN,LIKE, projections, aggregates, arithmetic andORDER BY. Text converts only when the whole string is a number, a float truncates to an integer, a number renders to text, null stays null. A cast that cannot be performed ('abc'toINTEGER,timetoINTEGER) makes InfluxDB 3 Core drop the connection mid-response, whichClient.HTTPreports as{:error, {:connection_error, %Mint.TransportError{reason: :closed}}}; the double reports{:error, {:connection_error, :closed}}.BOOLEANandTIMESTAMPtargets are outside the subset.LIMIT nandOFFSET m, in either order —OFFSETskips rows beforeLIMITtakes them, on plain, projected, grouped andDISTINCTrows alike;LIMIT 0returns no rows; a negative or non-numeric limit is rejected, as the engine rejects it$paramplaceholders viaparams: %{"$name" => value}in opts. ADateTime,NaiveDateTimeorDateparam renders as the ISO-8601 string Jason sends over HTTP, sotime >= $startworks the same on both clients; an integer param againsttimeis rejected on both.DATE_BIN(INTERVAL 'N unit', time)time bucketing- Aggregate functions:
AVG,SUM,COUNT,MIN,MAX,MEDIAN(the middle value; for an even count the mean of the two middle values in the column's type, so two integers average with integer division),STDDEV/STDDEV_SAMP(sample),STDDEV_POP,VAR/VAR_SAMP(sample),VAR_POP. The argument may be an arithmetic expression over fields and numeric literals (SUM(value * value),AVG(bid + ask)); two integer operands divide as integers (3 / 2 = 1), as in DataFusion. Division by zero is null in the double, where InfluxDB returns IEEE infinity for floats (serialised as JSONnullbut counted byCOUNT) and fails the query for integers. A sample statistic over one value is null.COUNT(DISTINCT col)counts distinct non-null values.MIN(time),MAX(time)andCOUNT(time)work (aDateTimeresult); every other aggregate overtime, and any arithmetic on it, is rejected as DataFusion rejects it. - Selector functions:
selector_first|last|min|max(field, time)['value']and['time'] - Ordered aggregates:
first_value(field ORDER BY col [ASC|DESC])andlast_value(field ORDER BY col [ASC|DESC])— the InfluxDB v3 SQL (DataFusion) spelling. TheORDER BYis required: without it the real engine returns an arbitrary row from the group, which the double cannot reproduce, so it rejects the query rather than certify a non-deterministic result. InfluxQL-styleFIRST(f, t)/LAST(f, t)are rejected because InfluxDB v3 SQL has no such functions. GROUP BY DATE_BIN(INTERVAL 'N unit', time)— optional. When omitted, aggregate queries return a single scalar row (COUNTover an empty result set is0; other aggregates returnnil).GROUP BY <col>[, <col>...]— bucket points by tag/field values, with or without an aggregate (SELECT host FROM m GROUP BY hostis one row per host). Bare column names (with optionalAS alias) are valid in theSELECTlist only when grouped; a projected column that is neither grouped nor aggregated is the engine's planning error ("must appear in the GROUP BY clause or must be part of an aggregate function").ORDER BYapplies to grouped rows too.- Interval units:
seconds,minutes,hours,days
A null column is omitted from the row rather than present as nil,
exactly as InfluxDB 3's JSON and JSONL responses do (COUNT is 0, never
null).
Anything outside this subset is rejected with
{:error, %{status: 400, body: "Client.Local: ..."}}. The Client.Local:
prefix marks the rejection as a limitation of the test double rather than
of InfluxDB — the real engine may well accept the query. check_sql/1
answers the same question without executing, so a test can skip with a
reason and the query can be covered in the integration tier instead.
SQL Param Types
params: values are serialised to SQL literals before query execution.
Supported types: binary, integer, float, boolean, and Decimal
(when the optional :decimal dependency is loaded — Decimal values
are emitted as bare numeric literals via Decimal.to_string(:normal)).
Gzip Decompression
If a write payload begins with gzip magic bytes (0x1F 0x8B) it is automatically decompressed before line protocol parsing.
Timestamp Precision
Pass precision: :nanosecond | :microsecond | :millisecond | :second
in opts to normalise stored timestamps to nanoseconds.
Summary
Functions
Reports whether the SQL subset can express sql, without executing it.
Creates a named bucket in this local instance.
Creates a named database in this local instance.
Creates a synthetic API token and stores it in ETS.
Deletes a bucket from this local instance.
Deletes a database from this local instance.
Deletes a token by its id field. Returns :ok even if the token was
not found, matching real InfluxDB delete semantics.
Executes a SQL statement and returns a summary map.
Returns a passing health status map with string keys, matching the JSON-decoded shape returned by the HTTP client.
Returns all buckets in this local instance as a list of maps with a
single :name key.
Returns all databases created in this local instance as a list of maps
with a single :name key.
Executes a Flux query with support for common predicates.
Executes an InfluxQL query.
Executes a SQL-like query against stored ETS points and returns rows.
Executes a SQL query and returns results as a lazy Stream.
Starts a new LocalClient instance with isolated ETS storage.
Stops a LocalClient instance and cleans up its ETS table.
Returns true if the given operation is supported by the connection's profile.
Parses line protocol binary and stores the resulting points in ETS.
Types
@type conn() :: %{ table: :ets.table(), databases: MapSet.t(binary()), database: binary() | nil, profile: profile() }
@type point_map() :: InfluxElixir.Client.Local.LineProtocolParser.point()
@type profile() :: :v3_core | :v3_enterprise | :v2
Functions
Reports whether the SQL subset can express sql, without executing it.
Returns :ok or the same {:error, %{status: 400, body: "Client.Local: ..."}}
that query_sql/3 would return. Use it to skip a test with a reason
instead of tagging it excluded:
case InfluxElixir.Client.Local.check_sql(sql) do
:ok -> run_against_local(sql)
{:error, %{body: why}} -> ExUnit.Callbacks.on_exit(fn -> :ok end); flunk(why)
endQueries outside the subset (CTEs, joins, window functions, median, ...)
belong in an integration test against a real InfluxDB; see the testing
guide.
@spec create_bucket( InfluxElixir.Client.connection(), binary(), keyword() ) :: :ok | {:error, term()}
Creates a named bucket in this local instance.
Creating an already-existing bucket is idempotent.
@spec create_database( InfluxElixir.Client.connection(), binary(), keyword() ) :: :ok | {:error, term()}
Creates a named database in this local instance.
Always succeeds — creating an already-existing database is idempotent.
@spec create_token( InfluxElixir.Client.connection(), binary(), keyword() ) :: {:ok, map()} | {:error, term()}
Creates a synthetic API token and stores it in ETS.
Returns {:ok, %{id: id, token: token_string, description: desc}}.
@spec delete_bucket(InfluxElixir.Client.connection(), binary()) :: :ok | {:error, term()}
Deletes a bucket from this local instance.
Returns :ok whether or not the bucket exists, matching the idempotent
delete semantics of the v2 API.
@spec delete_database(InfluxElixir.Client.connection(), binary()) :: :ok | {:error, term()}
Deletes a database from this local instance.
Returns {:error, %{status: 404, body: "database not found: name"}} if
the database does not exist.
@spec delete_token(InfluxElixir.Client.connection(), binary()) :: :ok | {:error, term()}
Deletes a token by its id field. Returns :ok even if the token was
not found, matching real InfluxDB delete semantics.
@spec execute_sql(InfluxElixir.Client.connection(), binary(), keyword()) :: {:ok, map()} | {:error, term()}
Executes a SQL statement and returns a summary map.
Supports DELETE FROM <measurement> and
DELETE FROM <measurement> WHERE ... — matching points are removed
from ETS and the count is returned in %{"rows_affected" => N}.
On :v3_core profile, DELETE is not supported (matches real InfluxDB v3
Core behavior) and returns {:error, :delete_not_supported}.
On :v3_enterprise profile, DELETE is supported.
Unknown statements return %{"rows_affected" => 0}.
@spec health(InfluxElixir.Client.connection()) :: {:ok, map()} | {:error, term()}
Returns a passing health status map with string keys, matching the JSON-decoded shape returned by the HTTP client.
@spec list_buckets(InfluxElixir.Client.connection()) :: {:ok, [map()]} | {:error, term()}
Returns all buckets in this local instance as a list of maps with a
single :name key.
@spec list_databases(InfluxElixir.Client.connection()) :: {:ok, [map()]} | {:error, term()}
Returns all databases created in this local instance as a list of maps
with a single :name key.
@spec query_flux(InfluxElixir.Client.connection(), binary(), keyword()) :: InfluxElixir.Client.query_result()
Executes a Flux query with support for common predicates.
Parses and applies:
from(bucket: "...")— scopes to a databaserange(start: -1h)— filters by timestamp (supports-Nh,-Nd,-Nm)filter(fn: (r) => r._measurement == "...")— filters by measurementfilter(fn: (r) => r._field == "...")— keeps only that fieldfilter(fn: (r) => r.<key> == "...")— filters by any tag/field equality
Rows use the same long shape real Flux returns — one row per field,
ordered by table then _time:
%{"result" => "_result", "table" => 0, "_time" => %DateTime{},
"_measurement" => "cpu", "_field" => "value", "_value" => 1.0,
"host" => "web01"}table numbers each series (measurement + tags + field) from 0.
@spec query_influxql( InfluxElixir.Client.connection(), binary(), keyword() ) :: InfluxElixir.Client.query_result()
Executes an InfluxQL query.
Supports InfluxQL-specific commands:
SHOW DATABASES— returns all databasesSHOW MEASUREMENTS— returns all measurement namesSHOW TAG KEYS FROM <measurement>— returns distinct tag keysSELECT ...— delegates to the SQL engine
@spec query_sql(InfluxElixir.Client.connection(), binary(), keyword()) :: InfluxElixir.Client.query_result()
Executes a SQL-like query against stored ETS points and returns rows.
Supports:
SELECT * FROM measurementSELECT DISTINCT column FROM measurementWHERE key = 'value'/WHERE key > N/WHERE key < NORDER BY time ASC|DESCLIMIT N$paramplaceholder substitution viaparams: %{"$name" => value}
@spec query_sql_stream( InfluxElixir.Client.connection(), binary(), keyword() ) :: Enumerable.t()
Executes a SQL query and returns results as a lazy Stream.
Delegates to query_sql/3 then wraps the list in a stream.
Mirrors the error semantics of the HTTP client's streaming query:
because the return type is an Enumerable.t(), errors cannot be returned as a
tuple. Instead a failure — an underlying query error or an operation the
connection's profile does not support — is raised as an
InfluxElixir.StreamError when the stream is enumerated, never swallowed as an
empty result. This keeps Client.Local a faithful drop-in test double for
Client.HTTP, so consumer code that rescues InfluxElixir.StreamError can be
exercised against it.
Starts a new LocalClient instance with isolated ETS storage.
Options
:database- connection-level default database name. Used when the caller does not passdatabase:in opts. Pre-created automatically.:databases- list of database names to pre-create (default:[]):profile- InfluxDB version profile to emulate. Determines which operations are available. Operations outside the profile return{:error, :unsupported_operation}. Valid values::v3_core(default) — write, SQL, InfluxQL, database CRUD:v3_enterprise— everything in v3_core plus token management:v2— write, Flux, bucket CRUD
Examples
iex> {:ok, conn} = InfluxElixir.Client.Local.start(databases: ["mydb"])
iex> conn.profile
:v3_core
iex> {:ok, conn} = InfluxElixir.Client.Local.start(profile: :v2)
iex> conn.profile
:v2
iex> {:ok, conn} = InfluxElixir.Client.Local.start(database: "metrics")
iex> conn.database
"metrics"
@spec stop(conn()) :: :ok
Stops a LocalClient instance and cleans up its ETS table.
Safe to call multiple times; a no-op if the table is already deleted.
Returns true if the given operation is supported by the connection's profile.
@spec write(InfluxElixir.Client.connection(), binary(), keyword()) :: InfluxElixir.Client.write_result()
Parses line protocol binary and stores the resulting points in ETS.
The database is read from opts[:database]. If the database does not
exist an {:error, %{status: 404, body: ...}} is returned. If line protocol
cannot be parsed an {:error, %{status: 400, body: ...}} is returned.
Payloads beginning with gzip magic bytes are automatically decompressed.
Pass precision: :nanosecond | :microsecond | :millisecond | :second to
control how numeric timestamps are interpreted (default: :nanosecond).