Module: PgPipeline::RequestOps

Defined in:
lib/pg_pipeline/request.rb

Constant Summary collapse

EMPTY_PARAMS =
[].freeze

Class Method Summary collapse

Class Method Details

.accept_result(req, result) ⇒ Object



183
184
185
186
187
188
189
190
191
192
# File 'lib/pg_pipeline/request.rb', line 183

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:



278
279
280
281
# File 'lib/pg_pipeline/request.rb', line 278

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



232
233
234
235
236
237
238
239
# File 'lib/pg_pipeline/request.rb', line 232

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



283
284
285
# File 'lib/pg_pipeline/request.rb', line 283

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

.finish!(req) ⇒ Object

Raises:



212
213
214
215
216
217
218
219
# File 'lib/pg_pipeline/request.rb', line 212

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

.immutable_values?(values) ⇒ Boolean

Returns:

  • (Boolean)


128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
# File 'lib/pg_pipeline/request.rb', line 128

def immutable_values?(values)
  index = 0
  size = values.size
  while index < size
    value = values[index]
    case value
    when Integer, Float, Symbol, NilClass, TrueClass, FalseClass
      nil
    when String
      return false unless value.frozen?
    else
      return false
    end
    index += 1
  end
  true
end

.park_waiter(req) ⇒ Object

Raises:



248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
# File 'lib/pg_pipeline/request.rb', line 248

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:



205
206
207
208
209
210
# File 'lib/pg_pipeline/request.rb', line 205

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



194
195
196
197
198
199
200
201
202
203
# File 'lib/pg_pipeline/request.rb', line 194

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



221
222
223
224
225
226
227
228
229
230
# File 'lib/pg_pipeline/request.rb', line 221

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)


159
160
161
162
163
164
# File 'lib/pg_pipeline/request.rb', line 159

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



166
167
168
169
170
171
172
173
# File 'lib/pg_pipeline/request.rb', line 166

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)


118
119
120
121
122
123
124
125
126
# File 'lib/pg_pipeline/request.rb', line 118

def snapshot_params(params)
  return EMPTY_PARAMS if params.nil?
  raise ArgumentError, "params must be an Array" unless params.is_a?(Array)
  return EMPTY_PARAMS if params.empty?

  return params if params.frozen? && immutable_values?(params)

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

.snapshot_sql(sql) ⇒ Object



111
112
113
114
# File 'lib/pg_pipeline/request.rb', line 111

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

.snapshot_value(value) ⇒ Object



146
147
148
149
150
151
152
153
154
155
156
157
# File 'lib/pg_pipeline/request.rb', line 146

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



175
176
177
178
179
180
181
# File 'lib/pg_pipeline/request.rb', line 175

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



241
242
243
244
245
246
# File 'lib/pg_pipeline/request.rb', line 241

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

  req.result
end

.wake_waiter(req) ⇒ Object



268
269
270
271
272
273
274
275
276
# File 'lib/pg_pipeline/request.rb', line 268

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