Class: Async::GRPC::XDS::DiscoveryClient

Inherits:
Object
  • Object
show all
Includes:
ADSStream::Delegate
Defined in:
lib/async/grpc/xds/discovery_client.rb

Overview

Client for xDS APIs (ADS or individual APIs) Implements Aggregated Discovery Service (ADS) protocol Acts as delegate for ADSStream, receiving discovery_response events

Constant Summary collapse

LISTENER_TYPE =

xDS API type URLs (v3 API)

"type.googleapis.com/envoy.config.listener.v3.Listener"
ROUTE_TYPE =
"type.googleapis.com/envoy.config.route.v3.RouteConfiguration"
CLUSTER_TYPE =
Cluster::TYPE_URL
ENDPOINT_TYPE =
Endpoint::TYPE_URL
SECRET_TYPE =
"type.googleapis.com/envoy.extensions.transport_sockets.tls.v3.Secret"

Instance Method Summary collapse

Constructor Details

#initialize(server_config, node: nil) ⇒ DiscoveryClient

Initialize xDS discovery client



40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
# File 'lib/async/grpc/xds/discovery_client.rb', line 40

def initialize(server_config, node: nil)
	@server_uri = server_config[:server_uri]
	@channel_creds = server_config[:channel_creds]
	@server_features = server_config[:server_features] || []
	@node_info = node || build_node_info
	@node = build_node_proto(@node_info)
	@grpc_client = nil
	@versions = {}  # Track version_info per type_url
	@nonces = {}     # Track nonces per type_url
	@mutex = Mutex.new
	@subscriptions = {}  # Track subscriptions by type_url
	@stream_task = nil
	@ads_stream = nil  # ADSStream instance when connected (owns stream state)
	@stream_ready_promise = nil  # Resolved when stream_opened runs
end

Instance Method Details

#closeObject

Close xDS discovery client



91
92
93
94
95
96
97
98
99
100
101
# File 'lib/async/grpc/xds/discovery_client.rb', line 91

def close
	@mutex.synchronize do
		@stream_task&.stop
		@grpc_client&.close
		@grpc_client = nil
		@subscriptions.clear
		@stream_task = nil
		@ads_stream = nil
		@stream_ready_promise = nil
	end
end

#discovery_response(response, stream) ⇒ Object

Process a discovery response received by an ADS stream.



174
175
176
# File 'lib/async/grpc/xds/discovery_client.rb', line 174

def discovery_response(response, stream)
	process_response(response, stream)
end

#stream_closed(stream) ⇒ Object

Record that an ADS stream has closed.



167
168
169
# File 'lib/async/grpc/xds/discovery_client.rb', line 167

def stream_closed(stream)
	@mutex.synchronize{@ads_stream = nil}
end

#stream_opened(stream) ⇒ Object

Record that an ADS stream has opened and unblock pending subscriptions.



160
161
162
163
# File 'lib/async/grpc/xds/discovery_client.rb', line 160

def stream_opened(stream)
	@mutex.synchronize{@ads_stream = stream}
	@stream_ready_promise&.resolve(stream)
end

#subscribe(type_url, resource_names, &block) ⇒ Object

Subscribe to resource type using ADS (Aggregated Discovery Service - single stream for all types)



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
# File 'lib/async/grpc/xds/discovery_client.rb', line 62

def subscribe(type_url, resource_names, &block)
	# Store subscription callback
	@mutex.synchronize do
		@subscriptions[type_url] = {
			resource_names: resource_names,
			callback: block
		}
	end
	
	# Ensure ADS stream is running
	ensure_stream_running
	
	# Wait for stream to be ready (event-driven, no polling)
	promise = @stream_ready_promise
	if promise && !promise.completed?
		begin
			promise.wait(timeout: 5)
		rescue Async::TimeoutError
			# Stream didn't open in time; send_discovery_request will no-op if @ads_stream is nil
		end
	end
	
	send_discovery_request(type_url, resource_names) if @ads_stream
	
	# Return the stream task (already running)
	@stream_task
end