Packages
postgrex
0.17.3
1.0.0-rc.1
retired
1.0.0-rc.0
retired
0.22.3
0.22.2
0.22.1
0.22.0
0.21.1
0.21.0
0.20.0
0.19.3
0.19.2
0.19.1
0.19.0
0.18.0
0.17.5
0.17.4
0.17.3
0.17.2
0.17.1
0.17.0
0.16.5
0.16.4
0.16.3
0.16.2
0.16.1
0.16.0
0.15.13
0.15.12
0.15.11
0.15.10
0.15.9
0.15.8
0.15.7
0.15.6
0.15.5
0.15.4
0.15.3
0.15.2
0.15.1
0.15.0
0.14.3
0.14.2
0.14.1
0.14.0
0.14.0-rc.1
0.14.0-rc.0
0.13.5
0.13.4
0.13.3
0.13.2
0.13.1
0.13.0
0.13.0-rc.0
0.12.2
0.12.1
0.12.0
0.11.2
0.11.1
0.11.0
0.10.0
0.9.1
0.9.0
0.8.4
0.8.3
0.8.2
0.8.1
0.8.0
0.7.0
0.6.0
0.5.5
0.5.4
0.5.3
0.5.2
0.5.1
0.5.0
0.4.2
PostgreSQL driver for Elixir
Security advisory:
This version has known vulnerabilities.
View advisories
Current section
Files
Jump to
Current section
Files
lib/postgrex/stream.ex
defmodule Postgrex.Stream do
@moduledoc """
Stream struct returned from stream commands.
All of its fields are private.
"""
@derive {Inspect, only: []}
defstruct [:conn, :query, :params, :options]
@type t :: %Postgrex.Stream{}
end
defmodule Postgrex.Cursor do
@moduledoc false
defstruct [:portal, :ref, :connection_id, :mode]
@type t :: %Postgrex.Cursor{}
end
defmodule Postgrex.Copy do
@moduledoc false
defstruct [:portal, :ref, :connection_id, :query]
@type t :: %Postgrex.Copy{}
end
defimpl Enumerable, for: Postgrex.Stream do
alias Postgrex.Query
def reduce(%Postgrex.Stream{query: %Query{} = query} = stream, acc, fun) do
%Postgrex.Stream{conn: conn, params: params, options: opts} = stream
stream = %DBConnection.Stream{conn: conn, query: query, params: params, opts: opts}
DBConnection.reduce(stream, acc, fun)
end
def reduce(%Postgrex.Stream{query: statement} = stream, acc, fun) do
%Postgrex.Stream{conn: conn, params: params, options: opts} = stream
query = %Query{name: "", statement: statement}
opts = Keyword.put(opts, :function, :prepare_open)
stream = %DBConnection.PrepareStream{conn: conn, query: query, params: params, opts: opts}
DBConnection.reduce(stream, acc, fun)
end
def member?(_, _) do
{:error, __MODULE__}
end
def count(_) do
{:error, __MODULE__}
end
def slice(_) do
{:error, __MODULE__}
end
end
defimpl Collectable, for: Postgrex.Stream do
alias Postgrex.Stream
alias Postgrex.Query
def into(%Stream{conn: %DBConnection{}} = stream) do
%Stream{conn: conn, query: query, params: params, options: opts} = stream
opts = Keyword.put(opts, :postgrex_copy, true)
case query do
%Query{} ->
copy = DBConnection.execute!(conn, query, params, opts)
{:ok, make_into(conn, stream, copy, opts)}
query ->
query = %Query{name: "", statement: query}
{_, copy} = DBConnection.prepare_execute!(conn, query, params, opts)
{:ok, make_into(conn, stream, copy, opts)}
end
end
def into(_) do
raise ArgumentError, "data can only be copied to database inside a transaction"
end
defp make_into(conn, stream, %Postgrex.Copy{ref: ref} = copy, opts) do
fn
:ok, {:cont, data} ->
_ = DBConnection.execute!(conn, copy, {:copy_data, ref, data}, opts)
:ok
:ok, close when close in [:done, :halt] ->
_ = DBConnection.execute!(conn, copy, {:copy_done, ref}, opts)
stream
end
end
end
defimpl DBConnection.Query, for: Postgrex.Copy do
alias Postgrex.Copy
import Postgrex.Messages
def parse(copy, _) do
raise "can not prepare #{inspect(copy)}"
end
def describe(copy, _) do
raise "can not describe #{inspect(copy)}"
end
def encode(%Copy{ref: ref}, {:copy_data, ref, data}, _) do
try do
encode_msg(msg_copy_data(data: data))
rescue
ArgumentError ->
reraise ArgumentError,
[message: "expected iodata to copy to database, got: " <> inspect(data)],
__STACKTRACE__
else
iodata ->
{:copy_data, iodata}
end
end
def encode(%Copy{ref: ref}, {:copy_done, ref}, _) do
:copy_done
end
def decode(copy, _result, _opts) do
raise "can not describe #{inspect(copy)}"
end
end
defimpl String.Chars, for: Postgrex.Copy do
def to_string(%Postgrex.Copy{query: query}) do
String.Chars.to_string(query)
end
end