Class: PgPipeline::ConnectionDriver

Inherits:
Object
  • Object
show all
Defined in:
lib/pg_pipeline/connection_driver.rb

Constant Summary collapse

DEFAULT_MAX_PENDING =
256
DEFAULT_MAX_IN_FLIGHT =
64

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(conn, max_pending: DEFAULT_MAX_PENDING, max_in_flight: DEFAULT_MAX_IN_FLIGHT) ⇒ ConnectionDriver

Returns a new instance of ConnectionDriver.



23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
# File 'lib/pg_pipeline/connection_driver.rb', line 23

def initialize(conn, max_pending: DEFAULT_MAX_PENDING, max_in_flight: DEFAULT_MAX_IN_FLIGHT)
  @max_pending = DriverOps.positive_integer!(max_pending, :max_pending)
  @max_in_flight = DriverOps.positive_integer!(max_in_flight, :max_in_flight)

  @conn = conn
  @caps = ServerCaps.from_connection(conn)
  @caps.assert_supported!

  @requests = BoundedQueue.new(@max_pending)
  @events = Async::Queue.new
  @reader_rearm = Async::Queue.new
  @writer_commands = Async::Queue.new

  @inflight = []
  @dispatching = nil
  @submitting = 0

  @accepting = false
  @running = false
  @draining = false
  @needs_flush = false
  @writer_armed = false
  @request_event_pending = false

  @socket = nil
  @owner_task = nil
  @reader_task = nil
  @writer_task = nil
end

Instance Attribute Details

#acceptingObject

Returns the value of attribute accepting.



18
19
20
# File 'lib/pg_pipeline/connection_driver.rb', line 18

def accepting
  @accepting
end

#capsObject (readonly)

Returns the value of attribute caps.



17
18
19
# File 'lib/pg_pipeline/connection_driver.rb', line 17

def caps
  @caps
end

#connObject (readonly)

Returns the value of attribute conn.



17
18
19
# File 'lib/pg_pipeline/connection_driver.rb', line 17

def conn
  @conn
end

#dispatchingObject

Returns the value of attribute dispatching.



18
19
20
# File 'lib/pg_pipeline/connection_driver.rb', line 18

def dispatching
  @dispatching
end

#drainingObject

Returns the value of attribute draining.



18
19
20
# File 'lib/pg_pipeline/connection_driver.rb', line 18

def draining
  @draining
end

#eventsObject

Returns the value of attribute events.



18
19
20
# File 'lib/pg_pipeline/connection_driver.rb', line 18

def events
  @events
end

#inflightObject

Returns the value of attribute inflight.



18
19
20
# File 'lib/pg_pipeline/connection_driver.rb', line 18

def inflight
  @inflight
end

#max_in_flightObject (readonly)

Returns the value of attribute max_in_flight.



17
18
19
# File 'lib/pg_pipeline/connection_driver.rb', line 17

def max_in_flight
  @max_in_flight
end

#max_pendingObject (readonly)

Returns the value of attribute max_pending.



17
18
19
# File 'lib/pg_pipeline/connection_driver.rb', line 17

def max_pending
  @max_pending
end

#needs_flushObject

Returns the value of attribute needs_flush.



18
19
20
# File 'lib/pg_pipeline/connection_driver.rb', line 18

def needs_flush
  @needs_flush
end

#owner_taskObject

Returns the value of attribute owner_task.



18
19
20
# File 'lib/pg_pipeline/connection_driver.rb', line 18

def owner_task
  @owner_task
end

#reader_rearmObject

Returns the value of attribute reader_rearm.



18
19
20
# File 'lib/pg_pipeline/connection_driver.rb', line 18

def reader_rearm
  @reader_rearm
end

#reader_taskObject

Returns the value of attribute reader_task.



18
19
20
# File 'lib/pg_pipeline/connection_driver.rb', line 18

def reader_task
  @reader_task
end

#request_event_pendingObject

Returns the value of attribute request_event_pending.



18
19
20
# File 'lib/pg_pipeline/connection_driver.rb', line 18

def request_event_pending
  @request_event_pending
end

#requestsObject

Returns the value of attribute requests.



18
19
20
# File 'lib/pg_pipeline/connection_driver.rb', line 18

def requests
  @requests
end

#runningObject

Returns the value of attribute running.



18
19
20
# File 'lib/pg_pipeline/connection_driver.rb', line 18

def running
  @running
end

#socketObject

Returns the value of attribute socket.



18
19
20
# File 'lib/pg_pipeline/connection_driver.rb', line 18

def socket
  @socket
end

#submittingObject

Returns the value of attribute submitting.



18
19
20
# File 'lib/pg_pipeline/connection_driver.rb', line 18

def submitting
  @submitting
end

#writer_armedObject

Returns the value of attribute writer_armed.



18
19
20
# File 'lib/pg_pipeline/connection_driver.rb', line 18

def writer_armed
  @writer_armed
end

#writer_commandsObject

Returns the value of attribute writer_commands.



18
19
20
# File 'lib/pg_pipeline/connection_driver.rb', line 18

def writer_commands
  @writer_commands
end

#writer_taskObject

Returns the value of attribute writer_task.



18
19
20
# File 'lib/pg_pipeline/connection_driver.rb', line 18

def writer_task
  @writer_task
end

Instance Method Details

#abort!(error = ConnectionLostError.new("connection aborted")) ⇒ Object



93
94
95
# File 'lib/pg_pipeline/connection_driver.rb', line 93

def abort!(error = ConnectionLostError.new("connection aborted"))
  DriverOps.abort!(self, error)
end

#available?Boolean

Returns:

  • (Boolean)


56
# File 'lib/pg_pipeline/connection_driver.rb', line 56

def available? = @accepting && @running

#dead?Boolean

Returns:

  • (Boolean)


57
# File 'lib/pg_pipeline/connection_driver.rb', line 57

def dead? = !@running && !@accepting

#graceful_closeObject



91
# File 'lib/pg_pipeline/connection_driver.rb', line 91

def graceful_close = DriverOps.graceful_close(self)

#health_check(timeout) ⇒ Object



70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
# File 'lib/pg_pipeline/connection_driver.rb', line 70

def health_check(timeout)
  return true unless available?

  probe = Request.new(sql: "SELECT 1", params: [])
  begin
    submit(probe)
    Async::Task.current.with_timeout(timeout) { probe.wait }
    true
  rescue Async::TimeoutError
    DriverOps.abort_timed_out_health_probe(self, probe)
  rescue QueryError, PipelineAbortedError
    true
  rescue ShutdownError, NotDispatchedError
    true
  rescue ConnectionLostError, PG::Error
    false
  ensure
    probe.cancel! unless probe.settled?
  end
end

#loadObject



55
# File 'lib/pg_pipeline/connection_driver.rb', line 55

def load = @requests.size + @inflight.size + @submitting + (@dispatching ? 1 : 0)

#start(parent: Async::Task.current) ⇒ Object



53
# File 'lib/pg_pipeline/connection_driver.rb', line 53

def start(parent: Async::Task.current) = DriverOps.start(self, parent)

#statsObject



59
60
61
62
63
64
65
66
67
68
# File 'lib/pg_pipeline/connection_driver.rb', line 59

def stats
  {
    available: available?,
    load: load,
    pending: @requests.size,
    in_flight: @inflight.size,
    submitting: @submitting,
    needs_flush: @needs_flush
  }
end

#submit(request) ⇒ Object



54
# File 'lib/pg_pipeline/connection_driver.rb', line 54

def submit(request) = DriverOps.submit(self, request)