Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 31 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,37 @@ sideband connection. See the [Realtime WebSocket guide](realtime.md) and the run
[text, transcription, voice, and sideband examples](examples/realtime/README.md)
for lifecycle, authentication, proxy, TLS, and custom transport details.

### Responses WebSockets

The Responses API also supports a persistent, block-scoped WebSocket for
sequential and multiplexed response creation. Add the optional
`async-websocket` gem, then send `response.create` events through
`client.responses.connect`:

```ruby
client.responses.connect do |connection|
connection.response.create(
model: "gpt-5.2",
input: "Say hello.",
stream_id: "turn_1"
)

connection.each do |event|
print(event.delta) if event.type == :"response.output_text.delta"
break if event.type == :"response.completed"
end
end
```

The connection is intentionally single-owner: callers serialize writes and
use one reader (`receive` or `each`). The SDK forwards response fields and
`stream_id` values without imposing additional client-side policy. Known
server events are decoded into generated models on a best-effort basis, while
newer event types remain observable as `UnknownServerEvent` values. The SDK
does not automatically reconnect or replay an ambiguous write. When the server
closes a connection, open a new connection and continue with
`previous_response_id` when the response was stored.

### Pagination

List methods in the OpenAI API are paginated.
Expand Down
2 changes: 2 additions & 0 deletions lib/openai.rb
Original file line number Diff line number Diff line change
Expand Up @@ -1315,4 +1315,6 @@
require_relative "openai/helpers/streaming/chat_events"
require_relative "openai/helpers/streaming/chat_completion_stream"
require_relative "openai/streaming"
require_relative "openai/helpers/websocket"
require_relative "openai/helpers/realtime"
require_relative "openai/helpers/responses_websocket"
158 changes: 37 additions & 121 deletions lib/openai/helpers/realtime/client_extension.rb
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,8 @@ module Helpers
module Realtime
# Realtime request integration kept outside the generated client implementation.
module ClientExtension
include OpenAI::WebSocket::ClientRequest

MAX_CALL_CLEANUP_SECONDS = 5.0
private_constant :MAX_CALL_CLEANUP_SECONDS

Expand Down Expand Up @@ -140,28 +142,20 @@ def realtime_connection_request(path:, query:, websocket_base_url: nil, options:
#
# @api private
def with_realtime_connection_request(path:, query:, websocket_base_url: nil, options: nil)
request, deadline = build_realtime_connection_request(
path: path,
query: query,
websocket_base_url: websocket_base_url,
options: options
)
handshake_completed = false
mark_handshake_completed = -> { handshake_completed = true }
yield(request, mark_handshake_completed)
rescue OpenAI::Errors::RealtimeConnectionError => e
raise if handshake_completed
raise unless e.http_status == 401 && @workload_identity_auth
build = lambda do |deadline|
build_realtime_connection_request(
path: path,
query: query,
websocket_base_url: websocket_base_url,
options: options,
deadline: deadline
)
end

@workload_identity_auth.invalidate_token
refreshed, = build_realtime_connection_request(
path: path,
query: query,
websocket_base_url: websocket_base_url,
options: options,
deadline: deadline
)
yield(refreshed, mark_handshake_completed)
with_websocket_connection_retry(
error_class: OpenAI::Errors::RealtimeConnectionError,
build: build
) { |request, marker| yield(request, marker) }
end

private def build_realtime_connection_request(
Expand All @@ -171,6 +165,28 @@ def with_realtime_connection_request(path:, query:, websocket_base_url: nil, opt
options:,
deadline: nil
)
build_shared_websocket_connection_request(
path: path,
query: query,
websocket_base_url: websocket_base_url,
options: options,
deadline: deadline,
validate: -> (value) { validate_realtime_websocket_request!(value) },
invalid_base_url_message: "`websocket_base_url` must be an absolute HTTP or WebSocket URL " \
"without credentials, query, or fragment",
malformed_base_url_message: "`websocket_base_url` is not a valid URL",
preserve_base_url_cause: true,
extra_query_message: "`request_options[:extra_query]` is not supported for Realtime WebSocket " \
"connections; omit it",
max_retries_message: "`request_options[:max_retries]` is not supported for Realtime WebSocket " \
"connections; use 0 or omit it",
timeout_error: lambda do |url, cause|
OpenAI::Errors::RealtimeConnectionError.new(url: url, cause: cause)
end
)
end

private def validate_realtime_websocket_request!(websocket_base_url)
if x509_identity?(@copy_options.fetch(:workload_identity))
raise OpenAI::Errors::Error, "X.509 workload identity does not support Realtime WebSocket connections"
end
Expand All @@ -184,106 +200,6 @@ def with_realtime_connection_request(path:, query:, websocket_base_url: nil, opt
if websocket_base_url && @provider_runtime
raise ArgumentError, "`websocket_base_url` cannot be combined with `provider`"
end

websocket_uri = parse_websocket_base_url(websocket_base_url)

opts = options.to_h.dup
OpenAI::RequestOptions.validate!(opts)
extra_query = opts.delete(:extra_query)
unless extra_query.nil? || (extra_query.respond_to?(:empty?) && extra_query.empty?)
message = "`request_options[:extra_query]` is not supported for Realtime WebSocket " \
"connections; omit it"
raise ArgumentError, message
end

max_retries = opts[:max_retries]
unless max_retries.nil? || max_retries == 0
message = "`request_options[:max_retries]` is not supported for Realtime WebSocket " \
"connections; use 0 or omit it"
raise ArgumentError, message
end

request = build_request(
{
method: :get,
path: path,
query: query,
security: {bearer_auth: true}
},
opts
)
error_request = if websocket_uri
with_websocket_base_url(request, path: path, base_url: websocket_uri)
else
request
end

error_url = websocket_url(error_request.fetch(:url))

if @workload_identity_auth
deadline ||= request[:timeout]&.then do |timeout|
OpenAI::Internal::Util.monotonic_secs + timeout
end
end

workload_identity_header = "Bearer #{OpenAI::Client::WORKLOAD_IDENTITY_API_KEY_PLACEHOLDER}"
if @workload_identity_auth && request.fetch(:headers)["authorization"] == workload_identity_header
token = @workload_identity_auth.get_token(deadline: deadline)
request = request.merge(
headers: request.fetch(:headers).merge("authorization" => "Bearer #{token}")
)
end

request = prepare_request(request, redirect_count: 0, retry_count: 0)
request = with_websocket_base_url(request, path: path, base_url: websocket_uri) if websocket_uri

url = websocket_url(request.fetch(:url))
headers = request.fetch(:headers).except("accept", "content-type").reject do |name, _value|
name.to_s.casecmp?("proxy-authorization")
end

request = request.merge(url: url, headers: headers)
request = request_with_remaining_timeout(request, deadline) unless deadline.nil?
[request, deadline]
rescue Timeout::Error => e
raise(
OpenAI::Errors::RealtimeConnectionError.new(
url: error_url,
cause: e
)
)
end

private def with_websocket_base_url(request, path:, base_url:)
url = OpenAI::Internal::Util.join_parsed_uri(
OpenAI::Internal::Util.parse_uri(base_url.to_s),
{path: OpenAI::Internal::Util.interpolate_path(path)}
)
url.query = request.fetch(:url).query
request.merge(url: url)
end

private def parse_websocket_base_url(value)
return if value.nil?

uri = URI(value.to_s)
valid_scheme = %w[http https ws wss].include?(uri.scheme)
ambiguous_component = uri.userinfo || uri.query || uri.fragment
unless uri.absolute? && uri.host && valid_scheme && !ambiguous_component
message = "`websocket_base_url` must be an absolute HTTP or WebSocket URL " \
"without credentials, query, or fragment"
raise ArgumentError, message
end

uri
rescue URI::Error => e
raise ArgumentError, "`websocket_base_url` is not a valid URL", cause: e
end

private def websocket_url(url)
url = url.dup
url.scheme = {"http" => "ws", "https" => "wss"}.fetch(url.scheme, url.scheme)
url
end
end
end
Expand Down
103 changes: 13 additions & 90 deletions lib/openai/helpers/realtime/connection.rb
Original file line number Diff line number Diff line change
Expand Up @@ -3,8 +3,8 @@
module OpenAI
module Realtime
# A live, typed Realtime WebSocket connection.
class Connection
include Enumerable
class Connection < OpenAI::WebSocket::Connection
include OpenAI::WebSocket::Protocol

# @return [OpenAI::Realtime::ConnectionResources::Session]
attr_reader :session
Expand All @@ -20,8 +20,7 @@ class Connection

# @api private
def initialize(socket:, url:)
@socket = socket
@url = url
super
@server_event_names = discriminator_values(OpenAI::Realtime::RealtimeServerEvent)
@client_event_names = discriminator_values(OpenAI::Realtime::RealtimeClientEvent)
@session = OpenAI::Realtime::ConnectionResources::Session.new(self)
Expand All @@ -30,34 +29,6 @@ def initialize(socket:, url:)
@input_audio_buffer = OpenAI::Realtime::ConnectionResources::InputAudioBuffer.new(self)
end

# @return [URI::Generic]
attr_reader :url

# Yield server events until the remote peer closes the connection.
def each
return enum_for(__method__) unless block_given?

while (event = receive)
yield(event)
end

self
end

# Receive and parse the next server event, or return nil after a clean close.
def receive
data = receive_raw
return nil if data.nil?

parse_event(data)
end

# Receive the next raw WebSocket message.
def receive_raw
message = @socket.read
message&.to_str
end

# Parse raw JSON as a typed server event. Valid events that are newer than this
# SDK remain observable as {UnknownServerEvent} values.
def parse_event(data)
Expand Down Expand Up @@ -115,65 +86,9 @@ def send_event(event)
raise ArgumentError.new("Invalid Realtime client event."), cause: e
end

# Send an already encoded text message.
def send_raw(data)
if closed?
raise(
OpenAI::Errors::RealtimeConnectionError.new(
url: @url,
message: "Cannot send on a closed Realtime WebSocket."
)
)
end

text = data.dup
text.force_encoding(Encoding::UTF_8) if text.encoding == Encoding::BINARY
text = text.encode(Encoding::UTF_8) unless text.encoding == Encoding::UTF_8
unless text.valid_encoding?
raise ArgumentError, "Realtime WebSocket text must contain valid UTF-8"
end

@socket.write(text)
nil
end

# Close the connection.
def close(code: 1000, reason: "")
return if closed?

@socket.close(code: code, reason: reason)
nil
end

# Abort without waiting for the WebSocket close handshake.
#
# @api private
def abort
return if closed?

@socket.abort
nil
end

# @return [Boolean]
def closed? = @socket.closed?

private def discriminator_values(union)
union.variants.to_h do |variant|
value = variant.fields.fetch(:type).fetch(:const)
[value.to_s, true]
end
end

private def event_type(event)
unless event.is_a?(Hash)
raise ArgumentError, "Realtime server event must be a JSON object"
end

type = event[:type]
return type if type.is_a?(String) || type.is_a?(Symbol)

raise ArgumentError, "Realtime server event type must be a string or symbol"
super(event, message: "Realtime server event must be a JSON object") unless event.is_a?(Hash)
super(event, message: "Realtime server event type must be a string or symbol")
end

private def validate_discriminator!(event, allowed, kind:)
Expand All @@ -194,6 +109,14 @@ def closed? = @socket.closed?

ArgumentError.new("Realtime event is missing required fields or contains invalid values")
end

private def connection_error(message)
OpenAI::Errors::RealtimeConnectionError.new(url: @url, message: message)
end

private def closed_send_message = "Cannot send on a closed Realtime WebSocket."

private def invalid_text_message = "Realtime WebSocket text must contain valid UTF-8"
end
end
end
Loading