Class: Async::GRPC::XDS::Service::Stream

Inherits:
Object
  • Object
show all
Defined in:
lib/async/grpc/xds/service.rb

Overview

Represents one ADS stream and its subscribed resources.

Instance Method Summary collapse

Constructor Details

#initialize(control_plane, output) ⇒ Stream

Initialize an ADS stream.



75
76
77
78
79
80
81
82
# File 'lib/async/grpc/xds/service.rb', line 75

def initialize(control_plane, output)
	@control_plane = control_plane
	@output = output
	@subscriptions = Hash.new{|hash, type_url| hash[type_url] = Set.new}
	@versions = {}
	@queue = Async::Queue.new
	@closed = false
end

Instance Method Details

#changed(type_url) ⇒ Object

Schedule a resource type for delivery after it changes.



105
106
107
# File 'lib/async/grpc/xds/service.rb', line 105

def changed(type_url)
	@queue << type_url unless @closed
end

#closeObject

Close the stream and stop waiting for changes.



133
134
135
136
# File 'lib/async/grpc/xds/service.rb', line 133

def close
	@closed = true
	@queue.close
end

#flush(type_url) ⇒ Object

Deliver the latest resource version for a subscribed type.



120
121
122
123
124
125
126
127
128
129
130
# File 'lib/async/grpc/xds/service.rb', line 120

def flush(type_url)
	names = @subscriptions[type_url]
	return unless names
	
	version = @control_plane.version(type_url)
	return if @versions[type_url] == version
	
	response = @control_plane.response(type_url, names)
	@output.write(response)
	@versions[type_url] = version
end

#request(request) ⇒ Object

Process a discovery request and update the stream's subscriptions.



86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
# File 'lib/async/grpc/xds/service.rb', line 86

def request(request)
	return if request.type_url.nil? || request.type_url.empty?
	
	if request.error_detail
		Console.warn(self, "Received xDS NACK.", type_url: request.type_url, error_detail: request.error_detail)
		return
	end
	
	if request.resource_names.any?
		@subscriptions[request.type_url].merge(request.resource_names)
	else
		@subscriptions[request.type_url]
	end
	
	@queue << request.type_url
end

#runObject

Deliver scheduled resource updates until the stream closes.



111
112
113
114
115
116
# File 'lib/async/grpc/xds/service.rb', line 111

def run
	until @closed
		type_url = @queue.dequeue
		flush(type_url)
	end
end