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

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

Overview

Represents one discovery stream and its subscribed resources.

Instance Method Summary collapse

Constructor Details

#initialize(control_plane, output, resource_type: nil) ⇒ Stream

Initialize a discovery stream.



20
21
22
23
24
25
26
27
28
# File 'lib/async/grpc/xds/stream.rb', line 20

def initialize(control_plane, output, resource_type: nil)
	@control_plane = control_plane
	@output = output
	@resource_type = resource_type
	@subscriptions = {}
	@versions = {}
	@queue = Async::Queue.new
	@closed = false
end

Instance Method Details

#changed(type_url) ⇒ Object

Schedule a resource type for delivery after it changes.



64
65
66
67
68
69
# File 'lib/async/grpc/xds/stream.rb', line 64

def changed(type_url)
	return if @resource_type && type_url != @resource_type
	return unless @subscriptions.key?(type_url)
	
	@queue << type_url unless @closed
end

#closeObject

Close the stream and stop waiting for changes.



95
96
97
98
# File 'lib/async/grpc/xds/stream.rb', line 95

def close
	@closed = true
	@queue.close
end

#flush(type_url) ⇒ Object

Deliver the latest resource version for a subscribed type.



82
83
84
85
86
87
88
89
90
91
92
# File 'lib/async/grpc/xds/stream.rb', line 82

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.



32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
# File 'lib/async/grpc/xds/stream.rb', line 32

def request(request)
	type_url = request.type_url
	
	if @resource_type
		if type_url.nil? || type_url.empty?
			type_url = @resource_type
		elsif type_url != @resource_type
			raise Protocol::GRPC::Error.new(
				Protocol::GRPC::Status::INVALID_ARGUMENT,
				"Expected resource type #{@resource_type.inspect}, but received #{type_url.inspect}."
			)
		end
	elsif type_url.nil? || type_url.empty?
		return
	end
	
	if request.error_detail
		Console.warn(self, "Received xDS NACK.", type_url: type_url, error_detail: request.error_detail)
		return
	end
	
	names = Set.new(request.resource_names)
	if @subscriptions[type_url] != names
		@subscriptions[type_url] = names
		@versions.delete(type_url)
	end
	
	@queue << type_url
end

#runObject

Deliver scheduled resource updates until the stream closes.



73
74
75
76
77
78
# File 'lib/async/grpc/xds/stream.rb', line 73

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