class Stream
Represents one discovery stream and its subscribed resources.
Definitions
def initialize(control_plane, output, resource_type: nil)
Initialize a discovery stream.
Signature
-
parameter
control_planeControlPlane The control plane that provides resources.
-
parameter
outputInterface(:write) The discovery response stream.
-
parameter
resource_typeString | Nil The fixed resource type, or
nilfor aggregated discovery.
Implementation
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
def request(request)
Process a discovery request and update the stream's subscriptions.
Signature
-
parameter
requestEnvoy::Service::Discovery::V3::DiscoveryRequest The discovery request.
Implementation
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
def changed(type_url)
Schedule a resource type for delivery after it changes.
Signature
-
parameter
type_urlString The changed xDS resource type URL.
Implementation
def changed(type_url)
return if @resource_type && type_url != @resource_type
return unless @subscriptions.key?(type_url)
@queue << type_url unless @closed
end
def run
- asynchronous
Deliver scheduled resource updates until the stream closes.
Signature
- asynchronous
Implementation
def run
until @closed
type_url = @queue.dequeue
flush(type_url)
end
end
def flush(type_url)
Deliver the latest resource version for a subscribed type.
Signature
-
parameter
type_urlString The xDS resource type URL.
Implementation
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
def close
Close the stream and stop waiting for changes.
Implementation
def close
@closed = true
@queue.close
end