diff --git a/lib/syskit/deployment.rb b/lib/syskit/deployment.rb index 975188965..5cffac5fc 100644 --- a/lib/syskit/deployment.rb +++ b/lib/syskit/deployment.rb @@ -40,6 +40,9 @@ class Deployment < ::Roby::Task # rubocop:disable Metrics/ClassLength # The underlying process object attr_reader :orocos_process + # The time after the latest task reconfiguration + attr_reader :latest_configuration_time + # An object describing the underlying pocess server # # @return [RobyApp::Configuration::ProcessServerConfig] @@ -500,7 +503,7 @@ def resolve_remote_task_handles(remote_tasks) RemoteTaskHandles = Struct.new( :handle, :state_reader, :state_getter, :default_properties, :configuring, :current_configuration, :needs_reconfiguration, - :in_fatal, :quarantined + :in_fatal, :quarantined, :configured_since ) # @api private @@ -522,6 +525,15 @@ def start_configuration(orocos_name) remote_task_handles[orocos_name].configuring = true end + # @api private + # + # Update a task handle's configured_since attribute + def update_configured_since(orocos_name) + t = Time.now + remote_task_handles[orocos_name].configured_since = t + @latest_configuration_time = t + end + # @api private # # Declare that the given task is being configured diff --git a/lib/syskit/interface/commands.rb b/lib/syskit/interface/commands.rb index a2bf5336c..14eb7d3da 100644 --- a/lib/syskit/interface/commands.rb +++ b/lib/syskit/interface/commands.rb @@ -19,7 +19,7 @@ def deployments # Return incremental update about deployments # # @return [Protocol::Deployment] - def poll_ready_deployments(known: []) + def poll_ready_deployments(known: [], reconfigured_since: nil) deployments = plan.find_tasks(Syskit::Deployment).running.find_all(&:ready?) deployment_ids = deployments.map { _1.droby_id.id } @@ -27,10 +27,25 @@ def poll_ready_deployments(known: []) deployments.find_all { !known.include?(_1.droby_id.id) } removed_deployments = known.find_all { |id| !deployment_ids.include?(id) } - [new_deployments, removed_deployments] + + if reconfigured_since + updated_deployments = + compute_reconfigured_deployments(known, reconfigured_since) + [new_deployments, removed_deployments, updated_deployments] + else + [new_deployments, removed_deployments] + end end command :poll_ready_deployments, - "incremental information about deployments" + "incremental information about deployments and deployed tasks" + + def compute_reconfigured_deployments(known_ids, reconfigured_since) + deployments.find_all do |d| + known_ids.include?(d.droby_id.id) && + (configuration_time = d.latest_configuration_time) && + configuration_time > reconfigured_since + end + end # Save the configuration of all running tasks of the given model to disk # diff --git a/lib/syskit/interface/v2/protocol.rb b/lib/syskit/interface/v2/protocol.rb index 8603ad3ed..6bb439f21 100644 --- a/lib/syskit/interface/v2/protocol.rb +++ b/lib/syskit/interface/v2/protocol.rb @@ -36,7 +36,7 @@ def pretty_print(pp) end DeployedTask = Struct.new( - :name, :ior, :orogen_model_name, keyword_init: true + :name, :ior, :orogen_model_name, :configured_since, keyword_init: true ) def self.register_marshallers(protocol) @@ -68,7 +68,8 @@ def self.marshal_remote_task_handle(name, remote_task_handle) ior = remote_task_handle.handle.ior model_name = remote_task_handle.handle.model.name DeployedTask.new( - name: name, ior: ior, orogen_model_name: model_name + name: name, ior: ior, orogen_model_name: model_name, + configured_since: remote_task_handle.configured_since ) end diff --git a/lib/syskit/task_context.rb b/lib/syskit/task_context.rb index f5b6d14fc..47c9aa7b1 100644 --- a/lib/syskit/task_context.rb +++ b/lib/syskit/task_context.rb @@ -1100,10 +1100,16 @@ def perform_setup(promise) log_remote_call(:configure, orocos_name) do orocos_task.configure(false) end + true else info "#{self} was already configured" + false end end + + promise.on_success do |called_configure| + execution_agent.update_configured_since(orocos_name) if called_configure + end end # (see Component#setting_up!)_ diff --git a/lib/syskit/telemetry/async.rb b/lib/syskit/telemetry/async.rb index 6156cb18f..61053a65e 100644 --- a/lib/syskit/telemetry/async.rb +++ b/lib/syskit/telemetry/async.rb @@ -1,5 +1,7 @@ # frozen_string_literal: true +require "syskit/interface/v2/protocol" + require "syskit/telemetry/async/main_thread_restrictions" require "syskit/telemetry/async/name_service" diff --git a/lib/syskit/telemetry/async/name_service.rb b/lib/syskit/telemetry/async/name_service.rb index 16611ec75..bf0d87997 100644 --- a/lib/syskit/telemetry/async/name_service.rb +++ b/lib/syskit/telemetry/async/name_service.rb @@ -90,9 +90,12 @@ def async_update_tasks(tasks) def queue_new_tasks_discovery(tasks) tasks.each do |t| next if @discovery[t.name] - next if t.ior == @registered_tasks[t.name]&.identity - async_discover_task(t) + async_task = @registered_tasks[t.name] + next if t.ior == async_task&.identity && + t.configured_since == async_task&.configured_since + + async_discover_task(t, async_task) end end @@ -121,9 +124,16 @@ def update_from_result(port_read_manager:) self.ior = ior return unless discovered - self.async_task = TaskContext.from_discovered_interface( - discovered, port_read_manager: port_read_manager - ) + async_task = self.async_task + if async_task && async_task.identity == ior + async_task.update_from_discovered_interface(discovered) + else + async_task = TaskContext.create_from_discovered_interface( + discovered, port_read_manager: port_read_manager + ) + end + async_task.configured_since = task.configured_since + self.async_task = async_task end def wait @@ -138,7 +148,7 @@ def resolved? # @api private # # Create a future that discovers a remote task - def async_discover_task(task) + def async_discover_task(task, async_task) ensure_in_main_thread future = Concurrent::Promises.future_on(@discovery_executor) do @@ -148,7 +158,9 @@ def async_discover_task(task) # set while the future was pending discover_task(task.name, ior, task.orogen_model_name) if ior end - @discovery[task.name] = AsyncDiscovery.new(task: task, future: future) + @discovery[task.name] = + AsyncDiscovery.new(task: task, future: future, + async_task: async_task) end def wait_and_resolve_all_pending_discoveries @@ -243,7 +255,7 @@ def finished_discovery_validate_ior(async_discovery) # started processing. Throw away the resolved task and start # again async_discovery.async_task&.dispose - async_discover_task(async_discovery.task) + async_discover_task(async_discovery.task, async_discovery.async_task) false end diff --git a/lib/syskit/telemetry/async/task_context.rb b/lib/syskit/telemetry/async/task_context.rb index 6d27fd0a7..2db364bbb 100644 --- a/lib/syskit/telemetry/async/task_context.rb +++ b/lib/syskit/telemetry/async/task_context.rb @@ -45,6 +45,12 @@ class TaskContext < TaskContextHooks # will be considered the same from the perspective of a hash key attr_reader :hash + # The time of the last configuration of this task + # + # This is used to trigger a discovery of the task's interface, since in + # syskit the interface may only change during reconfiguration + attr_accessor :configured_since + def states_index_to_symbols return @states_index_to_symbols if @states_index_to_symbols @@ -88,19 +94,23 @@ def self.discover_interface(name, ior, orogen_model) # Create a TaskContext object from the information returned by # {.discover} - def self.from_discovered_interface(discovered, port_read_manager:) + def self.create_from_discovered_interface(discovered, port_read_manager:) async_task = new( discovered.task.name, model: discovered.orogen_model, port_read_manager: port_read_manager ) - async_task.discover_attributes(discovered.attributes) - async_task.discover_properties(discovered.properties) - async_task.discover_ports(discovered.ports) + async_task.update_from_discovered_interface(discovered) async_task.reachable!(discovered.task) async_task end + def update_from_discovered_interface(discovered) + update_attributes(discovered.attributes) + update_properties(discovered.properties) + update_ports(discovered.ports) + end + # Synchronously discover remote task info and return the corresponding # {TaskContext} object # @@ -111,7 +121,7 @@ def self.discover(name, ior, orogen_model, port_read_manager:) discovered = Orocos.allow_blocking_calls do TaskContext.discover_interface(name, ior, orogen_model) end - TaskContext.from_discovered_interface( + TaskContext.create_from_discovered_interface( discovered, port_read_manager: port_read_manager ) end @@ -264,46 +274,61 @@ def port(name) @ports.fetch(name) end - def discover_attributes(raw_attributes) - @attributes = - raw_attributes.each_with_object({}) do |p, h| - async = Attribute.new(self, p.name, p.type) - async.reachable!(p) - h[p.name] = async - end + def update_attributes(raw_attributes) + update_interface_objects( + raw_attributes, @attributes, + :on_attribute_reachable, :on_attribute_unreachable + ) do |a| + Attribute.new(self, a.name, a.type) + end + end - @attributes.each_value { run_hook :on_attribute_reachable, _1 } + def update_properties(raw_properties) + update_interface_objects( + raw_properties, @properties, + :on_property_reachable, :on_property_unreachable + ) do |p| + Property.new(self, p.name, p.type) + end end - def discover_properties(raw_properties) - @properties = - raw_properties.each_with_object({}) do |p, h| - async = Property.new(self, p.name, p.type) - async.reachable!(p) - h[p.name] = async + def update_ports(raw_ports) + update_interface_objects( + raw_ports, @ports, + :on_port_reachable, :on_port_unreachable + ) do |p| + case p + when Orocos::InputPort + InputPort.new(self, p.name, p.type) + else + OutputPort.new( + self, p.name, p.type, @port_read_manager + ) end - - @properties.each_value { run_hook :on_property_reachable, _1 } + end end - def discover_ports(raw_ports) - @ports = - raw_ports.each_with_object({}) do |p, h| - async = - case p - when Orocos::InputPort - InputPort.new(self, p.name, p.type) - else - OutputPort.new( - self, p.name, p.type, @port_read_manager - ) - end + def update_interface_objects(new, actual, reachable_hook, + unreachable_hook) + new_names = Set.new + new.each do |obj| + name = obj.name + new_names << name + next if actual.key?(name) + + async = yield(obj) + async.reachable!(obj) + actual[name] = async + run_hook reachable_hook, name + end - async.reachable!(p) - h[p.name] = async - end + actual.delete_if do |k, v| + next if new_names.include?(k) - @ports.each_value { run_hook :on_port_reachable, _1 } + v.unreachable! + run_hook unreachable_hook, k + true + end end def dispose diff --git a/lib/syskit/telemetry/ui/runtime_state.rb b/lib/syskit/telemetry/ui/runtime_state.rb index d77f896e8..56c721b84 100644 --- a/lib/syskit/telemetry/ui/runtime_state.rb +++ b/lib/syskit/telemetry/ui/runtime_state.rb @@ -600,20 +600,30 @@ def process_current_deployments def query_deployment_update polling_call( ["syskit"], "poll_ready_deployments", - known: @current_deployments.map(&:id) - ) do |updated, removed| - update_current_deployments(updated, removed) + known: @current_deployments.map(&:id), + reconfigured_since: latest_configuration_time || Time.at(0) + ) do |new, removed, updated| + update_current_deployments(new, removed, updated || []) process_current_deployments end end - def update_current_deployments(updated, removed) + def update_current_deployments(new, removed, updated) + updated_ids = updated.map(&:id) @current_deployments.delete_if do |d| - removed.include?(d.id) + removed.include?(d.id) || updated_ids.include?(d.id) end + @current_deployments.concat(new) @current_deployments.concat(updated) end + def latest_configuration_time + times = @current_deployments.flat_map do + _1.deployed_tasks.map(&:configured_since) + end + times.compact.max + end + def reset_current_deployments @current_deployments = [] reset_task_inspector diff --git a/test/interface/v2/test_protocol.rb b/test/interface/v2/test_protocol.rb index 15116376b..2b92d5a17 100644 --- a/test/interface/v2/test_protocol.rb +++ b/test/interface/v2/test_protocol.rb @@ -29,15 +29,14 @@ module Protocol it "adds deployment-specific info" do flexmock(@deployment, pid: 200) - handles = { - "test" => Syskit::Deployment::RemoteTaskHandles.new( - flexmock( - ior: "some_ior", - model: flexmock(name: "orogen::Name") - ) + handle = Syskit::Deployment::RemoteTaskHandles.new( + flexmock( + ior: "some_ior", + model: flexmock(name: "orogen::Name") ) - } - flexmock(@deployment, remote_task_handles: handles) + ) + handle.configured_since = t0 = Time.now + flexmock(@deployment, remote_task_handles: { "test" => handle }) marshalled = @channel.marshal_filter_object(@deployment) assert_equal 200, marshalled.pid @@ -45,7 +44,8 @@ module Protocol expected_task = { name: "test", ior: "some_ior", - orogen_model_name: "orogen::Name" + orogen_model_name: "orogen::Name", + configured_since: t0 } assert_equal [expected_task], marshalled.deployed_tasks.map(&:to_h) diff --git a/test/telemetry/async/test_name_service.rb b/test/telemetry/async/test_name_service.rb index 605038c66..1843bb9a1 100644 --- a/test/telemetry/async/test_name_service.rb +++ b/test/telemetry/async/test_name_service.rb @@ -61,6 +61,42 @@ module Async assert_equal task2.ior, @ns.get("test").identity end + it "does not update an existing async task " \ + "if it has not been reconfigured" do + deployed_task, = make_deployed_task("test", "something") + + @ns.async_update_tasks([deployed_task]) + @ns.wait_for_task_discovery + @ns.resolve_discovered_tasks + + flexmock(@ns.get("test")) + .should_receive(:update_from_discovered_interface) + .never + @ns.async_update_tasks([deployed_task]) + @ns.wait_for_task_discovery + @ns.resolve_discovered_tasks + end + + it "updates an existing async task if it has been reconfigured" do + deployed_task, ruby_task = make_deployed_task("test", "something") + + @ns.async_update_tasks([deployed_task]) + @ns.wait_for_task_discovery + @ns.resolve_discovered_tasks + + flexmock(@ns.get("test")) + .should_receive(:update_from_discovered_interface) + .once.pass_thru + ruby_task.create_output_port("new_port", "/int32_t") + deployed_task.configured_since = Time.now + @ns.async_update_tasks([deployed_task]) + @ns.wait_for_task_discovery + @ns.resolve_discovered_tasks + + assert_equal %w[state new_port], + @ns.get("test").each_port.map(&:name) + end + it "does not register a task if it has been removed while it was " \ "being discovered" do deployed_task, = make_deployed_task("test", "something") @@ -265,8 +301,7 @@ module Async end def deployed_task_s - @deployed_task_s ||= - Struct.new(:name, :ior, :orogen_model_name, keyword_init: true) + @deployed_task_s ||= Syskit::Interface::V2::Protocol::DeployedTask end def make_deployed_task(name, orogen_model_name) diff --git a/test/telemetry/async/test_task_context.rb b/test/telemetry/async/test_task_context.rb index ed31c6b1a..ec494792d 100644 --- a/test/telemetry/async/test_task_context.rb +++ b/test/telemetry/async/test_task_context.rb @@ -311,6 +311,164 @@ module Async end end + describe "update_from_discovered_interface" do + it "adds new attributes" do + task, async = make_async_task "test" + Orocos.allow_blocking_calls do + task.create_attribute("new_a", "/int32_t") + end + + objects = [] + async.on_attribute_reachable { objects << _1 } + objects.clear + + async.update_from_discovered_interface(discover_interface(task)) + + assert_equal Set["new_a"], objects.to_set + assert_kind_of Attribute, async.attribute("new_a") + end + + it "removes attributes that disappeared" do + task, async = make_async_task "test" + + # NOTE: ruby task contexts do not have APIs to remove + # attributes, so modify the discovered data + discovery = discover_interface(task) + discovery.attributes.delete_if { |a| a.name == "attr" } + + removed = [] + async.on_attribute_unreachable { removed << _1 } + async.update_from_discovered_interface(discovery) + + assert_equal Set["attr"], removed.to_set + end + + it "does nothing for unchanged attributes" do + task, async = make_async_task "test" + names = [] + async.on_attribute_reachable { names << _1 } + async.on_attribute_unreachable { names << _1 } + names.clear + + flexmock(Attribute) + .new_instances.should_receive(:reachable!).never + async.update_from_discovered_interface(discover_interface(task)) + + assert_empty names + end + it "adds new properties" do + task, async = make_async_task "test" + Orocos.allow_blocking_calls do + task.create_property("new_p", "/int32_t") + end + + objects = [] + async.on_property_reachable { objects << _1 } + objects.clear + + async.update_from_discovered_interface(discover_interface(task)) + + assert_equal Set["new_p"], objects.to_set + assert_equal Set["prop", "new_p"], + async.each_property.to_set(&:name) + assert_kind_of Property, async.property("new_p") + end + + it "removes properties that disappeared" do + task, async = make_async_task "test" + + # NOTE: ruby task contexts do not have APIs to remove + # properties, so modify the discovered data + discovery = discover_interface(task) + discovery.properties.clear + + removed = [] + async.on_property_unreachable { removed << _1 } + async.update_from_discovered_interface(discovery) + + assert_equal Set["prop"], removed.to_set + end + + it "does nothing for unchanged properties" do + task, async = make_async_task "test" + names = [] + async.on_property_reachable { names << _1 } + async.on_property_unreachable { names << _1 } + names.clear + + flexmock(Property).new_instances.should_receive(:reachable!).never + async.update_from_discovered_interface(discover_interface(task)) + + assert_empty names + end + + it "adds new output ports" do + task, async = make_async_task "test" + Orocos.allow_blocking_calls do + task.create_output_port("new_out", "/int32_t") + end + + ports = [] + async.on_port_reachable { ports << _1 } + ports.clear + + async.update_from_discovered_interface(discover_interface(task)) + + assert_equal Set["new_out"], ports.to_set + assert_equal Set["state", "in", "out", "new_out"], + async.each_port.to_set(&:name) + assert_kind_of OutputPort, async.port("new_out") + end + + it "adds new input ports" do + task, async = make_async_task "test" + Orocos.allow_blocking_calls do + task.create_input_port("new_in", "/int32_t") + end + + ports = [] + async.on_port_reachable { ports << _1 } + ports.clear + + async.update_from_discovered_interface(discover_interface(task)) + + assert_equal Set["new_in"], ports.to_set + assert_equal Set["state", "in", "out", "new_in"], + async.each_port.to_set(&:name) + assert_kind_of InputPort, async.port("new_in") + end + + it "removes ports that disappeared" do + task, async = make_async_task "test" + Orocos.allow_blocking_calls do + task.remove_port(task.port("out")) + end + + ports = [] + async.on_port_unreachable { ports << _1 } + async.update_from_discovered_interface(discover_interface(task)) + + assert_equal Set["out"], ports.to_set + assert_equal Set["state", "in"], async.each_port.to_set(&:name) + end + + it "does nothing for unchanged ports" do + task, async = make_async_task "test" + ports = [] + async.on_port_reachable { ports << _1 } + async.on_port_unreachable { ports << _1 } + ports.clear + + flexmock(InputPort) + .new_instances.should_receive(:reachable!).never + flexmock(OutputPort) + .new_instances.should_receive(:reachable!).never + async.update_from_discovered_interface(discover_interface(task)) + + assert_empty ports + end + end + def make_ruby_task(name) ruby_task = Orocos.allow_blocking_calls do t = Orocos::RubyTasks::TaskContext.new(name) @@ -336,6 +494,12 @@ def discover_task(task) ) end + def discover_interface(task) + Orocos.allow_blocking_calls do + TaskContext.discover_interface(task.name, task.ior, task.model) + end + end + def assert_polling_eventually(period: 0.01, timeout: 2, &block) deadline = Time.now + timeout while Time.now < deadline