Module: PgPipeline::RequestOps

Defined in:
lib/pg_pipeline/request.rb

Class Method Summary collapse

Class Method Details

.accept_result(req, result) ⇒ Object



123
124
125
126
127
128
129
130
131
132
# File 'lib/pg_pipeline/request.rb', line 123

def accept_result(req, result)
  assert_result_slot!(req)
  req.result_seen = true

  if req.cancelled
    clear_result(result)
  else
    req.result = result
  end
end

.assert_result_slot!(req) ⇒ Object

Raises:



218
219
220
221
# File 'lib/pg_pipeline/request.rb', line 218

def assert_result_slot!(req)
  raise ProtocolError, "result arrived after query boundary" if req.query_boundary_seen
  raise ProtocolError, "multiple results for one pipeline unit" if req.result_seen
end

.cancel!(req) ⇒ Object



172
173
174
175
176
177
178
179
# File 'lib/pg_pipeline/request.rb', line 172

def cancel!(req)
  return if req.cancelled

  req.cancelled = true
  clear_result(req.result)
  req.result = nil
  req.error.clear_result! if req.error.respond_to?(:clear_result!)
end

.clear_result(result) ⇒ Object



223
224
225
# File 'lib/pg_pipeline/request.rb', line 223

def clear_result(result)
  result.clear if result.respond_to?(:clear)
end

.finish!(req) ⇒ Object

Raises:



152
153
154
155
156
157
158
159
# File 'lib/pg_pipeline/request.rb', line 152

def finish!(req)
  raise ProtocolError, "sync arrived before query boundary" unless req.query_boundary_seen
  return if req.settled

  req.settled = true
  req.state = :done
  wake_waiter(req) unless req.cancelled
end

.park_waiter(req) ⇒ Object

Raises:



188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
# File 'lib/pg_pipeline/request.rb', line 188

def park_waiter(req)
  raise ProtocolError, "request already has a waiter" if req.waiter

  scheduler = Fiber.scheduler
  raise Error, "request wait requires an active Fiber scheduler" unless scheduler

  waiter = Fiber.current
  req.waiter = waiter
  req.waiter_scheduler = scheduler

  begin
    scheduler.block(req, nil) until req.settled
  ensure
    if req.waiter.equal?(waiter)
      req.waiter = nil
      req.waiter_scheduler = nil
    end
  end
end

.query_boundary!(req) ⇒ Object

Raises:



145
146
147
148
149
150
# File 'lib/pg_pipeline/request.rb', line 145

def query_boundary!(req)
  raise ProtocolError, "query boundary before query result" unless req.result_seen
  raise ProtocolError, "duplicate query boundary" if req.query_boundary_seen

  req.query_boundary_seen = true
end

.record_error!(req, error, result) ⇒ Object



134
135
136
137
138
139
140
141
142
143
# File 'lib/pg_pipeline/request.rb', line 134

def record_error!(req, error, result)
  assert_result_slot!(req)
  req.result_seen = true

  if req.cancelled
    clear_result(result)
  else
    req.error ||= error
  end
end

.reject!(req, error) ⇒ Object



161
162
163
164
165
166
167
168
169
170
# File 'lib/pg_pipeline/request.rb', line 161

def reject!(req, error)
  return if req.settled

  clear_result(req.result)
  req.result = nil
  req.error ||= error
  req.settled = true
  req.state = :done
  wake_waiter(req) unless req.cancelled
end

.snapshot_name(name) ⇒ Object

Raises:

  • (ArgumentError)


99
100
101
102
103
104
# File 'lib/pg_pipeline/request.rb', line 99

def snapshot_name(name)
  value = name.to_s
  raise ArgumentError, "statement_name must not be empty" if value.empty?

  value.frozen? ? value : value.dup.freeze
end

.snapshot_param_types(param_types) ⇒ Object



106
107
108
109
110
111
112
113
# File 'lib/pg_pipeline/request.rb', line 106

def snapshot_param_types(param_types)
  return nil if param_types.nil?
  raise ArgumentError, "param_types must be an Array or nil" unless param_types.is_a?(Array)

  param_types.map { |oid| oid.nil? ? nil : Integer(oid) }.freeze
rescue ArgumentError, TypeError
  raise ArgumentError, "param_types must contain only integer OIDs or nil"
end

.snapshot_params(params) ⇒ Object

Raises:

  • (ArgumentError)


79
80
81
82
83
84
# File 'lib/pg_pipeline/request.rb', line 79

def snapshot_params(params)
  values = params.nil? ? [] : params
  raise ArgumentError, "params must be an Array" unless values.is_a?(Array)

  values.map { |value| snapshot_value(value) }.freeze
end

.snapshot_sql(sql) ⇒ Object



74
75
76
77
# File 'lib/pg_pipeline/request.rb', line 74

def snapshot_sql(sql)
  value = sql.to_s
  value.frozen? ? value : value.dup.freeze
end

.snapshot_value(value) ⇒ Object



86
87
88
89
90
91
92
93
94
95
96
97
# File 'lib/pg_pipeline/request.rb', line 86

def snapshot_value(value)
  case value
  when String
    value.frozen? ? value : value.dup.freeze
  when Hash
    value.each_with_object({}) do |(key, item), copy|
      copy[key] = item.is_a?(String) && !item.frozen? ? item.dup.freeze : item
    end.freeze
  else
    value
  end
end

.transition!(req, from, to) ⇒ Object



115
116
117
118
119
120
121
# File 'lib/pg_pipeline/request.rb', line 115

def transition!(req, from, to)
  unless req.state == from
    raise ProtocolError, "invalid request transition #{req.state.inspect} -> #{to.inspect}"
  end

  req.state = to
end

.wait(req) ⇒ Object



181
182
183
184
185
186
# File 'lib/pg_pipeline/request.rb', line 181

def wait(req)
  park_waiter(req) unless req.settled
  raise req.error if req.error

  req.result
end

.wake_waiter(req) ⇒ Object



208
209
210
211
212
213
214
215
216
# File 'lib/pg_pipeline/request.rb', line 208

def wake_waiter(req)
  waiter = req.waiter
  scheduler = req.waiter_scheduler

  req.waiter = nil
  req.waiter_scheduler = nil

  scheduler.unblock(req, waiter) if scheduler && waiter
end