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
- .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
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
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
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
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
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
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 |