Module: PgPipeline::RequestOps

Defined in:
lib/pg_pipeline/request.rb

Class Method Summary collapse

Class Method Details

.accept_result(req, result) ⇒ Object



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

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:



189
190
191
192
# File 'lib/pg_pipeline/request.rb', line 189

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



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

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



194
195
196
# File 'lib/pg_pipeline/request.rb', line 194

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

.finish!(req) ⇒ Object

Raises:



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

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
  req.condition.signal unless req.cancelled
end

.query_boundary!(req) ⇒ Object

Raises:



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

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



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

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



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

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

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

.snapshot_name(name) ⇒ Object

Raises:

  • (ArgumentError)


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

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



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

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)


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

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



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

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

.snapshot_value(value) ⇒ Object



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

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



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

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



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

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

  req.result
end