Current section

Files

Jump to
prestige lib prestige client.ex
Raw

lib/prestige/client.ex

defmodule Prestige.Client do
@moduledoc false
alias Prestige.Session
alias Prestige.Client.{Arguments, RequestStream, ResponseParser}
import Prestige.Client.RequestUtils
defmodule Request do
@moduledoc false
defstruct session: nil,
name: nil,
statement: nil,
args: nil,
headers: [],
prepare_statement: true
end
@spec execute(session :: Session.t(), statement :: String.t()) :: Enumerable.t()
def execute(session, statement) do
request = %Request{
session: session,
name: "stmt",
statement: statement,
args: [],
headers: [],
prepare_statement: false
}
RequestStream.stream(request)
|> ResponseParser.parse()
end
@spec execute(session :: Session.t(), name :: String.t(), statement :: String.t(), args :: list, headers :: list) ::
Enumerable.t()
def execute(session, name, statement, args, headers \\ []) do
request = %Request{session: session, name: name, statement: statement, args: args, headers: headers}
RequestStream.stream(request)
|> ResponseParser.parse()
end
def prepare_statement(session, name, statement, headers \\ []) do
prepare_statement = "PREPARE #{name} FROM #{statement}"
request = %Request{
session: session,
name: name,
statement: prepare_statement,
prepare_statement: false,
headers: headers
}
[result] =
RequestStream.stream(request)
|> ResponseParser.parse()
|> Enum.to_list()
presto_added_prepare_header = prefixed_header("added-prepare") |> String.downcase()
prepared_header = get_header(result.presto_headers, presto_added_prepare_header)
Session.add_prepared_statement(session, prepared_header)
end
@spec execute_statement(session :: Session.t(), name :: String.t(), args :: list, headers :: list) :: Enumerable.t()
def execute_statement(session, name, args, headers \\ []) do
execute_statement = "EXECUTE #{name} USING #{Arguments.to_arg_list(args)}"
request = %Request{
session: session,
name: name,
statement: execute_statement,
prepare_statement: false,
headers: headers
}
RequestStream.stream(request)
|> ResponseParser.parse()
end
def close_statement(session, name) do
deallocate_statement = "DEALLOCATE PREPARE #{name}"
request = %Request{
session: session,
statement: deallocate_statement,
prepare_statement: false
}
RequestStream.stream(request)
|> ResponseParser.parse()
|> Stream.run()
Session.remove_prepared_statement(session, name)
end
def start_transaction(session) do
transaction_id_header = prefixed_header("Transaction-Id")
presto_started_transaction_id_header = prefixed_header("started-transaction-id") |> String.downcase()
[result] =
execute(session, "stmt", "START TRANSACTION", [], [
{transaction_id_header, "none"}
])
|> Enum.to_list()
transaction_id = get_header(result.presto_headers, presto_started_transaction_id_header)
Session.set_transaction_id(session, transaction_id)
end
def rollback(session) do
execute(session, "stmt", "ROLLBACK", []) |> Enum.to_list()
end
def commit(session) do
execute(session, "stmt", "COMMIT", []) |> Enum.to_list()
end
defp get_header(headers, name) do
case Enum.find(headers, fn {key, _value} -> key == name end) do
{_key, value} -> value
nil -> nil
end
end
end