diff --git a/docker-compose.yml b/docker-compose.yml index 2e225817..eb029814 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -98,7 +98,7 @@ services: <<: *s3-client-env minio: - image: minio/minio:latest + image: pgsty/minio:latest volumes: - s3:/data ports: @@ -111,7 +111,7 @@ services: command: server /data --console-address ":9090" testbucket: - image: minio/mc:latest + image: pgsty/mc:latest depends_on: - minio environment: diff --git a/spec/mappings/control_system_modules_spec.cr b/spec/mappings/control_system_modules_spec.cr index 15845d44..f33fba6d 100644 --- a/spec/mappings/control_system_modules_spec.cr +++ b/spec/mappings/control_system_modules_spec.cr @@ -81,6 +81,54 @@ module PlaceOS::Core::Mappings end end + describe ".set_mappings" do + it "waits for an update of the same system that is still writing" do + driver = Model::Generator.driver(:device).save! + + modules = %w(alpha beta gamma delta epsilon).map do |name| + m = Model::Generator.module(driver) + m.custom_name = name + m.save! + end + + cs = Model::Generator.control_system + cs.modules = modules.compact_map &.id + cs.save! + + expected = modules.map { |m| "#{m.custom_name}/1" } + storage = Driver::RedisStorage.new(cs.id.as(String), "system") + + # an earlier update of this system is part way through writing + lock = ControlSystemModules.mapping_lock + lock.lock + + finished = false + done = Channel(Nil).new + spawn do + ControlSystemModules.set_mappings(cs, nil) + finished = true + ensure + done.send(nil) + end + + # plenty of time for the second update to run if nothing held it back + sleep 0.3.seconds + finished.should be_false + + # the earlier update's writes land while the second one waits + storage.clear + modules.reverse_each { |m| storage["#{m.custom_name}/1"] = m.id.as(String) } + storage.keys.should eq expected.reverse + + lock.unlock + done.receive + finished.should be_true + + # the second update rewrote the mapping in the system's order + storage.keys.should eq expected + end + end + describe ".update_logic_modules" do it "does not update if system is destroyed" do cs = Model::ControlSystem.new diff --git a/src/placeos-core/mappings/control_system_modules.cr b/src/placeos-core/mappings/control_system_modules.cr index bd0b8e01..23e5e345 100644 --- a/src/placeos-core/mappings/control_system_modules.cr +++ b/src/placeos-core/mappings/control_system_modules.cr @@ -85,6 +85,10 @@ module PlaceOS::Core updated_modules end + # Held while a system's mappings are written, so two updates write one after + # the other rather than interleaved + class_getter mapping_lock : Mutex = Mutex.new + # Set the module mappings for a ControlSystem # # Pass module_id and updated_name to overrride a lookup @@ -93,13 +97,19 @@ module PlaceOS::Core mod : Model::Module?, ) : Hash(String, String) system_id = control_system.id.as(String) - storage = Driver::RedisStorage.new(system_id, "system") + mapping_lock.synchronize { write_mappings(control_system, mod, system_id) } + end - # Clear out the ControlSystem's mapping - storage.clear + protected def self.write_mappings( + control_system : Model::ControlSystem, + mod : Model::Module?, + system_id : String, + ) : Hash(String, String) + storage = Driver::RedisStorage.new(system_id, "system") # No mappings to set if ControlSystem has been destroyed if control_system.destroyed? + storage.clear Log.info { { message: "module mappings deleted", system_id: control_system.id, @@ -122,7 +132,9 @@ module PlaceOS::Core end end - # Set the mappings in redis + # Replace the ControlSystem's mapping. The lookups above can take a + # while, so the hash is only empty for the time it takes to write it. + storage.clear mappings.each do |mapping, module_id| storage[mapping] = module_id end