Class: Async::GRPC::XDS::Stream
- Inherits:
-
Object
- Object
- Async::GRPC::XDS::Stream
- Defined in:
- lib/async/grpc/xds/stream.rb
Overview
Represents one discovery stream and its subscribed resources.
Instance Method Summary collapse
-
#changed(type_url) ⇒ Object
Schedule a resource type for delivery after it changes.
-
#close ⇒ Object
Close the stream and stop waiting for changes.
-
#flush(type_url) ⇒ Object
Deliver the latest resource version for a subscribed type.
-
#initialize(control_plane, output, resource_type: nil) ⇒ Stream
constructor
Initialize a discovery stream.
-
#request(request) ⇒ Object
Process a discovery request and update the stream's subscriptions.
-
#run ⇒ Object
Deliver scheduled resource updates until the stream closes.
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 |
#close ⇒ Object
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 |
#run ⇒ Object
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 |