Class: OmfCommon::Comm::AMQP::FileBroadcaster
- Inherits:
-
Object
- Object
- OmfCommon::Comm::AMQP::FileBroadcaster
- Includes:
- MonitorMixin
- Defined in:
- lib/omf_common/comm/amqp/amqp_file_transfer.rb
Overview
Distributes a local file to a set of receivers subscribed to the same topic but may join a various stages.
Constant Summary collapse
- DEF_CHUNK_SIZE =
2**16
- DEF_IDLE_TIME =
60
Instance Method Summary collapse
- #_send(f, chunk_size, chunk_count, exchange, idle_time) ⇒ Object
- #_wait_for_closedown(idle_time) ⇒ Object
-
#initialize(file_path, channel, topic, opts = {}, &block) ⇒ FileBroadcaster
constructor
A new instance of FileBroadcaster.
Constructor Details
#initialize(file_path, channel, topic, opts = {}, &block) ⇒ FileBroadcaster
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 |
# File 'lib/omf_common/comm/amqp/amqp_file_transfer.rb', line 27 def initialize(file_path, channel, topic, opts = {}, &block) super() # init monitor mixin @block = block unless File.readable?(file_path) raise "Can't read file '#{file_path}'" end @mime_type = `file -b --mime-type #{file_path}`.strip unless $?.success? raise "Can't determine file's mime-type (#{$?})" end @file_path = file_path f = File.open(file_path, 'rb') chunk_size = opts[:chunk_size] || DEF_CHUNK_SIZE chunk_count = (f.size / chunk_size) + 1 @outstanding_chunks = Set.new @running = true @semaphore = new_cond() idle_time = opts[:idle_time] || DEF_IDLE_TIME #chunk_count.times.each {|i| @outstanding_chunks << i} exchange = channel.topic(topic, :auto_delete => true) OmfCommon.eventloop.defer do _send(f, chunk_size, chunk_count, exchange, idle_time) end control_topic = "#{topic}_control" control_exchange = channel.topic(control_topic, :auto_delete => true) channel.queue("", :exclusive => false) do |queue| queue.bind(control_exchange) debug "Subscribing to control channel '#{control_topic}'" queue.subscribe do |headers, payload| hdrs = headers.headers debug "Incoming control message '#{hdrs}'" from = hdrs['request_from'] from = 0 if from < 0 to = hdrs['request_to'] to = chunk_count - 1 if !to || to >= chunk_count synchronize do (from .. to).each { |i| @outstanding_chunks << i} @semaphore.signal end end @control_queue = queue end end |
Instance Method Details
#_send(f, chunk_size, chunk_count, exchange, idle_time) ⇒ Object
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 |
# File 'lib/omf_common/comm/amqp/amqp_file_transfer.rb', line 75 def _send(f, chunk_size, chunk_count, exchange, idle_time) chunks_to_send = nil @sent_chunk = false _wait_for_closedown(idle_time) loop do synchronize do @semaphore.wait_while { @outstanding_chunks.empty? && @running } return unless @running # done! chunks_to_send = @outstanding_chunks.to_a end chunks_to_send.each do |chunk_id| #sleep 3 synchronize do @outstanding_chunks.delete(chunk_id) @sent_chunk = true end offset = chunk_id * chunk_size f.seek(offset, IO::SEEK_SET) chunk = f.read(chunk_size) payload = Base64.encode64(chunk) headers = {chunk_id: chunk_id, chunk_count: chunk_count, chunk_offset: offset, chunk_size: payload.size, path: f.path, file_size: f.size, mime_type: @mime_type} debug "Sending chunk #{chunk_id}" exchange.publish(payload, {headers: headers}) end end end |
#_wait_for_closedown(idle_time) ⇒ Object
105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 |
# File 'lib/omf_common/comm/amqp/amqp_file_transfer.rb', line 105 def _wait_for_closedown(idle_time) OmfCommon.eventloop.after(idle_time) do done = false synchronize do done = !@sent_chunk && @outstanding_chunks.empty? @sent_chunk = false end if done @control_queue.unsubscribe if @control_queue @block.call({action: :done}) if @block else # there was activity in last interval, wait a bit longer _wait_for_closedown(idle_time) end end end |