Class: OpenAI::Helpers::Streaming::ResponseStreamState
- Inherits:
-
Object
- Object
- OpenAI::Helpers::Streaming::ResponseStreamState
- Defined in:
- lib/openai/helpers/streaming/response_stream.rb
Instance Attribute Summary collapse
-
#completed_response ⇒ Object
readonly
Returns the value of attribute completed_response.
Instance Method Summary collapse
- #accumulate_event(event:, current_snapshot:) ⇒ Object
- #handle_event(event, text_only: false) ⇒ Object
-
#initialize(text_format:, starting_after: nil) ⇒ ResponseStreamState
constructor
A new instance of ResponseStreamState.
Constructor Details
#initialize(text_format:, starting_after: nil) ⇒ ResponseStreamState
93 94 95 96 97 98 99 |
# File 'lib/openai/helpers/streaming/response_stream.rb', line 93 def initialize(text_format:, starting_after: nil) @current_snapshot = nil @completed_response = nil @text_format = text_format @resumed = !starting_after.nil? @unexposed_buffers = {}.compare_by_identity end |
Instance Attribute Details
#completed_response ⇒ Object (readonly)
Returns the value of attribute completed_response.
91 92 93 |
# File 'lib/openai/helpers/streaming/response_stream.rb', line 91 def completed_response @completed_response end |
Instance Method Details
#accumulate_event(event:, current_snapshot:) ⇒ Object
192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 |
# File 'lib/openai/helpers/streaming/response_stream.rb', line 192 def accumulate_event(event:, current_snapshot:) if current_snapshot.nil? if event.is_a?(OpenAI::Models::Responses::ResponseCreatedEvent) return isolated_value(event.response) end unless @resumed raise "Expected first event to be response.created" end if event.is_a?(OpenAI::Models::Responses::ResponseCompletedEvent) @completed_response = event.response return event.response end return nil end case event when OpenAI::Models::Responses::ResponseOutputItemAddedEvent current_snapshot.output.push(isolated_value(event.item)) when OpenAI::Models::Responses::ResponseContentPartAddedEvent output = current_snapshot.output[event.output_index] if output.is_a?(OpenAI::Models::Responses::ResponseOutputMessage) output.content.push(isolated_value(event.part)) current_snapshot.output[event.output_index] = output end when OpenAI::Models::Responses::ResponseTextDeltaEvent output = current_snapshot.output[event.output_index] if output.is_a?(OpenAI::Models::Responses::ResponseOutputMessage) content = output.content[event.content_index] if content.is_a?(OpenAI::Models::Responses::ResponseOutputText) content.text = append_delta(content.text, event.delta) output.content[event.content_index] = content current_snapshot.output[event.output_index] = output end end when OpenAI::Models::Responses::ResponseFunctionCallArgumentsDeltaEvent output = current_snapshot.output[event.output_index] if output.is_a?(OpenAI::Models::Responses::ResponseFunctionToolCall) output.arguments = append_delta(output.arguments || "", event.delta) current_snapshot.output[event.output_index] = output end when OpenAI::Models::Responses::ResponseCompletedEvent @completed_response = event.response end current_snapshot end |
#handle_event(event, text_only: false) ⇒ Object
101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 |
# File 'lib/openai/helpers/streaming/response_stream.rb', line 101 def handle_event(event, text_only: false) return [event] if event.is_a?(UnknownStreamEvent) @current_snapshot = accumulate_event( event: event, current_snapshot: @current_snapshot ) events_to_yield = [] case event when OpenAI::Models::Responses::ResponseTextDeltaEvent return [event] if text_only snapshot = nil if @current_snapshot output = @current_snapshot.output[event.output_index] if output.is_a?(OpenAI::Models::Responses::ResponseOutputMessage) content = output.content[event.content_index] snapshot = content.text if content.is_a?(OpenAI::Models::Responses::ResponseOutputText) end end @unexposed_buffers.delete(snapshot) # A server-directed resumed stream or an unknown future snapshot value # has no complete prefix from which to build a truthful snapshot. events_to_yield << OpenAI::Streaming::ResponseTextDeltaEvent.new(event.to_h.merge(snapshot: snapshot)) when OpenAI::Models::Responses::ResponseTextDoneEvent text = if @current_snapshot output = @current_snapshot.output[event.output_index] if output.is_a?(OpenAI::Models::Responses::ResponseOutputMessage) content = output.content[event.content_index] content.text if content.is_a?(OpenAI::Models::Responses::ResponseOutputText) end else event.text end parsed = parse_structured_text(text) events_to_yield << OpenAI::Streaming::ResponseTextDoneEvent.new( content_index: event.content_index, item_id: event.item_id, output_index: event.output_index, sequence_number: event.sequence_number, text: event.text, type: event.type, **event.to_h.slice(:logprobs), parsed: parsed ) when OpenAI::Models::Responses::ResponseFunctionCallArgumentsDeltaEvent return [event] if text_only snapshot = nil if @current_snapshot output = @current_snapshot.output[event.output_index] if output.is_a?(OpenAI::Models::Responses::ResponseFunctionToolCall) snapshot = output.arguments end end @unexposed_buffers.delete(snapshot) # See the text-delta branch above: a partial server resume or an # unknown future output item has no truthful accumulated prefix. events_to_yield << OpenAI::Streaming::ResponseFunctionCallArgumentsDeltaEvent.new( event.to_h.merge(snapshot: snapshot) ) when OpenAI::Models::Responses::ResponseCompletedEvent events_to_yield << OpenAI::Streaming::ResponseCompletedEvent.new( sequence_number: event.sequence_number, type: event.type, response: event.response ) else # Pass through other events unchanged. events_to_yield << event end events_to_yield end |