Module: PgPipeline::RequestOps
- Defined in:
- lib/pg_pipeline/request.rb
Constant Summary collapse
- EMPTY_PARAMS =
[].freeze
Class Method Summary collapse
- .accept_result(req, result) ⇒ Object
- .assert_result_slot!(req) ⇒ Object
- .cancel!(req) ⇒ Object
- .clear_result(result) ⇒ Object
- .finish!(req) ⇒ Object
- .immutable_values?(values) ⇒ Boolean
- .park_waiter(req) ⇒ Object
- .query_boundary!(req) ⇒ Object
- .record_error!(req, error, result) ⇒ Object
- .reject!(req, error) ⇒ Object
- .snapshot_name(name) ⇒ Object
- .snapshot_param_types(param_types) ⇒ Object
- .snapshot_params(params) ⇒ Object
- .snapshot_sql(sql) ⇒ Object
- .snapshot_value(value) ⇒ Object
- .transition!(req, from, to) ⇒ Object
- .wait(req) ⇒ Object
- .wake_waiter(req) ⇒ Object
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
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
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
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
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
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
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
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 |