Module: SDN::CLI::MQTT::Read
- Defined in:
- lib/sdn/cli/mqtt/read.rb
Overview
Reader loop that translates SDN responses into MQTT state updates.
Instance Method Summary collapse
-
#read ⇒ void
Continuously consumes SDN messages and republishes state to MQTT.
Instance Method Details
#read ⇒ void
This method returns an undefined value.
Continuously consumes SDN messages and republishes state to MQTT.
11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 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 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 |
# File 'lib/sdn/cli/mqtt/read.rb', line 11 def read loop do @sdn.receive do || @mqtt.batch_publish do src = Message.print_address(.src) # ignore the UAI Plus and ourselves if src != "7F.7F.7F" && !Message.group_address?(.src) && !(motor = @motors[src.delete(".")]) SDN.logger.info "Found new motor #{src}" @motors_found = true motor = publish_motor(src.delete("."), .node_type) end follow_ups = [] case when Message::PostNodeLabel publish("#{motor.addr}/$name", .label) if motor.publish(:label, .label) && @homie when Message::PostMotorPosition, Message::ILT2::PostMotorPosition if .is_a?(Message::ILT2::PostMotorPosition) # keep polling while it's still moving; check prior two positions if motor.position_pulses == .position_pulses && motor.last_position_pulses == .position_pulses motor.publish(:state, :stopped) else motor.publish(:state, :running) if motor.position_pulses && motor.position_pulses != .position_pulses motor.publish(:last_direction, (motor.position_pulses < .position_pulses) ? :down : :up) end follow_ups << Message::ILT2::GetMotorPosition.new(.src) end motor.last_position_pulses = motor.position_pulses ip = (1..16).find do |i| # divide by 5 for some leniency motor["ip#{i}_pulses"].to_i / 5 == .position_pulses / 5 end motor.publish(:ip, ip) end motor.publish(:position_percent, .position_percent) motor.publish(:position_pulses, .position_pulses) motor.publish(:ip, .ip) if .respond_to?(:ip) motor.group_objects.each do |group| positions_percent = group.motor_objects.map(&:position_percent) positions_pulses = group.motor_objects.map(&:position_pulses) ips = group.motor_objects.map(&:ip) position_percent = nil # calculate an average, but only if we know a position for # every shade if !positions_percent.include?(:nil) && !positions_percent.include?(nil) position_percent = positions_percent.sum / positions_percent.length end position_pulses = nil if !positions_pulses.include?(:nil) && !positions_pulses.include?(nil) position_pulses = positions_pulses.sum / positions_pulses.length end ip = nil ip = ips.first if ips.uniq.length == 1 ip = nil if ip == :nil group.publish(:position_percent, position_percent) group.publish(:position_pulses, position_pulses) group.publish(:ip, ip) end when Message::PostMotorStatus handle_post_motor_status(, motor, follow_ups) when Message::PostMotorLimits motor.publish(:up_limit, .up_limit) motor.publish(:down_limit, .down_limit) when Message::ILT2::PostMotorSettings motor.publish(:down_limit, .limit) when Message::PostMotorDirection motor.publish(:direction, .direction) when Message::PostMotorRollingSpeed motor.publish(:up_speed, .up_speed) motor.publish(:down_speed, .down_speed) motor.publish(:slow_speed, .slow_speed) when Message::PostMotorIP, Message::ILT2::PostMotorIP motor.publish(:"ip#{.ip}_pulses", .position_pulses) if .respond_to?(:position_percent) motor.publish(:"ip#{.ip}_percent", .position_percent) elsif motor.down_limit motor.publish(:"ip#{.ip}_percent", .position_pulses.to_f / motor.down_limit * 100) end when Message::PostGroupAddr motor.add_group(.group_index, .group_address) end @mutex.synchronize do = Message.group_address?(@prior_message&.&.src) if @prior_message correct_response = @response_pending && @prior_message&.&.class&.expected_response?() correct_response = false if ! && .src != @prior_message&.&.dest correct_response = false if && .dest != @prior_message&.&.src if && correct_response @pending_group_motors.delete(Message.print_address(.src).delete(".")) correct_response = false unless @pending_group_motors.empty? end signal = correct_response || !follow_ups.empty? @response_pending = @broadcast_pending if correct_response follow_ups.each do |follow_up| unless @queue.any? { |mr| mr. == follow_up } @queue.push(MessageAndRetries.new(follow_up, 5, 1)) end end @cond.signal if signal end rescue EOFError SDN.logger.fatal "EOF reading" exit 2 rescue MalformedMessage => e SDN.logger.warn "Ignoring malformed message: #{e}" unless e.to_s.include?("issing data") rescue => e SDN.logger.error "Got garbage: #{e}; #{e.backtrace}" end end end end |