Packages
electric
1.1.9
1.7.8
1.7.7
1.7.6
1.7.5
1.7.4
1.7.3
1.7.2
1.7.1
1.7.0
1.6.10
1.6.9
1.6.8
1.6.7
1.6.6
1.6.5
1.6.4
1.6.3
1.6.2
1.6.1
1.6.0
1.5.1
1.5.0
1.4.16
1.4.16-beta-1
1.4.15
1.4.14
1.4.13
1.4.12
1.4.11
1.4.10
1.4.8
1.4.7
1.4.6
1.4.5
1.4.4
1.4.3
1.4.2
1.4.1
1.4.0
1.3.4
1.3.3
1.3.2
1.2.4
1.2.3
1.2.2
1.2.1
1.2.0
1.1.14
1.1.13
1.1.12
1.1.11
1.1.10
1.1.9
1.1.8
1.1.7
1.1.6
retired
1.1.5
retired
1.1.4
retired
1.1.3
retired
1.1.2
1.1.1
1.1.0
1.0.24
1.0.23
1.0.22
1.0.21
1.0.20
1.0.19
1.0.18
1.0.17
1.0.15
1.0.13
1.0.12
1.0.11
1.0.10
1.0.9
1.0.5
1.0.4
1.0.3
1.0.2
1.0.1
1.0.0
1.0.0-beta.23
1.0.0-beta.22
1.0.0-beta.20
1.0.0-beta.19
1.0.0-beta.18
1.0.0-beta.17
1.0.0-beta.16
1.0.0-beta.15
1.0.0-beta.14
1.0.0-beta.13
1.0.0-beta.12
1.0.0-beta.11
1.0.0-beta.10
1.0.0-beta.9
1.0.0-beta.8
1.0.0-beta.7
1.0.0-beta.6
1.0.0-beta.5
1.0.0-beta.4
1.0.0-beta.3
1.0.0-beta.2
1.0.0-beta.1
0.9.5
0.9.4
0.9.3
0.9.2
0.9.1
0.9.0
0.8.1
0.8.0
0.7.7
0.7.6
0.7.5
0.7.4
0.7.3
0.7.2
0.7.1
0.7.0
0.6.3
0.6.2
0.6.1
0.5.2
0.4.4
Postgres sync engine. Sync little subsets of your Postgres data into local apps and services.
Current section
Files
Jump to
Current section
Files
lib/electric/config.ex
defmodule Electric.Config.Defaults do
@moduledoc false
# we want the default storage and kv implementations to honour the
# `:storage_dir` configuration setting so we need to use runtime-evaluated
# functions to get them. Since you can't embed anoymous functions these
# functions are used instead.
@doc false
def storage(opts \\ []) do
storage_dir = Keyword.get_lazy(opts, :storage_dir, fn -> storage_dir("shapes") end)
{Electric.ShapeCache.PureFileStorage, storage_dir: storage_dir}
end
@doc false
def persistent_kv(opts \\ []) do
storage_dir = Keyword.get_lazy(opts, :storage_dir, fn -> storage_dir("state") end)
{Electric.PersistentKV.Filesystem, :new!, root: storage_dir}
end
defp storage_dir(sub_dir) do
Path.join(storage_dir(), sub_dir)
end
defp storage_dir do
Electric.Config.get_env(:storage_dir)
end
def process_registry_partitions do
System.schedulers_online()
end
end
defmodule Electric.Config do
require Logger
@type instance_id :: String.t()
@build_env Mix.env()
@known_feature_flags ~w[allow_subqueries]
@defaults [
## Database
provided_database_id: "single_stack",
db_pool_size: 20,
replication_stream_id: "default",
replication_slot_temporary?: false,
replication_slot_temporary_random_name?: false,
max_txn_size: 250 * 1024 * 1024,
manual_table_publishing?: false,
## HTTP API
# set enable_http_api: false to turn off the HTTP server totally
enable_http_api: true,
long_poll_timeout: 20_000,
http_api_num_acceptors: nil,
tcp_send_timeout: :timer.seconds(30),
cache_max_age: 60,
cache_stale_age: 60 * 5,
chunk_bytes_threshold: Electric.ShapeCache.LogChunker.default_chunk_size_threshold(),
allow_shape_deletion?: false,
service_port: 3000,
listen_on_ipv6?: false,
stack_ready_timeout: 5_000,
send_cache_headers?: true,
max_shapes: nil,
## Storage
storage_dir: "./persistent",
storage: &Electric.Config.Defaults.storage/0,
persistent_kv: &Electric.Config.Defaults.persistent_kv/0,
## Telemetry
instance_id: nil,
prometheus_port: nil,
call_home_telemetry?: @build_env == :prod,
telemetry_statsd_host: nil,
telemetry_url: URI.new!("https://checkpoint.electric-sql.com"),
system_metrics_poll_interval: :timer.seconds(5),
otel_export_period: :timer.seconds(30),
otel_per_process_metrics?: false,
otel_sampling_ratio: 0.01,
metrics_sampling_ratio: 1,
telemetry_top_process_count: 5,
telemetry_long_gc_threshold: 500,
telemetry_long_schedule_threshold: 500,
telemetry_long_message_queue_enable_threshold: 1000,
telemetry_long_message_queue_disable_threshold: 100,
## Memory
shape_hibernate_after: :timer.seconds(30),
## Performance tweaks
publication_alter_debounce_ms: 0,
## Misc
process_registry_partitions: &Electric.Config.Defaults.process_registry_partitions/0,
feature_flags: if(Mix.env() == :test, do: @known_feature_flags, else: []),
schema_reconciler_period: 60_000
]
@installation_id_key "electric_installation_id"
def default(key) do
case Keyword.fetch!(@defaults, key) do
fun when is_function(fun, 0) -> fun.()
value -> value
end
end
@doc false
@spec ensure_instance_id() :: instance_id()
# the instance id needs to be consistent across calls, so we do need to have
# a value in the config, even if it's not configured by the user.
def ensure_instance_id do
case Application.get_env(:electric, :instance_id) do
nil ->
instance_id = generate_instance_id()
Logger.info("Setting electric instance_id: #{instance_id}")
Application.put_env(:electric, :instance_id, instance_id)
instance_id
id when is_binary(id) ->
id
end
end
defp generate_instance_id do
Electric.Utils.uuid4()
end
# the installation id is persisted to disk to remain the same between restarts of the sync service
@spec persist_installation_id(term, binary) :: instance_id()
def persist_installation_id(persistent_kv, instance_id) when is_binary(instance_id) do
case Electric.PersistentKV.get(persistent_kv, @installation_id_key) do
{:ok, id} when is_binary(id) ->
id
{:error, :not_found} ->
:ok = Electric.PersistentKV.set(persistent_kv, @installation_id_key, instance_id)
instance_id
end
end
@spec installation_id!(term) :: binary | no_return
def installation_id!(kv) do
case Electric.PersistentKV.get(kv, @installation_id_key) do
{:ok, id} when is_binary(id) -> id
{:error, :not_found} -> raise "Electric's installation_id not set"
end
end
@spec get_env(Application.key()) :: Application.value()
def get_env(key) do
# handle the case where the config value was set in runtime.exs but to
# `nil` because of a missing env var. This allows us to just use `nil`
# as the default config values in runtime.exs so avoiding hard-coding
# defaults all over the place.
case Application.get_env(:electric, key) do
nil -> default(key)
value -> value
end
end
def get_env_lazy(key, fun) when is_function(fun, 0) do
case Application.fetch_env(:electric, key) do
{:ok, nil} -> fun.()
{:ok, value} -> value
:error -> fun.()
end
end
@spec fetch_env!(Application.key()) :: Application.value()
def fetch_env!(key) do
Application.fetch_env!(:electric, key)
end
def persistent_kv do
with {m, f, a} <- get_env(:persistent_kv) do
apply(m, f, [a])
end
end
@doc ~S"""
Parse a PostgreSQL URI into a keyword list.
## Examples
iex> parse_postgresql_uri("postgresql://postgres:password@example.com/app-db") |> deobfuscate()
{:ok, [
hostname: "example.com",
port: 5432,
database: "app-db",
username: "postgres",
password: "password",
]}
iex> parse_postgresql_uri("postgresql://electric@192.168.111.33:81/__shadow")
{:ok, [
hostname: "192.168.111.33",
port: 81,
database: "__shadow",
username: "electric"
]}
iex> parse_postgresql_uri("postgresql://pg@[2001:db8::1234]:4321")
{:ok, [
hostname: "2001:db8::1234",
port: 4321,
database: "pg",
username: "pg"
]}
iex> parse_postgresql_uri("postgresql://user@localhost:5433/")
{:ok, [
hostname: "localhost",
port: 5433,
database: "user",
username: "user"
]}
iex> parse_postgresql_uri("postgresql://user%2Btesting%40gmail.com:weird%2Fpassword@localhost:5433/my%2Bdb%2Bname") |> deobfuscate()
{:ok, [
hostname: "localhost",
port: 5433,
database: "my+db+name",
username: "user+testing@gmail.com",
password: "weird/password"
]}
iex> parse_postgresql_uri("postgres://super_user@localhost:7801/postgres?sslmode=disable")
{:ok, [
hostname: "localhost",
port: 7801,
database: "postgres",
username: "super_user",
sslmode: :disable
]}
iex> parse_postgresql_uri("postgres://super_user@localhost:7801/postgres?sslmode=require")
{:ok, [
hostname: "localhost",
port: 7801,
database: "postgres",
username: "super_user",
sslmode: :require
]}
iex> parse_postgresql_uri("postgres://super_user@localhost:7801/postgres?sslmode=yesplease")
{:error, "invalid \"sslmode\" value: \"yesplease\""}
iex> parse_postgresql_uri("postgrex://localhost")
{:error, "invalid URL scheme: \"postgrex\""}
iex> parse_postgresql_uri("postgresql://localhost")
{:error, "invalid or missing username"}
iex> parse_postgresql_uri("postgresql://:@localhost")
{:error, "invalid or missing username"}
iex> parse_postgresql_uri("postgresql://:password@localhost")
{:error, "invalid or missing username"}
iex> parse_postgresql_uri("postgresql://user:password")
{:error, "invalid or missing username"}
iex> parse_postgresql_uri("postgresql://user:password@")
{:error, "missing host"}
iex> parse_postgresql_uri("postgresql://user@localhost:5433/mydb?opts=-c%20synchronous_commit%3Doff&foo=bar")
{:error, "unsupported query options: \"foo\", \"opts\""}
iex> parse_postgresql_uri("postgresql://electric@localhost/db?replication=database")
{:error, "unsupported \"replication\" query option. Electric opens both a replication connection and regular connections to Postgres as needed"}
iex> parse_postgresql_uri("postgresql://electric@localhost/db?replication=off")
{:error, "unsupported \"replication\" query option. Electric opens both a replication connection and regular connections to Postgres as needed"}
"""
@spec parse_postgresql_uri(binary) :: {:ok, keyword} | {:error, binary}
def parse_postgresql_uri(uri_str) do
%URI{scheme: scheme, host: host, port: port, path: path, userinfo: userinfo, query: query} =
URI.parse(uri_str)
with :ok <- validate_url_scheme(scheme),
:ok <- validate_url_host(host),
{:ok, {username, password}} <- parse_url_userinfo(userinfo),
{:ok, options} <- parse_url_query(query) do
conn_params =
Enum.reject(
[
hostname: host,
port: port || 5432,
database: parse_database(path, username) |> URI.decode(),
username: URI.decode(username),
password: if(password, do: password |> URI.decode() |> Electric.Utils.wrap_in_fun())
] ++ options,
fn {_key, val} -> is_nil(val) end
)
{:ok, conn_params}
end
end
def parse_postgresql_uri!(uri_str) do
case parse_postgresql_uri(uri_str) do
{:ok, results} -> results
{:error, message} -> raise Dotenvy.Error, message: message
end
end
defp validate_url_scheme(scheme) when scheme in ["postgres", "postgresql"], do: :ok
defp validate_url_scheme(scheme), do: {:error, "invalid URL scheme: #{inspect(scheme)}"}
defp validate_url_host(str) do
if is_binary(str) and String.trim(str) != "" do
:ok
else
{:error, "missing host"}
end
end
defp parse_url_userinfo(str) do
with false <- is_nil(str),
{:ok, {username, password}} <- split_userinfo(str),
false <- String.trim(username) == "" do
{:ok, {username, password}}
else
_ -> {:error, "invalid or missing username"}
end
end
defp split_userinfo(str) do
case String.split(str, ":") do
[username] -> {:ok, {username, nil}}
[username, password] -> {:ok, {username, password}}
_ -> :error
end
end
defp parse_url_query(nil), do: {:ok, []}
defp parse_url_query(query_str) do
case URI.decode_query(query_str) do
empty when map_size(empty) == 0 ->
{:ok, []}
%{"sslmode" => sslmode} when sslmode in ~w[disable allow prefer require] ->
{:ok, sslmode: String.to_existing_atom(sslmode)}
%{"sslmode" => sslmode} when sslmode in ~w[verify-ca verify-full] ->
{:error,
"unsupported \"sslmode\" value #{inspect(sslmode)}. Use sslmode=require and set the ELECTRIC_DATABASE_CA_CERTIFICATE_FILE config to ensure Electric verifies database server identity"}
%{"sslmode" => sslmode} ->
{:error, "invalid \"sslmode\" value: #{inspect(sslmode)}"}
%{"replication" => _} ->
{:error,
"unsupported \"replication\" query option. Electric opens both a replication connection and regular connections to Postgres as needed"}
map ->
{:error,
"unsupported query options: " <>
(map |> Map.keys() |> Enum.sort() |> Enum.map_join(", ", &inspect/1))}
end
end
defp parse_database(nil, username), do: username
defp parse_database("/", username), do: username
defp parse_database("/" <> dbname, _username), do: dbname
@log_levels ~w[emergency alert critical error warning warn notice info debug]
@public_log_levels ~w[error warning info debug]
@spec parse_log_level(binary) :: {:ok, Logger.level()} | {:error, binary}
def parse_log_level(str) when str in @log_levels do
{:ok, String.to_existing_atom(str)}
end
def parse_log_level(str) do
{:error, "invalid log level: #{inspect(str)}. Must be one of #{inspect(@public_log_levels)}"}
end
def parse_log_level!(str) when str in @log_levels, do: String.to_existing_atom(str)
def parse_log_level!(_str) do
raise Dotenvy.Error, message: "Must be one of #{inspect(@public_log_levels)}"
end
@spec parse_telemetry_url(binary) :: {:ok, binary} | {:error, binary}
def parse_telemetry_url(str) do
case URI.new(str) do
{:ok, %URI{scheme: scheme}} when scheme in ["http", "https"] -> {:ok, str}
_ -> {:error, "invalid URL format: \"#{str}\""}
end
end
def parse_telemetry_url!(str) do
case parse_telemetry_url(str) do
{:ok, url} -> url
{:error, message} -> raise Dotenvy.Error, message: message
end
end
@time_units ~w[ms msec s sec m min]
@spec parse_human_readable_time(binary | nil) :: {:ok, pos_integer} | {:error, binary}
def parse_human_readable_time(str) do
with {num, suffix} <- Float.parse(str),
true <- num > 0,
suffix = String.trim(suffix),
true <- suffix == "" or suffix in @time_units do
{:ok, trunc(num * time_multiplier(suffix))}
else
_ -> {:error, "invalid time unit: #{inspect(str)}. Must be one of #{inspect(@time_units)}"}
end
end
defp time_multiplier(""), do: 1
defp time_multiplier(millisecond) when millisecond in ["ms", "msec"], do: 1
defp time_multiplier(second) when second in ["s", "sec"], do: 1000
defp time_multiplier(minute) when minute in ["m", "min"], do: 1000 * 60
def parse_human_readable_time!(str) do
case parse_human_readable_time(str) do
{:ok, result} -> result
{:error, message} -> raise Dotenvy.Error, message: message
end
end
def validate_security_config!(secret, insecure) do
cond do
insecure && secret != nil ->
raise "You cannot set both ELECTRIC_SECRET and ELECTRIC_INSECURE=true"
!insecure && secret == nil ->
raise "You must set ELECTRIC_SECRET unless ELECTRIC_INSECURE=true. Setting ELECTRIC_INSECURE=true risks exposing your database, only use insecure mode in development or you've otherwise secured the Electric API"
true ->
if insecure do
Logger.warning(
"Electric is running in insecure mode - this risks exposing your database - only use insecure mode in development or if you've otherwise secured the Electric API."
)
end
:ok
end
end
@doc false
# helper function for use in doc tests
def deobfuscate({:ok, connection_opts}),
do: {:ok, Electric.Utils.deobfuscate_password(connection_opts)}
def deobfuscate(other), do: other
def parse_feature_flags(str) do
str
|> String.split(",")
|> Enum.map(&String.trim/1)
|> Enum.reject(&(&1 == ""))
|> Enum.split_with(&(&1 in @known_feature_flags))
|> case do
{known, []} ->
known
{_, unknown} ->
raise Dotenvy.Error,
message:
"Unknown feature flags specified: #{inspect(unknown)}. Known feature flags: #{inspect(@known_feature_flags)}"
end
end
end