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.



33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
# File 'lib/pg_pipeline/connection_driver.rb', line 33

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!
  DriverOps.warn_flush_coupling_once(@caps)

  @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

  @readable_events = 0
  @results_read = 0
  @units_completed = 0
  @flush_calls = 0
  @flush_incomplete = 0
  @dispatches = 0

  @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

#dispatchesObject



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

def dispatches = @dispatches || 0

#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

#flush_callsObject



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

def flush_calls = @flush_calls || 0

#flush_incompleteObject



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

def flush_incomplete = @flush_incomplete || 0

#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

#readable_eventsObject



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

def readable_events = @readable_events || 0

#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

#results_readObject



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

def results_read = @results_read || 0

#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

#units_completedObject



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

def units_completed = @units_completed || 0

#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



121
122
123
# File 'lib/pg_pipeline/connection_driver.rb', line 121

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

#available?Boolean

Returns:

  • (Boolean)


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

def available? = @accepting && @running

#dead?Boolean

Returns:

  • (Boolean)


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

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

#graceful_closeObject



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

def graceful_close = DriverOps.graceful_close(self)

#health_check(timeout) ⇒ Object



98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
# File 'lib/pg_pipeline/connection_driver.rb', line 98

def health_check(timeout)
  return true unless available?

  probe = Request.build("SELECT 1", nil)
  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



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

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

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



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

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

#statsObject



77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
# File 'lib/pg_pipeline/connection_driver.rb', line 77

def stats
  {
    available: available?,
    load: load,
    pending: @requests.size,
    in_flight: @inflight.size,
    submitting: @submitting,
    needs_flush: @needs_flush,
    fast_sync: @caps.fast_sync?,
    readable_events: readable_events,
    results_read: results_read,
    units_completed: units_completed,
    flush_calls: flush_calls,
    flush_incomplete: flush_incomplete,
    dispatches: dispatches,
    units_per_readable: DriverOps.ratio(units_completed, readable_events),
    results_per_readable: DriverOps.ratio(results_read, readable_events),
    flush_calls_per_unit: DriverOps.ratio(flush_calls, units_completed)
  }
end

#submit(request) ⇒ Object



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

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