Packages
electric
0.7.5
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/application.ex
defmodule Electric.Application do
use Application
@process_registry_name Electric.Registry.Processes
@spec process_name(atom(), atom()) :: {:via, atom(), atom()}
def process_name(electric_instance_id, module) when is_atom(module) do
{:via, Registry, {@process_registry_name, {module, electric_instance_id}}}
end
@spec process_name(atom(), atom(), term()) :: {:via, atom(), {atom(), term()}}
def process_name(electric_instance_id, module, id) when is_atom(module) do
{:via, Registry, {@process_registry_name, {module, electric_instance_id, id}}}
end
@impl true
def start(_type, _args) do
:erlang.system_flag(:backtrace_depth, 50)
{storage_module, storage_opts} = Application.fetch_env!(:electric, :storage)
{kv_module, kv_fun, kv_params} = Application.fetch_env!(:electric, :persistent_kv)
persistent_kv = apply(kv_module, kv_fun, [kv_params])
replication_stream_id = Application.fetch_env!(:electric, :replication_stream_id)
publication_name = "electric_publication_#{replication_stream_id}"
slot_name = "electric_slot_#{replication_stream_id}"
with {:ok, storage_opts} <- storage_module.shared_opts(storage_opts) do
storage = {storage_module, storage_opts}
get_pg_version = fn ->
Electric.ConnectionManager.get_pg_version(Electric.ConnectionManager)
end
get_service_status = fn ->
Electric.ServiceStatus.check(
get_connection_status: fn ->
Electric.ConnectionManager.get_status(Electric.ConnectionManager)
end
)
end
prepare_tables_fn =
{Electric.Postgres.Configuration, :configure_tables_for_replication!,
[get_pg_version, publication_name]}
inspector =
{Electric.Postgres.Inspector.EtsInspector,
server: Electric.Postgres.Inspector.EtsInspector}
core_processes = [
{Registry,
name: @process_registry_name, keys: :unique, partitions: System.schedulers_online()}
]
per_env_processes =
if Application.fetch_env!(:electric, :environment) != :test do
electric_instance_id = Application.fetch_env!(:electric, :electric_instance_id)
shape_log_collector = Electric.Replication.ShapeLogCollector.name(electric_instance_id)
shape_cache =
{Electric.ShapeCache,
electric_instance_id: electric_instance_id,
storage: storage,
inspector: inspector,
prepare_tables_fn: prepare_tables_fn,
chunk_bytes_threshold: Application.fetch_env!(:electric, :chunk_bytes_threshold),
log_producer: shape_log_collector,
consumer_supervisor: Electric.Shapes.ConsumerSupervisor.name(electric_instance_id),
registry: Registry.ShapeChanges}
connection_manager_opts = [
electric_instance_id: electric_instance_id,
connection_opts: Application.fetch_env!(:electric, :connection_opts),
replication_opts: [
publication_name: publication_name,
try_creating_publication?: true,
slot_name: slot_name,
transaction_received:
{Electric.Replication.ShapeLogCollector, :store_transaction,
[shape_log_collector]},
relation_received:
{Electric.Replication.ShapeLogCollector, :handle_relation_msg,
[shape_log_collector]}
],
pool_opts: [
name: Electric.DbPool,
pool_size: Application.fetch_env!(:electric, :db_pool_size),
types: PgInterop.Postgrex.Types
],
timeline_opts: [
shape_cache: {Electric.ShapeCache, []},
persistent_kv: persistent_kv
],
log_collector:
{Electric.Replication.ShapeLogCollector,
electric_instance_id: electric_instance_id, inspector: inspector},
shape_cache: shape_cache
]
[
Electric.Telemetry,
{Registry,
name: Registry.ShapeChanges, keys: :duplicate, partitions: System.schedulers_online()},
{Electric.ConnectionManager, connection_manager_opts},
{Electric.Postgres.Inspector.EtsInspector, pool: Electric.DbPool},
{Bandit,
plug:
{Electric.Plug.Router,
storage: storage,
registry: Registry.ShapeChanges,
shape_cache: shape_cache,
get_service_status: get_service_status,
inspector: inspector,
long_poll_timeout: 20_000,
max_age: Application.fetch_env!(:electric, :cache_max_age),
stale_age: Application.fetch_env!(:electric, :cache_stale_age),
allow_shape_deletion: Application.get_env(:electric, :allow_shape_deletion, false)},
port: Application.fetch_env!(:electric, :service_port),
thousand_island_options: http_listener_options()}
]
|> add_prometheus_router(Application.fetch_env!(:electric, :prometheus_port))
else
[]
end
Supervisor.start_link(core_processes ++ per_env_processes,
strategy: :one_for_one,
name: Electric.Supervisor
)
end
end
defp add_prometheus_router(children, nil), do: children
defp add_prometheus_router(children, port) do
children ++
[
{
Bandit,
plug: {Electric.Plug.UtilityRouter, []},
port: port,
thousand_island_options: http_listener_options()
}
]
end
defp http_listener_options do
if Application.get_env(:electric, :listen_on_ipv6?, false) do
[transport_options: [:inet6]]
else
[]
end
end
end