InfluxElixir.Client.Local.Store (InfluxElixir v0.1.35)

Copy Markdown View Source

The ETS store behind InfluxElixir.Client.Local: the one module that knows the key layout.

One :ordered_set per instance, :public so async: true tests can write from any process. Every mutation is a single insert or delete of its own key, so concurrent writers — tests sharing a database, BatchWriter flushes racing direct writes — never read-modify-write a shared value and no write is lost:

  • {:database, name} => true
  • {:bucket, name} => %{retention: seconds}
  • {:token, id} => the token map
  • {:point, database, measurement, seq} => the point — seq is a monotonic integer, so points scan in insertion order
  • {:series_time, database, measurement, tags, timestamp} — one per point written; a second write of the same key adds
  • {:duplicates, database, measurement} — reads merge that measurement's duplicate points (same tags and time) only when this marker exists
  • {: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

The store is policy-free: what a write may contain, which errors the engine returns and how a query reads rows live in Client.Local and the modules it delegates to.

Summary

Types

A stored point.

t()

A store instance.

Functions

Whether a bucket is registered.

The buckets as {name, meta}, sorted by name.

A column's registered kind, or nil — a read that registers nothing.

Every column in a database as {measurement, column, kind}, in that order.

Whether a database is registered.

The registered databases.

Removes a bucket; :error when it was not registered.

Deletes the points of a measurement that match? accepts, judged on the merged points as the engine sees them; every stored object behind a matching point is deleted by its own key, so a concurrent write is never lost. Returns the number of (merged) points deleted.

Removes a token (a missing one is fine).

Deletes the store. Safe to call again: the table dies with its owner, and an on_exit callback can run after that, so a missing table is :ok.

Drops a database with everything in it — points, schema, the duplicate index — so a re-created database starts empty. :error when it was not registered.

Whether any point was written to the measurement.

The database's measurements that hold points, first-written first.

Creates a store with databases registered.

The double's "now", in nanoseconds: used to stamp untimed points and for SQL and Flux now(). System.os_time/1 and System.system_time/1 can differ by microseconds (time warp), and a test stamps points with either; taking the later of the two means a point written a moment ago is never after now(), as it never is on a real server.

A measurement's points in insertion order, duplicates merged. InfluxDB — both versions, verified — treats points with the same tags and time as one point whose fields merge, the later write winning per field.

Every point in a database, duplicates merged.

Registers a bucket with its metadata (replacing any earlier one).

Registers a database (idempotent).

Stores a token under its id.

Registers a column's kind if it is new. The first writer fixes it atomically (insert_new), so a concurrent writer never loses a column. Returns :ok when the kind matches (or was just registered) and {:conflict, existing} when the column already has another kind.

Stores a point as written, one insert of its own key. A point without a timestamp gets the server's time, as on the engine.

Whether the measurement has any column registered (the table exists).

A measurement's tag columns.

Types

point()

A stored point.

t()

@type t() :: :ets.table()

A store instance.

Functions

bucket?(table, name)

@spec bucket?(t(), binary()) :: boolean()

Whether a bucket is registered.

buckets(table)

@spec buckets(t()) :: [{binary(), map()}]

The buckets as {name, meta}, sorted by name.

column_kind(table, database, measurement, column)

@spec column_kind(t(), binary(), binary(), binary()) :: binary() | nil

A column's registered kind, or nil — a read that registers nothing.

columns(table, database)

@spec columns(t(), binary()) :: [{binary(), binary(), binary()}]

Every column in a database as {measurement, column, kind}, in that order.

database?(table, name)

@spec database?(t(), binary()) :: boolean()

Whether a database is registered.

databases(table)

@spec databases(t()) :: MapSet.t(binary())

The registered databases.

delete_bucket(table, name)

@spec delete_bucket(t(), binary()) :: :ok | :error

Removes a bucket; :error when it was not registered.

delete_points(table, database, measurement, match?)

@spec delete_points(t(), binary(), binary(), (point() -> boolean())) ::
  non_neg_integer()

Deletes the points of a measurement that match? accepts, judged on the merged points as the engine sees them; every stored object behind a matching point is deleted by its own key, so a concurrent write is never lost. Returns the number of (merged) points deleted.

delete_token(table, id)

@spec delete_token(t(), binary()) :: true

Removes a token (a missing one is fine).

drop(table)

@spec drop(t()) :: :ok

Deletes the store. Safe to call again: the table dies with its owner, and an on_exit callback can run after that, so a missing table is :ok.

drop_database(table, name)

@spec drop_database(t(), binary()) :: :ok | :error

Drops a database with everything in it — points, schema, the duplicate index — so a re-created database starts empty. :error when it was not registered.

measurement?(table, database, measurement)

@spec measurement?(t(), binary(), binary()) :: boolean()

Whether any point was written to the measurement.

measurements(table, database)

@spec measurements(t(), binary()) :: [binary()]

The database's measurements that hold points, first-written first.

new(databases)

@spec new(Enumerable.t()) :: t()

Creates a store with databases registered.

now_ns()

@spec now_ns() :: integer()

The double's "now", in nanoseconds: used to stamp untimed points and for SQL and Flux now(). System.os_time/1 and System.system_time/1 can differ by microseconds (time warp), and a test stamps points with either; taking the later of the two means a point written a moment ago is never after now(), as it never is on a real server.

points(table, database, measurement)

@spec points(t(), binary(), binary()) :: [point()]

A measurement's points in insertion order, duplicates merged. InfluxDB — both versions, verified — treats points with the same tags and time as one point whose fields merge, the later write winning per field.

points_in_db(table, database)

@spec points_in_db(t(), binary()) :: [point()]

Every point in a database, duplicates merged.

put_bucket(table, name, meta)

@spec put_bucket(t(), binary(), map()) :: true

Registers a bucket with its metadata (replacing any earlier one).

put_database(table, name)

@spec put_database(t(), binary()) :: true

Registers a database (idempotent).

put_token(table, id, token)

@spec put_token(t(), binary(), map()) :: true

Stores a token under its id.

register_column(table, database, measurement, column, kind)

@spec register_column(t(), binary(), binary(), binary(), binary()) ::
  :ok | {:conflict, binary()}

Registers a column's kind if it is new. The first writer fixes it atomically (insert_new), so a concurrent writer never loses a column. Returns :ok when the kind matches (or was just registered) and {:conflict, existing} when the column already has another kind.

store_point(table, database, point)

@spec store_point(t(), binary(), point()) :: true

Stores a point as written, one insert of its own key. A point without a timestamp gets the server's time, as on the engine.

Storing a measurement's points as one list meant every write read the list, prepended and wrote it back: concurrent writers overwrote each other (#15: 159 of 480 writes survived) and each insert copied the whole list. A plain insert is atomic and O(log n).

table?(table, database, measurement)

@spec table?(t(), binary(), binary()) :: boolean()

Whether the measurement has any column registered (the table exists).

tag_columns(table, database, measurement)

@spec tag_columns(t(), binary(), binary()) :: MapSet.t(binary())

A measurement's tag columns.