Module: PgPipeline::RequestOps
- Defined in:
- lib/pg_pipeline/request.rb
Class Method Summary collapse
- .accept_result(req, result) ⇒ Object
- .assert_result_slot!(req) ⇒ Object
- .cancel!(req) ⇒ Object
- .clear_result(result) ⇒ Object
- .finish!(req) ⇒ Object
- .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
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
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
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
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
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
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
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 |