Packages
ecto
0.6.0
3.14.1
3.14.0
3.13.6
3.13.5
3.13.4
3.13.3
3.13.2
3.13.1
3.13.0
3.12.6
3.12.5
3.12.4
3.12.3
3.12.2
3.12.1
3.12.0
3.11.2
3.11.1
3.11.0
3.10.3
3.10.2
3.10.1
3.10.0
3.9.6
3.9.5
3.9.4
3.9.3
3.9.2
3.9.1
3.9.0
3.8.4
3.8.3
3.8.2
3.8.1
3.8.0
3.7.2
3.7.1
3.7.0
3.6.2
3.6.1
3.6.0
3.5.8
3.5.7
3.5.6
3.5.5
3.5.4
3.5.3
3.5.2
3.5.1
3.5.0
3.5.0-rc.1
3.5.0-rc.0
3.4.6
3.4.5
3.4.4
3.4.3
3.4.2
3.4.1
3.4.0
3.3.4
3.3.3
3.3.2
3.3.1
3.3.0
3.2.5
3.2.4
3.2.3
3.2.2
3.2.1
3.2.0
3.1.7
3.1.6
3.1.5
3.1.4
3.1.3
3.1.2
3.1.1
3.1.0
3.0.9
3.0.8
3.0.7
3.0.6
3.0.5
3.0.4
3.0.3
3.0.2
3.0.1
3.0.0
3.0.0-rc.1
3.0.0-rc.0
2.2.12
2.2.11
2.2.10
2.2.9
2.2.8
2.2.7
2.2.6
2.2.5
2.2.4
2.2.3
2.2.2
2.2.1
2.2.0
2.2.0-rc.1
2.2.0-rc.0
2.1.6
2.1.5
2.1.4
2.1.3
2.1.2
2.1.1
2.1.0
2.1.0-rc.5
2.1.0-rc.4
2.1.0-rc.3
2.1.0-rc.2
2.1.0-rc.1
2.1.0-rc.0
2.0.6
2.0.5
2.0.4
2.0.3
2.0.2
2.0.1
2.0.0
2.0.0-rc.6
2.0.0-rc.5
2.0.0-rc.4
2.0.0-rc.3
2.0.0-rc.2
2.0.0-rc.1
2.0.0-rc.0
2.0.0-beta.2
2.0.0-beta.1
2.0.0-beta.0
1.1.9
1.1.8
1.1.7
1.1.6
1.1.5
1.1.4
1.1.3
1.1.2
1.1.1
1.1.0
1.0.7
1.0.6
1.0.5
1.0.4
1.0.3
1.0.2
1.0.1
1.0.0
0.16.0
0.15.0
0.14.3
0.14.2
0.14.1
0.14.0
0.13.1
0.13.0
0.12.1
0.12.0
0.12.0-rc
0.11.3
0.11.2
0.11.1
0.11.0
0.10.3
0.10.2
0.10.1
0.10.0
0.9.0
0.8.1
0.8.0
0.7.2
0.7.1
0.7.0
0.6.0
0.5.1
0.5.0
0.4.0
0.3.0
0.2.8
0.2.7
0.2.6
0.2.5
0.2.4
0.2.3
0.2.2
0.2.1
0.2.0
0.1.0
A toolkit for data mapping and language integrated query for Elixir
Current section
Files
Jump to
Current section
Files
lib/ecto/adapters/postgres/worker.ex
if Code.ensure_loaded?(Postgrex.Connection) do
defmodule Ecto.Adapters.Postgres.Worker do
@moduledoc false
use GenServer
def start(args) do
GenServer.start(__MODULE__, args)
end
def start_link(args) do
GenServer.start_link(__MODULE__, args)
end
def query!(worker, sql, params, opts) do
case GenServer.call(worker, {:query, sql, params, opts}, opts[:timeout]) do
{:ok, res} -> res
{:error, %Postgrex.Error{} = err} -> raise err
end
end
def begin!(worker, opts) do
case GenServer.call(worker, {:begin, opts}, opts[:timeout]) do
:ok -> :ok
{:error, %Postgrex.Error{} = err} -> raise err
end
end
def commit!(worker, opts) do
case GenServer.call(worker, {:commit, opts}, opts[:timeout]) do
:ok -> :ok
{:error, %Postgrex.Error{} = err} -> raise err
end
end
def rollback!(worker, opts) do
case GenServer.call(worker, {:rollback, opts}, opts[:timeout]) do
:ok -> :ok
{:error, %Postgrex.Error{} = err} -> raise err
end
end
def monitor_me(worker) do
GenServer.cast(worker, {:monitor, self})
end
def demonitor_me(worker) do
GenServer.cast(worker, {:demonitor, self})
end
def init(opts) do
Process.flag(:trap_exit, true)
lazy? = Keyword.get(opts, :lazy, true)
unless lazy? do
case Postgrex.Connection.start_link(opts) do
{:ok, conn} ->
conn = conn
_ ->
:ok
end
end
{:ok, %{conn: conn, params: opts, monitor: nil, transactions: 0}}
end
# Connection is disconnected, reconnect before continuing
def handle_call(request, from, %{conn: nil, params: params} = s) do
case Postgrex.Connection.start_link(params) do
{:ok, conn} ->
handle_call(request, from, %{s | conn: conn})
{:error, err} ->
{:reply, {:error, err}, s}
end
end
# TODO: Move query out of the worker to reduce copying
def handle_call({:query, sql, params, opts}, _from, %{conn: conn} = s) do
{:reply, Postgrex.Connection.query(conn, sql, params, opts), s}
end
def handle_call({:begin, opts}, _from, %{conn: conn, transactions: trans} = s) do
sql =
if trans == 0 do
"BEGIN"
else
"SAVEPOINT ecto_#{trans}"
end
reply =
case Postgrex.Connection.query(conn, sql, [], opts) do
{:ok, _} -> :ok
{:error, _} = err -> err
end
{:reply, reply, %{s | transactions: trans + 1}}
end
def handle_call({:commit, opts}, _from, %{conn: conn, transactions: trans} = s) when trans >= 1 do
reply =
case trans do
1 ->
case Postgrex.Connection.query(conn, "COMMIT", [], opts) do
{:ok, _} -> :ok
{:error, _} = err -> err
end
_ ->
:ok
end
{:reply, reply, %{s | transactions: trans - 1}}
end
def handle_call({:rollback, opts}, _from, %{conn: conn, transactions: trans} = s) when trans >= 1 do
sql =
case trans do
1 -> "ROLLBACK"
_ -> "ROLLBACK TO SAVEPOINT ecto_#{trans-1}"
end
reply =
case Postgrex.Connection.query(conn, sql, [], opts) do
{:ok, _} -> :ok
{:error, _} = err -> err
end
{:reply, reply, %{s | transactions: trans - 1}}
end
def handle_cast({:monitor, pid}, %{monitor: nil} = s) do
ref = Process.monitor(pid)
{:noreply, %{s | monitor: {pid, ref}}}
end
def handle_cast({:demonitor, pid}, %{monitor: {pid, ref}} = s) do
Process.demonitor(ref)
{:noreply, %{s | monitor: nil}}
end
def handle_info({:EXIT, conn, _reason}, %{conn: conn} = s) do
{:stop, :normal, %{s | conn: nil}}
end
def handle_info({:DOWN, ref, :process, pid, _info}, %{monitor: {pid, ref}} = s) do
{:stop, :normal, s}
end
def handle_info(_info, s) do
{:noreply, s}
end
def terminate(_reason, %{conn: conn}) do
if conn && Process.alive?(conn) do
Postgrex.Connection.stop(conn)
end
end
end
end