Skip to content

Commit 06c7da9

Browse files
committed
v0.4.1: Integrate WebSocket with standard chat.ask interface
Add transport: :websocket support via with_params so chat.ask works over WebSocket without API changes. Add WebSocket#call for pre-built payloads. Fix token count preservation in WebSocket responses.
1 parent ebe213c commit 06c7da9

8 files changed

Lines changed: 206 additions & 75 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,17 @@ All notable changes to this project will be documented in this file.
55
The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/),
66
and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).
77

8+
## [0.4.1] - 2026-02-24
9+
10+
### Added
11+
12+
- `chat.with_params(transport: :websocket)` integration with standard `chat.ask` interface
13+
- `WebSocket#call` for accepting pre-built payloads from the provider
14+
15+
### Fixed
16+
17+
- WebSocket responses now preserve token counts from `StreamAccumulator`
18+
819
## [0.4.0] - 2026-02-24
920

1021
### Added

‎README.md‎

Lines changed: 18 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -269,59 +269,40 @@ Requires the `websocket-client-simple` gem:
269269
gem 'websocket-client-simple'
270270
```
271271

272-
### Basic usage
272+
### Usage
273273

274-
```ruby
275-
ws = RubyLLM::ResponsesAPI::WebSocket.new(api_key: ENV['OPENAI_API_KEY'])
276-
ws.connect
274+
Just add `transport: :websocket` to your params -- the standard `chat.ask` API works as-is:
277275

278-
# Stream a response
279-
message = ws.create_response(
280-
model: 'gpt-4o',
281-
input: [{ type: 'message', role: 'user', content: 'Hello!' }]
282-
) do |chunk|
283-
print chunk.content if chunk.content
284-
end
276+
```ruby
277+
chat = RubyLLM.chat(model: 'gpt-4o', provider: :openai_responses)
278+
chat.with_params(transport: :websocket)
285279

286-
puts "\n#{message.content}"
280+
chat.ask("Hello!")
281+
chat.ask("What's 2+2?") # reuses the same WebSocket connection
287282
```
288283

289-
### Multi-turn conversations
290-
291-
`previous_response_id` is tracked automatically across turns:
284+
Streaming works the same way:
292285

293286
```ruby
294-
ws.create_response(model: 'gpt-4o', input: [
295-
{ type: 'message', role: 'user', content: 'My name is Alice.' }
296-
])
297-
298-
ws.create_response(model: 'gpt-4o', input: [
299-
{ type: 'message', role: 'user', content: "What's my name?" }
300-
])
301-
# => "Alice" (auto-chained via previous_response_id)
287+
chat.ask("Tell me a story") { |chunk| print chunk.content }
302288
```
303289

304-
### With tools
290+
### Direct WebSocket access
291+
292+
For advanced use cases (raw Responses API format, warmup, explicit connection management):
305293

306294
```ruby
295+
ws = RubyLLM::ResponsesAPI::WebSocket.new(api_key: ENV['OPENAI_API_KEY'])
296+
ws.connect
297+
307298
ws.create_response(
308299
model: 'gpt-4o',
309-
input: [{ type: 'message', role: 'user', content: 'Search for Ruby 3.4 release notes' }],
310-
tools: [{ type: 'web_search_preview' }]
311-
)
312-
```
313-
314-
### Warmup
315-
316-
Pre-cache model weights without generating output:
300+
input: [{ type: 'message', role: 'user', content: 'Hello!' }]
301+
) { |chunk| print chunk.content }
317302

318-
```ruby
303+
# Pre-cache model weights
319304
ws.warmup(model: 'gpt-4o')
320-
```
321-
322-
### Cleanup
323305

324-
```ruby
325306
ws.disconnect
326307
```
327308

‎lib/ruby_llm/providers/openai_responses.rb‎

Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,16 @@ def api_base
1616
@config.openai_api_base || 'https://api.openai.com/v1'
1717
end
1818

19+
# Override to support WebSocket transport via with_params(transport: :websocket)
20+
def complete(messages, tools:, temperature:, model:, params: {}, headers: {}, schema: nil, thinking: nil, &block) # rubocop:disable Metrics/ParameterLists
21+
if params[:transport]&.to_sym == :websocket
22+
ws_complete(messages, tools: tools, temperature: temperature, model: model,
23+
params: params.except(:transport), schema: schema, thinking: thinking, &block)
24+
else
25+
super
26+
end
27+
end
28+
1929
def headers
2030
{
2131
'Authorization' => "Bearer #{@config.openai_api_key}",
@@ -137,6 +147,35 @@ def retrieve_container_file_content(container_id, file_id)
137147

138148
private
139149

150+
def ws_complete(messages, tools:, temperature:, model:, params:, schema:, thinking:, &block)
151+
normalized_temperature = maybe_normalize_temperature(temperature, model)
152+
153+
payload = Utils.deep_merge(
154+
render_payload(
155+
messages,
156+
tools: tools,
157+
temperature: normalized_temperature,
158+
model: model,
159+
stream: true,
160+
schema: schema,
161+
thinking: thinking
162+
),
163+
params
164+
)
165+
166+
ws_connection.connect unless ws_connection.connected?
167+
ws_connection.call(payload, &block)
168+
end
169+
170+
def ws_connection
171+
@ws_connection ||= WebSocket.new(
172+
api_key: @config.openai_api_key,
173+
api_base: api_base,
174+
organization_id: @config.openai_organization_id,
175+
project_id: @config.openai_project_id
176+
)
177+
end
178+
140179
# DELETE request via the underlying Faraday connection
141180
# RubyLLM::Connection only exposes get/post, so we use Faraday directly
142181
def delete_request(url)

‎lib/ruby_llm/providers/openai_responses/web_socket.rb‎

Lines changed: 40 additions & 36 deletions
Original file line numberDiff line numberDiff line change
@@ -11,14 +11,15 @@ class OpenAIResponses
1111
#
1212
# Requires the `websocket-client-simple` gem (soft dependency).
1313
#
14-
# Usage:
14+
# Integrated usage (recommended):
15+
# chat = RubyLLM.chat(model: 'gpt-4o', provider: :openai_responses)
16+
# chat.with_params(transport: :websocket)
17+
# chat.ask("Hello!")
18+
#
19+
# Standalone usage (advanced):
1520
# ws = RubyLLM::ResponsesAPI::WebSocket.new(api_key: ENV['OPENAI_API_KEY'])
1621
# ws.connect
17-
#
18-
# ws.create_response(model: 'gpt-4o', input: [{ type: 'message', role: 'user', content: 'Hi' }]) do |chunk|
19-
# print chunk.content if chunk.content
20-
# end
21-
#
22+
# ws.create_response(model: 'gpt-4o', input: [...]) { |chunk| ... }
2223
# ws.disconnect
2324
class WebSocket
2425
WEBSOCKET_PATH = '/v1/responses'
@@ -73,7 +74,6 @@ def connect(timeout: 10)
7374
end
7475
end
7576

76-
# Route all messages to the current queue (swapped per request)
7777
@ws.on(:message) do |msg|
7878
q = @mutex.synchronize { @message_queue }
7979
q&.push(msg.data)
@@ -89,35 +89,47 @@ def connect(timeout: 10)
8989
self
9090
end
9191

92-
# Send a response.create request and stream chunks via block.
93-
# @param model [String] model ID
94-
# @param input [Array<Hash>] input items in Responses API format
95-
# @param tools [Array<Hash>, nil] tool definitions
96-
# @param previous_response_id [String, nil] chain to a prior response
97-
# @param instructions [String, nil] system/developer instructions
98-
# @param extra [Hash] additional top-level fields forwarded to the API
92+
# Send a pre-built payload over WebSocket, streaming chunks via block.
93+
# This is the integration point for Provider#complete -- it accepts the
94+
# same payload hash that render_payload returns.
95+
#
96+
# @param payload [Hash] Responses API payload (model, input, tools, etc.)
9997
# @yield [RubyLLM::Chunk] each streamed chunk
10098
# @return [RubyLLM::Message] the assembled final message
101-
# @raise [ConcurrencyError] if another response is already in flight
102-
# @raise [ConnectionError] if not connected
103-
def create_response(model:, input:, tools: nil, previous_response_id: nil, instructions: nil, **extra, &block)
99+
def call(payload, &block)
104100
ensure_connected!
105101
acquire_flight!
106102

107103
queue = Queue.new
108104
@mutex.synchronize { @message_queue = queue }
109105

110-
payload = build_payload(
106+
envelope = { type: 'response.create', response: payload.except(:stream) }
107+
send_json(envelope)
108+
accumulate_response(queue, &block)
109+
ensure
110+
@mutex.synchronize { @message_queue = nil }
111+
release_flight!
112+
end
113+
114+
# Send a response.create request using raw Responses API format.
115+
# Useful for standalone usage outside the RubyLLM chat interface.
116+
#
117+
# @param model [String] model ID
118+
# @param input [Array<Hash>] input items in Responses API format
119+
# @param tools [Array<Hash>, nil] tool definitions
120+
# @param previous_response_id [String, nil] chain to a prior response
121+
# @param instructions [String, nil] system/developer instructions
122+
# @param extra [Hash] additional fields forwarded to the API
123+
# @yield [RubyLLM::Chunk] each streamed chunk
124+
# @return [RubyLLM::Message] the assembled final message
125+
def create_response(model:, input:, tools: nil, previous_response_id: nil, instructions: nil, **extra, &block)
126+
payload = build_standalone_payload(
111127
model: model, input: input, tools: tools,
112128
previous_response_id: previous_response_id,
113129
instructions: instructions, **extra
114130
)
115131

116-
send_json(payload)
117-
accumulate_response(queue, &block)
118-
ensure
119-
@mutex.synchronize { @message_queue = nil }
120-
release_flight!
132+
call(payload, &block)
121133
end
122134

123135
# Warm up the connection by sending a response.create with generate: false.
@@ -209,7 +221,7 @@ def build_headers
209221
headers
210222
end
211223

212-
def build_payload(model:, input:, tools: nil, previous_response_id: nil, instructions: nil, **extra)
224+
def build_standalone_payload(model:, input:, tools: nil, previous_response_id: nil, instructions: nil, **extra)
213225
prev_id = previous_response_id || @last_response_id
214226
response = { model: model, input: input }
215227
response[:tools] = tools.map { |t| Tools.tool_for(t) } if tools&.any?
@@ -220,7 +232,7 @@ def build_payload(model:, input:, tools: nil, previous_response_id: nil, instruc
220232
Compaction.apply_compaction(response, extra)
221233

222234
forwarded = extra.reject { |k, _| KNOWN_PARAMS.include?(k) }
223-
{ type: 'response.create', response: response.merge(forwarded) }
235+
response.merge(forwarded)
224236
end
225237

226238
def send_json(payload)
@@ -247,24 +259,16 @@ def accumulate_response(queue, &block)
247259
end
248260
end
249261

250-
build_final_message(accumulator)
262+
message = accumulator.to_message(nil)
263+
message.response_id = @last_response_id
264+
message
251265
end
252266

253267
def track_response_id(data)
254268
resp_id = data.dig('response', 'id')
255269
@mutex.synchronize { @last_response_id = resp_id } if resp_id
256270
end
257271

258-
def build_final_message(accumulator)
259-
Message.new(
260-
role: :assistant,
261-
content: accumulator.content,
262-
tool_calls: accumulator.tool_calls.empty? ? nil : accumulator.tool_calls,
263-
model_id: accumulator.model_id,
264-
response_id: @last_response_id
265-
)
266-
end
267-
268272
def ensure_connected!
269273
raise ConnectionError, 'WebSocket is not connected. Call #connect first.' unless connected?
270274
end

‎lib/rubyllm_responses_api.rb‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -37,7 +37,7 @@
3737
module RubyLLM
3838
# ResponsesAPI namespace for direct access to helpers and version
3939
module ResponsesAPI
40-
VERSION = '0.4.0'
40+
VERSION = '0.4.1'
4141

4242
# Shorthand access to built-in tool helpers
4343
BuiltInTools = Providers::OpenAIResponses::BuiltInTools

‎ruby_llm-responses_api.gemspec‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@
22

33
Gem::Specification.new do |spec|
44
spec.name = 'ruby_llm-responses_api'
5-
spec.version = '0.4.0'
5+
spec.version = '0.4.1'
66
spec.authors = ['Chris Hasinski']
77
spec.email = ['krzysztof.hasinski@gmail.com']
88

‎spec/ruby_llm/providers/openai_responses/web_socket_spec.rb‎

Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -279,6 +279,46 @@ def last_sent_payload
279279
end
280280
end
281281

282+
describe '#call' do
283+
before { connect_ws! }
284+
285+
it 'accepts a pre-built payload and wraps in WS envelope' do
286+
payload = { model: 'gpt-4o', input: [{ type: 'message', role: 'user', content: 'Hi' }] }
287+
288+
with_events(standard_events) do
289+
ws.call(payload)
290+
end
291+
292+
sent = last_sent_payload
293+
expect(sent['type']).to eq('response.create')
294+
expect(sent['response']['model']).to eq('gpt-4o')
295+
end
296+
297+
it 'strips :stream from the payload' do
298+
payload = { model: 'gpt-4o', input: [], stream: true }
299+
300+
with_events(standard_events) do
301+
ws.call(payload)
302+
end
303+
304+
expect(last_sent_payload['response']).not_to have_key('stream')
305+
end
306+
307+
it 'yields chunks and returns assembled Message' do
308+
payload = { model: 'gpt-4o', input: [] }
309+
chunks = []
310+
311+
message = with_events(standard_events(text_deltas: %w[Hello world])) do
312+
ws.call(payload) { |chunk| chunks << chunk }
313+
end
314+
315+
expect(chunks.select(&:content).map(&:content)).to eq(%w[Hello world])
316+
expect(message).to be_a(RubyLLM::Message)
317+
expect(message.content).to eq('Helloworld')
318+
expect(message.response_id).to eq('resp_ws_test')
319+
end
320+
end
321+
282322
describe '#warmup' do
283323
before { connect_ws! }
284324

0 commit comments

Comments
 (0)