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
Returns a new instance of 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 |