Skip to content

Commit d4bbcbf

Browse files
authored
Merge pull request #1108 from code0-tech/882-add-opentelemetry-support
Add OpenTelemetry support
2 parents dd97a9c + cedb303 commit d4bbcbf

12 files changed

Lines changed: 755 additions & 21 deletions

File tree

Gemfile

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -95,3 +95,10 @@ gem 'json-schema', '~> 6.0'
9595
gem 'triangulum', '0.26.2'
9696

9797
gem 'benchmark'
98+
99+
# OpenTelemetry
100+
gem 'opentelemetry-exporter-otlp', '~> 0.34.0' # we need this to get traces
101+
gem 'opentelemetry-exporter-otlp-logs', '~> 0.5.1' # we need this to get logs
102+
gem 'opentelemetry-exporter-otlp-metrics', '~> 0.10.0' # we need this to get metrics
103+
gem 'opentelemetry-instrumentation-all', '~> 0.94.0' # we need this for logs and traces
104+
gem 'opentelemetry-instrumentation-logger', '~> 0.4.0' # we need this to get logs

Gemfile.lock

Lines changed: 213 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -212,6 +212,214 @@ GEM
212212
mini_portile2 (~> 2.8.2)
213213
racc (~> 1.4)
214214
open3 (0.2.1)
215+
opentelemetry-api (1.10.0)
216+
logger
217+
opentelemetry-common (0.25.0)
218+
opentelemetry-api (~> 1.0)
219+
opentelemetry-exporter-otlp (0.34.0)
220+
google-protobuf (>= 3.18)
221+
googleapis-common-protos-types (~> 1.3)
222+
opentelemetry-api (~> 1.1)
223+
opentelemetry-common (~> 0.20)
224+
opentelemetry-sdk (~> 1.10)
225+
opentelemetry-semantic_conventions
226+
opentelemetry-exporter-otlp-logs (0.5.1)
227+
google-protobuf (>= 3.18)
228+
googleapis-common-protos-types (~> 1.3)
229+
opentelemetry-api (~> 1.1)
230+
opentelemetry-common (~> 0.20)
231+
opentelemetry-logs-api (~> 0.1)
232+
opentelemetry-logs-sdk (~> 0.1)
233+
opentelemetry-sdk
234+
opentelemetry-semantic_conventions
235+
opentelemetry-exporter-otlp-metrics (0.10.0)
236+
google-protobuf (>= 3.18, < 5.0)
237+
googleapis-common-protos-types (~> 1.3)
238+
opentelemetry-api (~> 1.1)
239+
opentelemetry-common (~> 0.20)
240+
opentelemetry-metrics-api (~> 0.2)
241+
opentelemetry-metrics-sdk (~> 0.5)
242+
opentelemetry-sdk (~> 1.2)
243+
opentelemetry-semantic_conventions
244+
opentelemetry-helpers-mysql (0.6.0)
245+
opentelemetry-api (~> 1.7)
246+
opentelemetry-common (~> 0.21)
247+
opentelemetry-helpers-sql (0.4.0)
248+
opentelemetry-api (~> 1.7)
249+
opentelemetry-helpers-sql-processor (0.5.0)
250+
opentelemetry-api (~> 1.0)
251+
opentelemetry-common (~> 0.21)
252+
opentelemetry-instrumentation-action_mailer (0.8.1)
253+
opentelemetry-instrumentation-active_support (~> 0.10)
254+
opentelemetry-instrumentation-action_pack (0.18.0)
255+
opentelemetry-instrumentation-rack (~> 0.29)
256+
opentelemetry-instrumentation-action_view (0.13.0)
257+
opentelemetry-instrumentation-active_support (~> 0.10)
258+
opentelemetry-instrumentation-active_job (0.13.0)
259+
opentelemetry-instrumentation-base (~> 0.25)
260+
opentelemetry-instrumentation-active_model_serializers (0.25.0)
261+
opentelemetry-instrumentation-active_support (>= 0.7.0)
262+
opentelemetry-instrumentation-active_record (0.13.0)
263+
opentelemetry-instrumentation-base (~> 0.25)
264+
opentelemetry-instrumentation-active_storage (0.5.1)
265+
opentelemetry-instrumentation-active_support (~> 0.10)
266+
opentelemetry-instrumentation-active_support (0.12.0)
267+
opentelemetry-instrumentation-base (~> 0.25)
268+
opentelemetry-instrumentation-all (0.94.0)
269+
opentelemetry-instrumentation-active_model_serializers (~> 0.25.0)
270+
opentelemetry-instrumentation-anthropic (~> 0.5.0)
271+
opentelemetry-instrumentation-aws_lambda (~> 0.7.0)
272+
opentelemetry-instrumentation-aws_sdk (~> 0.12.0)
273+
opentelemetry-instrumentation-bunny (~> 0.25.0)
274+
opentelemetry-instrumentation-concurrent_ruby (~> 0.25.0)
275+
opentelemetry-instrumentation-dalli (~> 0.30.0)
276+
opentelemetry-instrumentation-delayed_job (~> 0.26.0)
277+
opentelemetry-instrumentation-ethon (~> 0.29.0)
278+
opentelemetry-instrumentation-excon (~> 0.29.0)
279+
opentelemetry-instrumentation-faraday (~> 0.33.0)
280+
opentelemetry-instrumentation-grape (~> 0.7.0)
281+
opentelemetry-instrumentation-graphql (~> 0.32.0)
282+
opentelemetry-instrumentation-grpc (~> 0.5.0)
283+
opentelemetry-instrumentation-gruf (~> 0.6.0)
284+
opentelemetry-instrumentation-http (~> 0.30.0)
285+
opentelemetry-instrumentation-http_client (~> 0.29.0)
286+
opentelemetry-instrumentation-httpx (~> 0.8.0)
287+
opentelemetry-instrumentation-koala (~> 0.24.0)
288+
opentelemetry-instrumentation-lmdb (~> 0.26.0)
289+
opentelemetry-instrumentation-mongo (~> 0.26.0)
290+
opentelemetry-instrumentation-mysql2 (~> 0.34.0)
291+
opentelemetry-instrumentation-net_http (~> 0.29.0)
292+
opentelemetry-instrumentation-pg (~> 0.36.0)
293+
opentelemetry-instrumentation-que (~> 0.13.0)
294+
opentelemetry-instrumentation-racecar (~> 0.7.0)
295+
opentelemetry-instrumentation-rack (~> 0.31.0)
296+
opentelemetry-instrumentation-rails (~> 0.42.0)
297+
opentelemetry-instrumentation-rake (~> 0.6.0)
298+
opentelemetry-instrumentation-rdkafka (~> 0.10.0)
299+
opentelemetry-instrumentation-redis (~> 0.29.0)
300+
opentelemetry-instrumentation-resque (~> 0.9.0)
301+
opentelemetry-instrumentation-restclient (~> 0.28.0)
302+
opentelemetry-instrumentation-ruby_kafka (~> 0.25.0)
303+
opentelemetry-instrumentation-sidekiq (~> 0.29.0)
304+
opentelemetry-instrumentation-sinatra (~> 0.30.0)
305+
opentelemetry-instrumentation-trilogy (~> 0.69.0)
306+
opentelemetry-instrumentation-anthropic (0.5.0)
307+
opentelemetry-instrumentation-base (~> 0.25)
308+
opentelemetry-instrumentation-aws_lambda (0.7.0)
309+
opentelemetry-instrumentation-base (~> 0.25)
310+
opentelemetry-instrumentation-aws_sdk (0.12.0)
311+
opentelemetry-instrumentation-base (~> 0.25)
312+
opentelemetry-instrumentation-base (0.26.1)
313+
opentelemetry-api (~> 1.7)
314+
opentelemetry-common (~> 0.21)
315+
opentelemetry-registry (~> 0.1)
316+
opentelemetry-instrumentation-bunny (0.25.0)
317+
opentelemetry-instrumentation-base (~> 0.25)
318+
opentelemetry-instrumentation-concurrent_ruby (0.25.0)
319+
opentelemetry-instrumentation-base (~> 0.25)
320+
opentelemetry-instrumentation-dalli (0.30.0)
321+
opentelemetry-instrumentation-base (~> 0.25)
322+
opentelemetry-instrumentation-delayed_job (0.26.0)
323+
opentelemetry-instrumentation-base (~> 0.25)
324+
opentelemetry-instrumentation-ethon (0.29.0)
325+
opentelemetry-instrumentation-base (~> 0.25)
326+
opentelemetry-instrumentation-excon (0.29.1)
327+
opentelemetry-instrumentation-base (~> 0.25)
328+
opentelemetry-instrumentation-faraday (0.33.0)
329+
opentelemetry-instrumentation-base (~> 0.25)
330+
opentelemetry-instrumentation-grape (0.7.0)
331+
opentelemetry-instrumentation-rack (~> 0.29)
332+
opentelemetry-instrumentation-graphql (0.32.0)
333+
opentelemetry-instrumentation-base (~> 0.25)
334+
opentelemetry-instrumentation-grpc (0.5.1)
335+
opentelemetry-instrumentation-base (~> 0.25)
336+
opentelemetry-instrumentation-gruf (0.6.1)
337+
opentelemetry-instrumentation-base (~> 0.25)
338+
opentelemetry-instrumentation-http (0.30.0)
339+
opentelemetry-instrumentation-base (~> 0.25)
340+
opentelemetry-instrumentation-http_client (0.29.0)
341+
opentelemetry-instrumentation-base (~> 0.25)
342+
opentelemetry-instrumentation-httpx (0.8.0)
343+
opentelemetry-instrumentation-base (~> 0.25)
344+
opentelemetry-instrumentation-koala (0.24.0)
345+
opentelemetry-instrumentation-base (~> 0.25)
346+
opentelemetry-instrumentation-lmdb (0.26.0)
347+
opentelemetry-instrumentation-base (~> 0.25)
348+
opentelemetry-instrumentation-logger (0.4.0)
349+
opentelemetry-instrumentation-base (~> 0.25)
350+
opentelemetry-logs-api (~> 0.1)
351+
opentelemetry-instrumentation-mongo (0.26.0)
352+
opentelemetry-instrumentation-base (~> 0.25)
353+
opentelemetry-instrumentation-mysql2 (0.34.0)
354+
opentelemetry-helpers-mysql
355+
opentelemetry-helpers-sql
356+
opentelemetry-helpers-sql-processor
357+
opentelemetry-instrumentation-base (~> 0.25)
358+
opentelemetry-instrumentation-net_http (0.29.0)
359+
opentelemetry-instrumentation-base (~> 0.25)
360+
opentelemetry-instrumentation-pg (0.36.0)
361+
opentelemetry-helpers-sql
362+
opentelemetry-helpers-sql-processor
363+
opentelemetry-instrumentation-base (~> 0.25)
364+
opentelemetry-instrumentation-que (0.13.0)
365+
opentelemetry-instrumentation-base (~> 0.25)
366+
opentelemetry-instrumentation-racecar (0.7.0)
367+
opentelemetry-instrumentation-base (~> 0.25)
368+
opentelemetry-instrumentation-rack (0.31.1)
369+
opentelemetry-instrumentation-base (~> 0.25)
370+
opentelemetry-instrumentation-rails (0.42.0)
371+
opentelemetry-instrumentation-action_mailer (~> 0.7)
372+
opentelemetry-instrumentation-action_pack (~> 0.17)
373+
opentelemetry-instrumentation-action_view (~> 0.12)
374+
opentelemetry-instrumentation-active_job (~> 0.11)
375+
opentelemetry-instrumentation-active_record (~> 0.12)
376+
opentelemetry-instrumentation-active_storage (~> 0.4)
377+
opentelemetry-instrumentation-active_support (~> 0.11)
378+
opentelemetry-instrumentation-concurrent_ruby (~> 0.25)
379+
opentelemetry-instrumentation-rake (0.6.0)
380+
opentelemetry-instrumentation-base (~> 0.25)
381+
opentelemetry-instrumentation-rdkafka (0.10.0)
382+
opentelemetry-instrumentation-base (~> 0.25)
383+
opentelemetry-instrumentation-redis (0.29.0)
384+
opentelemetry-instrumentation-base (~> 0.25)
385+
opentelemetry-instrumentation-resque (0.9.0)
386+
opentelemetry-instrumentation-base (~> 0.25)
387+
opentelemetry-instrumentation-restclient (0.28.0)
388+
opentelemetry-instrumentation-base (~> 0.25)
389+
opentelemetry-instrumentation-ruby_kafka (0.25.0)
390+
opentelemetry-instrumentation-base (~> 0.25)
391+
opentelemetry-instrumentation-sidekiq (0.29.0)
392+
opentelemetry-instrumentation-base (~> 0.25)
393+
opentelemetry-instrumentation-sinatra (0.30.0)
394+
opentelemetry-instrumentation-rack (~> 0.29)
395+
opentelemetry-instrumentation-trilogy (0.69.0)
396+
opentelemetry-helpers-mysql
397+
opentelemetry-helpers-sql
398+
opentelemetry-helpers-sql-processor
399+
opentelemetry-instrumentation-base (~> 0.25)
400+
opentelemetry-semantic_conventions (>= 1.8.0)
401+
opentelemetry-logs-api (0.4.0)
402+
opentelemetry-api (~> 1.0)
403+
opentelemetry-logs-sdk (0.6.0)
404+
opentelemetry-api (~> 1.2)
405+
opentelemetry-logs-api (~> 0.1)
406+
opentelemetry-sdk (~> 1.3)
407+
opentelemetry-metrics-api (0.6.0)
408+
opentelemetry-api (~> 1.0)
409+
opentelemetry-metrics-sdk (0.15.0)
410+
opentelemetry-api (~> 1.1)
411+
opentelemetry-metrics-api (~> 0.2)
412+
opentelemetry-sdk (~> 1.2)
413+
opentelemetry-registry (0.6.0)
414+
opentelemetry-api (~> 1.1)
415+
opentelemetry-sdk (1.12.0)
416+
logger
417+
opentelemetry-api (~> 1.1)
418+
opentelemetry-common (~> 0.20)
419+
opentelemetry-registry (~> 0.2)
420+
opentelemetry-semantic_conventions
421+
opentelemetry-semantic_conventions (1.41.0)
422+
opentelemetry-api (~> 1.0)
215423
parallel (1.28.0)
216424
parser (3.3.11.1)
217425
ast (~> 2.4.1)
@@ -429,6 +637,11 @@ DEPENDENCIES
429637
image_processing (>= 1.2)
430638
json-schema (~> 6.0)
431639
lograge (~> 0.14.0)
640+
opentelemetry-exporter-otlp (~> 0.34.0)
641+
opentelemetry-exporter-otlp-logs (~> 0.5.1)
642+
opentelemetry-exporter-otlp-metrics (~> 0.10.0)
643+
opentelemetry-instrumentation-all (~> 0.94.0)
644+
opentelemetry-instrumentation-logger (~> 0.4.0)
432645
pg (~> 1.1)
433646
pry (~> 0.16.0)
434647
pry-byebug (~> 3.10)

app/grpc/concerns/grpc_stream_handler.rb

Lines changed: 54 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,12 @@ module GrpcStreamHandler
44
include Code0::ZeroTrack::Loggable
55
extend ActiveSupport::Concern
66

7+
StreamItem = Struct.new(:data, :otel_context, keyword_init: true)
8+
9+
def self.tracer
10+
@tracer ||= ::OpenTelemetry.tracer_provider.tracer('sagittarius-grpc-stream')
11+
end
12+
713
class_methods do
814
def grpc_stream(method)
915
define_method(method) do |_, call|
@@ -17,17 +23,23 @@ def grpc_stream(method)
1723
grpc_object: grpc_object)
1824

1925
encoded_data = send('encoders')[method].call(grpc_object)
20-
encoded_data64 = Base64.encode64(encoded_data)
26+
encoded_data64 = Base64.encode64(encoded_data).delete("\n")
2127

2228
GrpcStreamHandler.logger.info(message: 'Encoded data', runtime_id: runtime_id, method: method,
2329
encoded_data: encoded_data64)
2430

31+
carrier = {}
32+
OpenTelemetry.propagation.inject(carrier)
33+
trace_context64 = Base64.encode64(carrier.to_json).delete("\n")
34+
35+
notification_payload = "#{self},#{method},#{runtime_id},#{encoded_data64},#{trace_context64}"
36+
2537
ActiveRecord::Base.connection.raw_connection
26-
.exec("NOTIFY grpc_streams, '#{self},#{method},#{runtime_id},#{encoded_data64}'")
38+
.exec("NOTIFY grpc_streams, '#{notification_payload}'")
2739
end
2840
define_singleton_method("end_#{method}") do |runtime_id|
2941
ActiveRecord::Base.connection.raw_connection
30-
.exec("NOTIFY grpc_streams, '#{self},#{method},#{runtime_id},end'")
42+
.exec("NOTIFY grpc_streams, '#{self},#{method},#{runtime_id},end,'")
3143
end
3244
end
3345
end
@@ -48,11 +60,18 @@ def self.listen!
4860

4961
conn.wait_for_notify(1) do |_, _, payload|
5062
logger.info(message: 'Received notification', payload: payload)
51-
class_name, method_name, runtime_id, encoded_data64 = payload.split(',')
63+
parts = payload.split(',')
64+
class_name = parts[0]
65+
method_name = parts[1]
66+
runtime_id = parts[2]
67+
encoded_data64 = parts[3]
68+
trace_context64 = parts[4]
5269

5370
clazz = class_name.constantize
5471
method_name = method_name.to_sym
5572

73+
otel_context = extract_otel_context(trace_context64)
74+
5675
if encoded_data64 == 'end'
5776
decoded_data = :end
5877
else
@@ -62,7 +81,7 @@ def self.listen!
6281

6382
queues = GrpcStreamHandler.yielders.dig(clazz, method_name, runtime_id.to_i)
6483
queues&.each do |queue|
65-
queue << decoded_data
84+
queue << StreamItem.new(data: decoded_data, otel_context: otel_context)
6685
rescue StandardError => e
6786
logger.error(message: 'Error while yielding data', error: e.message, backtrace: e.backtrace)
6887
end
@@ -73,6 +92,15 @@ def self.listen!
7392
end
7493
end
7594

95+
def self.extract_otel_context(trace_context64)
96+
return nil if trace_context64.blank?
97+
98+
carrier = JSON.parse(Base64.decode64(trace_context64))
99+
OpenTelemetry.propagation.extract(carrier)
100+
rescue StandardError
101+
nil
102+
end
103+
76104
def create_enumerator(clazz, method, runtime_id, _call)
77105
logger.debug(message: 'Creating enumerator', runtime_id: runtime_id, clazz: clazz, method: method)
78106

@@ -82,7 +110,7 @@ def create_enumerator(clazz, method, runtime_id, _call)
82110
method_queues = queues[method] ||= {}
83111
runtime_queues = method_queues[runtime_id] ||= []
84112

85-
runtime_queues.each { |existing_queue| existing_queue << :end }
113+
runtime_queues.each { |existing_queue| existing_queue << StreamItem.new(data: :end, otel_context: nil) }
86114
runtime_queues.clear
87115
runtime_queues << queue
88116

@@ -92,19 +120,27 @@ def create_enumerator(clazz, method, runtime_id, _call)
92120
loop do
93121
item = queue.pop(timeout: 1)
94122
next if item.nil?
95-
break if item == :end
96-
97-
begin
98-
y << item
99-
ApplicationRecord.connection_pool.with_connection do
100-
Runtime.update(runtime_id, last_heartbeat: Time.zone.now)
123+
break if item.data == :end
124+
125+
otel_context = item.otel_context || OpenTelemetry::Context.current
126+
OpenTelemetry::Context.with_current(otel_context) do
127+
GrpcStreamHandler.tracer.in_span("#{clazz}/#{method} send") do |span|
128+
y << item.data
129+
ApplicationRecord.connection_pool.with_connection do
130+
Runtime.update(runtime_id, last_heartbeat: Time.zone.now)
131+
end
132+
rescue ActiveRecord::ActiveRecordError => e
133+
logger.warn(message: 'Failed to update runtime heartbeat', exception: e.message, backtrace: e.backtrace)
134+
rescue GRPC::Core::CallError => e
135+
logger.info(message: 'Stream was closed from client side (probably)')
136+
137+
span.set_attribute('rpc.response.status_code', 'UNKNOWN')
138+
span.set_attribute('error.type', e.class.name)
139+
span.status = ::OpenTelemetry::Trace::Status.error(e.message)
140+
span.record_exception(e)
141+
142+
raise
101143
end
102-
rescue ActiveRecord::ActiveRecordError => e
103-
logger.warn(message: 'Failed to update runtime heartbeat', exception: e.message, backtrace: e.backtrace)
104-
rescue GRPC::Core::CallError
105-
logger.info(message: 'Stream was closed from client side (probably)')
106-
107-
raise
108144
end
109145
end
110146
ensure

app/grpc/execution_handler.rb

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -38,7 +38,7 @@ def test(requests, call)
3838
rescue StandardError => e
3939
logger.error(message: 'Error reading execution stream', error: e.message,
4040
backtrace: e.backtrace)
41-
outbound_queue << :end if outbound_queue
41+
outbound_queue << GrpcStreamHandler::StreamItem.new(data: :end, otel_context: nil) if outbound_queue
4242
ensure
4343
logger.info(message: 'Execution runtime request stream closed')
4444
end

0 commit comments

Comments
 (0)