Class: Crystil::Wrappers::OpenAI

Inherits:
Object
  • Object
show all
Defined in:
lib/crystil/wrappers/openai.rb,
sig/crystil/wrappers/openai.rbs

Overview

Wrapper for OpenAI Ruby client (ruby-openai gem)

Instance Method Summary collapse

Constructor Details

#initialize(config, collector, sentinel = nil) ⇒ OpenAI

Returns a new instance of OpenAI.

Parameters:



9
10
11
12
13
# File 'lib/crystil/wrappers/openai.rb', line 9

def initialize(config, collector, sentinel = nil)
  @config = config
  @collector = collector
  @sentinel = sentinel
end

Instance Method Details

#register(client) ⇒ Object

Parameters:

  • client (Object)

Returns:

  • (Object)


15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
# File 'lib/crystil/wrappers/openai.rb', line 15

def register(client)
  validate_client!(client)

  # Prevent double registration
  return client if client.instance_variable_defined?(:@crystil_registered)

  # Store references in client instance
  client.instance_variable_set(:@crystil_config, @config)
  client.instance_variable_set(:@crystil_collector, @collector)
  client.instance_variable_set(:@crystil_sentinel, @sentinel)
  client.instance_variable_set(:@crystil_registered, true)

  # Wrap the chat method
  wrap_chat_method(client)

  # Wrap responses.create if the client supports the Responses API (ruby-openai >= 8.0)
  wrap_responses_method(client) if client.respond_to?(:responses)

  client
end

#validate_client!(client) ⇒ void

This method returns an undefined value.

Parameters:

  • client (Object)

Raises:



160
161
162
163
164
165
166
167
168
# File 'lib/crystil/wrappers/openai.rb', line 160

def validate_client!(client)
  # When register an OpenAI client, it currently only checks if there
  # is the classic "chat" interface.  OpenAI now supports the newer
  # "responses" interface.  So, if an OpenAI client later only has
  # the newer "responses" interface, this validate_client will fail.
  return if client.respond_to?(:chat)

  raise RegistrationError, "Client does not appear to be a valid OpenAI client (missing chat method)"
end

#wrap_chat_method(client) ⇒ void

This method returns an undefined value.

Parameters:

  • client (Object)


170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
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
245
246
# File 'lib/crystil/wrappers/openai.rb', line 170

def wrap_chat_method(client)
  client.singleton_class.class_eval do
    include Base

    alias_method :original_chat, :chat

    define_method(:chat) do |parameters: {}|
      start_time = Time.now
      version = defined?(::OpenAI::VERSION) ? ::OpenAI::VERSION : nil

      sentinel = instance_variable_get(:@crystil_sentinel)
      sentinel&.raise_if_irrelevant!(
        title: OPENAI_CLIENT_TITLE,
        request: parameters,
        version: version
      )

      # Detect streaming and wrap callback to accumulate chunks
      streaming = parameters[:stream].is_a?(Proc)
      accumulated_response = {} if streaming

      if streaming
        # Configure streaming to include usage (matches Python SDK behavior).
        # A caller-set value (true or false) wins — we only force `true`
        # when the key is absent.
        parameters[:stream_options] ||= {}
        parameters[:stream_options][:include_usage] = true unless parameters[:stream_options].key?(:include_usage)

        user_callback = parameters[:stream]
        parameters[:stream] = proc do |chunk, bytesize|
          # Normalize chunk to match Python SDK format (add missing delta keys with nil)
          normalized_chunk = crystil_normalize_openai_chunk(chunk)
          # Accumulate chunk (merge into accumulated_response)
          crystil_merge_streaming_chunk(accumulated_response, normalized_chunk)
          # Call user's original callback with original chunk (don't modify user's data)
          user_callback.call(chunk, bytesize)
        end
      end

      # Call the original method
      response = original_chat(parameters: parameters)

      # Use accumulated response for streaming, otherwise use returned response
      final_response = streaming ? accumulated_response : response

      # Submit analytics
      crystil_submit_analytics(
        method: :chat,
        args: [],
        kwargs: parameters,
        response: final_response,
        start_time: start_time,
        end_time: Time.now,
        title: OPENAI_CLIENT_TITLE,
        version: version
      )

      response
    rescue CrystilRequestInterceptedError => e
      # We don't want to send intercepts to collector
      raise e
    rescue StandardError => e
      crystil_submit_error_analytics(
        method: :chat,
        args: [],
        kwargs: parameters,
        error: e,
        start_time: start_time,
        end_time: Time.now,
        title: OPENAI_CLIENT_TITLE,
        version: version
      )

      raise e
    end
  end
end

#wrap_responses_method(client) ⇒ void

This method returns an undefined value.

Parameters:

  • client (Object)


38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
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
# File 'lib/crystil/wrappers/openai.rb', line 38

def wrap_responses_method(client)
  responses_obj = client.responses

  # Copy Crystil references onto the responses sub-object so Base helpers can access them
  responses_obj.instance_variable_set(:@crystil_config, client.instance_variable_get(:@crystil_config))
  responses_obj.instance_variable_set(:@crystil_collector, client.instance_variable_get(:@crystil_collector))
  responses_obj.instance_variable_set(:@crystil_sentinel, client.instance_variable_get(:@crystil_sentinel))

  responses_obj.singleton_class.class_eval do
    include Base

    alias_method :original_create, :create

    define_method(:create) do |parameters: {}|
      start_time = Time.now
      version = defined?(::OpenAI::VERSION) ? ::OpenAI::VERSION : nil

      sentinel = instance_variable_get(:@crystil_sentinel)
      sentinel&.raise_if_irrelevant!(
        title: OPENAI_CLIENT_TITLE,
        request: parameters,
        version: version
      )

      # Dup before mutating so we don't replace the caller's :stream key with our
      # internal wrapping proc — the caller's hash should be unchanged after the call.
      parameters = parameters.dup

      streaming = parameters[:stream].respond_to?(:call)

      if streaming
        accumulated_events = []
        user_callback = parameters[:stream]
        # Determine how many arguments to forward to the user's callback.
        # .arity returns negative values for methods with optional/splat params
        # (e.g. def call(*args) => -1, def call(a, b=nil) => -2), so .abs gives
        # us the minimum required argument count. For fixed-arity procs/lambdas
        # this is exact. Note: a variadic callable (def call(*args)) has arity -1,
        # so .abs yields 1 — the event_type argument will be silently dropped for
        # that signature. Use def call(event, event_type = nil) to receive both.
        user_callback_arity =
          case user_callback
          when Proc
            user_callback.arity.abs
          else
            user_callback.method(:call).arity.abs
          end

        parameters[:stream] = proc do |chunk, event_type|
          accumulated_events << chunk if chunk.is_a?(Hash)
          user_callback.call(*[chunk, event_type].first(user_callback_arity))
        end
      end

      response = original_create(parameters: parameters)

      # Default to the returned response; for streaming, replace it with the
      # full response from the terminal response.completed event.
      final_response = response

      if streaming
        error_event = accumulated_events.find { |e| e["type"] == "error" }
        if error_event
          final_response = {
            "status" => "failed",
            "error" => error_event["error"]
          }
        end

        # Extract the full response from the terminal response.completed event.
        # A missing terminal event means the stream was interrupted.
        completed = accumulated_events.find { |e| e["type"] == "response.completed" }
        unless completed || error_event
          raise StreamError,
                "Responses API stream ended without response.completed event"
        end

        final_response = completed["response"] if completed && !error_event
      end

      # Treat only an explicit failed response status as failed analytics.
      # Other non-exceptional statuses are still recorded as succeeded.
      response_status = final_response.is_a?(Hash) ? final_response["status"] : nil
      analytics_status = response_status == "failed" ? "failed" : "succeeded"
      analytics_exception =
        if analytics_status == "failed"
          final_response.dig("error", "message") ||
            "OpenAI response status: #{response_status || "unknown"}"
        end

      crystil_submit_analytics(
        method: :create,
        args: [],
        kwargs: parameters,
        response: final_response,
        start_time: start_time,
        end_time: Time.now,
        title: OPENAI_CLIENT_TITLE,
        version: version,
        status: analytics_status,
        exception: analytics_exception
      )

      response
    rescue CrystilRequestInterceptedError => e
      raise e
    rescue StandardError => e
      crystil_submit_error_analytics(
        method: :create,
        args: [],
        kwargs: parameters,
        error: e,
        start_time: start_time,
        end_time: Time.now,
        title: OPENAI_CLIENT_TITLE,
        version: version
      )
      raise e
    end
  end
end