InfluxElixir.Client.Local (InfluxElixir v0.1.28)

Copy Markdown View Source

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.

ProfileWriteSQLInfluxQLFluxDB CRUDBucket CRUDTokens
:v3_coreyesyesyesnoyesnono
:v3_enterpriseyesyesyesnoyesnoyes
:v2yesnonoyesnoyesno

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}
end

Checking 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")
end

ETS 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() — seq is a monotonic integer, so points scan in insertion order
  • {:column, database, measurement, column} => the column's kind (iox::column_type::tag or iox::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 one data entry 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=1i then v=2.0 is "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.
  • time is a reserved column; a key cannot be both a tag and a field on one line; an integer must fit in 64 bits (7u is 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 measurement
  • SELECT col1, col2 [, ...] FROM measurement with optional AS alias (projects fields and tags; time is selectable). time and DATE_BIN buckets are DateTime values with microsecond precision, the same as the HTTP and Flight transports return; compare them with DateTime.compare/2 or 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 BY may 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 final SELECT may read from any of them (FROM w). A CTE's output columns are its fields (time stays time).
  • FROM a CROSS JOIN b — every row of a paired with every row of b (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, HAVING and window functions are rejected by name rather than silently ignored.
  • Table qualifiers and aliases: FROM q AS w / FROM q w, and w.time, q.bid in any clause — one table per query, so the prefix is dropped.
  • WHERE with =, != / <>, <, <=, >, >=, combined with AND, OR, NOT and parentheses (AND binds tighter than OR, 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 (so amount >= '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 = 2 matches the tag "2"; rack > 3 does 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 (a nil param) is never true.
  • WHERE col IN (v1, v2, ...) and WHERE 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 NULL and WHERE col IS NOT NULL
  • WHERE col [NOT] BETWEEN low AND high (inclusive; time too)
  • WHERE col [NOT] LIKE 'pattern' and ILIKE (% any run, _ one character; LIKE is case-sensitive, ILIKE is not). LIKE over 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), or now() 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 BY must name a selected column, as in DataFusion)
  • ORDER BY a [ASC|DESC][, b [ASC|DESC] ...] — each term time, a column, an output alias, or (on raw and projected rows) an expression such as CAST(level AS INTEGER) DESC; every term applies, each with its own direction
  • CAST(expr AS INTEGER | INT | BIGINT | DOUBLE | FLOAT | VARCHAR | STRING) and DataFusion's col::TYPE shorthand, wherever an expression is allowed: WHERE (CAST(level AS INTEGER) <= 20 compares a numeric tag numerically), BETWEEN, LIKE, projections, aggregates, arithmetic and ORDER 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' to INTEGER, time to INTEGER) makes InfluxDB 3 Core drop the connection mid-response, which Client.HTTP reports as {:error, {:connection_error, %Mint.TransportError{reason: :closed}}}; the double reports {:error, {:connection_error, :closed}}. BOOLEAN and TIMESTAMP targets are outside the subset.

  • LIMIT n and OFFSET m, in either order — OFFSET skips rows before LIMIT takes them, on plain, projected, grouped and DISTINCT rows alike; LIMIT 0 returns no rows; a negative or non-numeric limit is rejected, as the engine rejects it
  • $param placeholders via params: %{"$name" => value} in opts. A DateTime, NaiveDateTime or Date param renders as the ISO-8601 string Jason sends over HTTP, so time >= $start works the same on both clients; an integer param against time is 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 JSON null but counted by COUNT) 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) and COUNT(time) work (a DateTime result); every other aggregate over time, 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]) and last_value(field ORDER BY col [ASC|DESC]) — the InfluxDB v3 SQL (DataFusion) spelling. The ORDER BY is 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-style FIRST(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 (COUNT over an empty result set is 0; other aggregates return nil).
  • GROUP BY <col>[, <col>...] — bucket points by tag/field values, with or without an aggregate (SELECT host FROM m GROUP BY host is one row per host). Bare column names (with optional AS alias) are valid in the SELECT list 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 BY applies 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

conn()

@type conn() :: %{
  table: :ets.table(),
  databases: MapSet.t(binary()),
  database: binary() | nil,
  profile: profile()
}

point_map()

profile()

@type profile() :: :v3_core | :v3_enterprise | :v2

Functions

check_sql(sql)

@spec check_sql(binary()) :: :ok | {:error, %{status: 400, body: binary()}}

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)
end

Queries outside the subset (CTEs, joins, window functions, median, ...) belong in an integration test against a real InfluxDB; see the testing guide.

create_bucket(conn, name, opts \\ [])

@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.

create_database(conn, name, opts \\ [])

@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.

create_token(conn, description, opts \\ [])

@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}}.

delete_bucket(conn, name)

@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.

delete_database(conn, name)

@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.

delete_token(conn, token_id)

@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.

execute_sql(conn, sql, opts \\ [])

@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}.

health(conn)

@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.

list_buckets(conn)

@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.

list_databases(conn)

@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.

query_flux(conn, flux, opts \\ [])

Executes a Flux query with support for common predicates.

Parses and applies:

  • from(bucket: "...") — scopes to a database
  • range(start: -1h) — filters by timestamp (supports -Nh, -Nd, -Nm)
  • filter(fn: (r) => r._measurement == "...") — filters by measurement
  • filter(fn: (r) => r._field == "...") — keeps only that field
  • filter(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.

query_influxql(conn, influxql, opts \\ [])

Executes an InfluxQL query.

Supports InfluxQL-specific commands:

  • SHOW DATABASES — returns all databases
  • SHOW MEASUREMENTS — returns all measurement names
  • SHOW TAG KEYS FROM <measurement> — returns distinct tag keys
  • SELECT ... — delegates to the SQL engine

query_sql(conn, sql, opts \\ [])

Executes a SQL-like query against stored ETS points and returns rows.

Supports:

  • SELECT * FROM measurement
  • SELECT DISTINCT column FROM measurement
  • WHERE key = 'value' / WHERE key > N / WHERE key < N
  • ORDER BY time ASC|DESC
  • LIMIT N
  • $param placeholder substitution via params: %{"$name" => value}

query_sql_stream(conn, sql, opts \\ [])

@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.

start(opts \\ [])

@spec start(keyword()) :: {:ok, conn()}

Starts a new LocalClient instance with isolated ETS storage.

Options

  • :database - connection-level default database name. Used when the caller does not pass database: 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"

stop(map)

@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.

supports?(map, operation)

@spec supports?(conn(), atom()) :: boolean()

Returns true if the given operation is supported by the connection's profile.

write(conn, payload, opts \\ [])

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).