Skip to content
17 changes: 9 additions & 8 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -88,8 +88,8 @@ $ bundle exec rake modelgen:latest

## Options

* **server** sets address (and port) of a Trino coordinator server.
* **ssl** enables https.
* **server** sets address (and port) of a Trino coordinator server. Changes after client initialization are not applied; create a new client to change this option.
* **ssl** enables https. Changes after client initialization are not applied; create a new client to change this option.
* Setting `true` enables SSL and verifies server certificate using system's built-in certificates.
* Setting `{verify: false}` enables SSL but doesn't verify server certificate.
* Setting a Hash object enables SSL and verify server certificate with options:
Expand All @@ -104,20 +104,21 @@ $ bundle exec rake modelgen:latest
* **client_info** sets client info to queries. It can be a string to pass a raw string, or an object that can be encoded to JSON.
* **client_tags** sets client tags to queries. It needs to be an array of strings. The tags are shown on web interface.
* **user** sets user name to connect to a Trino.
* **password** sets a password to connect to Trino using basic auth.
* **password** sets a password to connect to Trino using basic auth. Requires the client's connection to use HTTPS.
* **time_zone** sets time zone of queries. Time zone affects some functions such as `format_datetime`.
* **language** sets language of queries. Language affects some functions such as `format_datetime`.
* **properties** set session properties. Session properties affect internal behavior such as `hive.force_local_scheduling: true`, `raptor.reader_stream_buffer_size: "32MB"`, etc.
* **query_timeout** sets timeout in seconds for the entire query execution (from the first API call until there're no more output data). If timeout happens, client raises TrinoQueryTimeoutError. Default is nil (disabled).
* **plan_timeout** sets timeout in seconds for query planning execution (from the first API call until result columns become available). If timeout happens, client raises TrinoQueryTimeoutError. Default is nil (disabled).
* **http_headers** sets custom HTTP headers. It must be a Hash of string to string.
* **http_proxy** sets host:port of a HTTP proxy server.
* **http_debug** enables debug message to STDOUT for each HTTP requests.
* **http_debug_logger** sets a custom `Logger` instance for HTTP debug logs. Requires **http_debug** to be `true`.
* **http_proxy** sets host:port of a HTTP proxy server. Changes after client initialization are not applied; create a new client to change this option.
* **http_debug** enables debug message to STDOUT for each HTTP requests. Changes after client initialization are not applied; create a new client to change this option.
* **http_debug_logger** sets a custom `Logger` instance for HTTP debug logs. Requires **http_debug** to be `true`. Changes after client initialization are not applied; create a new client to change this option.
* **http_open_timeout** sets timeout in seconds to open new HTTP connection.
* **http_timeout** sets timeout in seconds to read data from a server.
* **gzip** enables gzip compression.
* **follow_redirect** enables HTTP redirection support.
* **gzip** enables gzip compression. Changes after client initialization are not applied; create a new client to change this option.
* **follow_redirect** enables HTTP redirection support. Changes after client initialization are not applied; create a new client to change this option.
* **faraday_adapter** sets the Faraday adapter to use. Default is `Faraday.default_adapter`. To reuse persistent HTTP connections, install `faraday-net_http_persistent` and specify `:net_http_persistent`. Changes after client initialization are not applied; create a new client to change this option.
* **model_version** set the Trino version to which a job is submitted. Supported versions are 351, 316, 303, 0.205, 0.178, 0.173, 0.153 and 0.149. Default is 351.

See [RDoc](http://www.rubydoc.info/gems/presto-client/) for the full documentation.
Expand Down
16 changes: 12 additions & 4 deletions lib/trino/client/client.rb
Original file line number Diff line number Diff line change
Expand Up @@ -19,12 +19,20 @@ module Trino::Client
require 'trino/client/query'

class Client
# Creates a reusable Faraday connection using the initial options.
# Connection settings are not reapplied after initialization.
# These include server, ssl, proxy, and faraday_adapter; see the README for the full list.
# Create a new Client to change these settings.
#
# Query-specific headers, including Basic Auth, use the options at query start.
# Password authentication requires the actual connection to use HTTPS, regardless of later changes to options[:ssl].
def initialize(options)
@options = options
@faraday = Trino::Client.faraday_client(options)
end

def query(query, &block)
q = Query.start(query, @options)
q = Query.start(query, @options, @faraday)
if block
begin
yield q
Expand All @@ -37,15 +45,15 @@ def query(query, &block)
end

def resume_query(next_uri)
return Query.resume(next_uri, @options)
return Query.resume(next_uri, @options, @faraday)
end

def kill(query_id)
return Query.kill(query_id, @options)
return Query.kill(query_id, @options, @faraday)
end

def run(query)
q = Query.start(query, @options)
q = Query.start(query, @options, @faraday)
begin
columns = q.columns
if columns.empty?
Expand Down
53 changes: 27 additions & 26 deletions lib/trino/client/faraday_client.rb
Original file line number Diff line number Diff line change
Expand Up @@ -14,9 +14,8 @@
# limitations under the License.
#
module Trino::Client
FARADAY1_USED = Faraday::VERSION.start_with?("1.")
private_constant :FARADAY1_USED

require 'base64'
require 'cgi'

module TrinoHeaders
Expand Down Expand Up @@ -77,26 +76,17 @@ def self.faraday_client(options)
faraday_options[:ssl] = ssl if ssl

faraday = Faraday.new(faraday_options) do |faraday|
if options[:user] && options[:password]
# https://lostisland.github.io/faraday/middleware/authentication
if FARADAY1_USED
faraday.request(:basic_auth, options[:user], options[:password])
else
faraday.request :authorization, :basic, options[:user], options[:password]
end
end
if options[:follow_redirect]
faraday.response :follow_redirects
end
if options[:gzip]
faraday.request :gzip
end
faraday.response :logger, options[:http_debug_logger] if options[:http_debug]
faraday.adapter Faraday.default_adapter
faraday.adapter(options[:faraday_adapter] || Faraday.default_adapter)
end

faraday.headers.merge!(HEADERS)
faraday.headers.merge!(optional_headers(options))

return faraday
end
Expand Down Expand Up @@ -129,71 +119,82 @@ def self.faraday_ssl_options(options)
return ssl
end

def self.optional_headers(options)
usePrestoHeader = false
if options[:model_version] && options[:model_version] < 351
usePrestoHeader = true
def self.build_query_headers(options, faraday:)
if options[:password] && faraday.url_prefix.scheme != "https"
raise ArgumentError, "Protocol must be https when passing a password"
end
use_presto_headers = false
if options[:model_version] && options[:model_version].to_i < 351
use_presto_headers = true
end

headers = {}

if options[:user] && options[:password]
credentials = Base64.strict_encode64(
"#{options[:user]}:#{options[:password]}"
)
headers["Authorization"] = "Basic #{credentials}"
end

if v = options[:user]
if usePrestoHeader
if use_presto_headers
headers[PrestoHeaders::PRESTO_USER] = v
else
headers[TrinoHeaders::TRINO_USER] = v
end
end
if v = options[:source]
if usePrestoHeader
if use_presto_headers
headers[PrestoHeaders::PRESTO_SOURCE] = v
else
headers[TrinoHeaders::TRINO_SOURCE] = v
end
end
if v = options[:catalog]
if usePrestoHeader
if use_presto_headers
headers[PrestoHeaders::PRESTO_CATALOG] = v
else
headers[TrinoHeaders::TRINO_CATALOG] = v
end
end
if v = options[:schema]
if usePrestoHeader
if use_presto_headers
headers[PrestoHeaders::PRESTO_SCHEMA] = v
else
headers[TrinoHeaders::TRINO_SCHEMA] = v
end
end
if v = options[:time_zone]
if usePrestoHeader
if use_presto_headers
headers[PrestoHeaders::PRESTO_TIME_ZONE] = v
else
headers[TrinoHeaders::TRINO_TIME_ZONE] = v
end
end
if v = options[:language]
if usePrestoHeader
if use_presto_headers
headers[PrestoHeaders::PRESTO_LANGUAGE] = v
else
headers[TrinoHeaders::TRINO_LANGUAGE] = v
end
end
if v = options[:properties]
if usePrestoHeader
if use_presto_headers
headers[PrestoHeaders::PRESTO_SESSION] = encode_properties(v)
else
headers[TrinoHeaders::TRINO_SESSION] = encode_properties(v)
end
end
if v = options[:client_info]
if usePrestoHeader
if use_presto_headers
headers[PrestoHeaders::PRESTO_CLIENT_INFO] = encode_client_info(v)
else
headers[TrinoHeaders::TRINO_CLIENT_INFO] = encode_client_info(v)
end
end
if v = options[:client_tags]
if usePrestoHeader
if use_presto_headers
headers[PrestoHeaders::PRESTO_CLIENT_TAGS] = encode_client_tags(v)
else
headers[TrinoHeaders::TRINO_CLIENT_TAGS] = encode_client_tags(v)
Expand Down Expand Up @@ -245,6 +246,6 @@ def self.encode_client_tags(tags)
Array(tags).join(",")
end

private_class_method :faraday_ssl_options, :optional_headers, :encode_properties, :encode_client_info, :encode_client_tags
private_class_method :faraday_ssl_options, :encode_properties, :encode_client_info, :encode_client_tags

end
16 changes: 10 additions & 6 deletions lib/trino/client/query.rb
Original file line number Diff line number Diff line change
Expand Up @@ -25,17 +25,21 @@ module Trino::Client
require 'trino/client/statement_client'

class Query
def self.start(query, options)
new StatementClient.new(faraday_client(options), query, options)
def self.start(query, options, faraday = nil)
faraday ||= faraday_client(options)
new StatementClient.new(faraday, query, options)
end

def self.resume(next_uri, options)
new StatementClient.new(faraday_client(options), nil, options, next_uri)
def self.resume(next_uri, options, faraday = nil)
faraday ||= faraday_client(options)
new StatementClient.new(faraday, nil, options, next_uri)
end

def self.kill(query_id, options)
faraday = faraday_client(options)
def self.kill(query_id, options, faraday = nil)
faraday ||= faraday_client(options)
headers = Trino::Client.build_query_headers(options, faraday: faraday)
response = faraday.delete do |req|
req.headers.merge!(headers)
req.url "/v1/query/#{query_id}"
end
return response.status / 100 == 2
Expand Down
8 changes: 7 additions & 1 deletion lib/trino/client/statement_client.rb
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ class StatementClient

def initialize(faraday, query, options, next_uri=nil)
@faraday = faraday
@headers = Trino::Client.build_query_headers(options, faraday: @faraday)

@options = options
@query = query
Expand Down Expand Up @@ -72,6 +73,7 @@ def post_query_request!
begin
r = @faraday.post do |req|
req.url uri
req.headers.merge!(@headers)

req.body = @query
init_request(req)
Expand Down Expand Up @@ -229,7 +231,9 @@ def with_retry_loop
def faraday_get_with_retry(uri)
with_retry_loop do
begin
response = @faraday.get(uri)
response = @faraday.get(uri) do |req|
req.headers.merge!(@headers)
end
rescue Faraday::TimeoutError, Faraday::ConnectionFailed
throw :retry_with_backoff
rescue => e
Expand Down Expand Up @@ -278,6 +282,7 @@ def raise_timeout_error!
def cancel_leaf_stage
if uri = @results.partial_cancel_uri
@faraday.delete do |req|
req.headers.merge!(@headers)
req.url uri
end
end
Expand All @@ -291,6 +296,7 @@ def close
begin
if uri = @results.next_uri
@faraday.delete do |req|
req.headers.merge!(@headers)
req.url uri
end
end
Expand Down
Loading
Loading