Class: Aspera::Agent::Transferd

Inherits:
Base
  • Object
show all
Defined in:
lib/aspera/agent/transferd.rb

Instance Attribute Summary collapse

Class Method Summary collapse

Instance Method Summary collapse

Methods inherited from Base

#last_job_id, #wait_for_completion

Constructor Details

#initialize(url: AUTO_LOCAL_TCP_PORT, start: true, stop: true, **base) ⇒ Transferd

Returns a new instance of Transferd.

Parameters:

  • url (String) (defaults to: AUTO_LOCAL_TCP_PORT) —

    URL of the transfer manager daemon

  • start (Boolean) (defaults to: true) —

    If false, expect that an external daemon is already running

  • stop (Boolean) (defaults to: true) —

    If false, do not shutdown daemon on exit

  • base (Hash) —

    Base class options



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
134
135
136
137
138
139
# File 'lib/aspera/agent/transferd.rb', line 59

def initialize(
  url:   AUTO_LOCAL_TCP_PORT,
  start: true,
  stop:  true,
  **base
)
  super(**base)
  @transfer_id = nil
  @stop = stop
  is_local_auto_port = url.eql?(AUTO_LOCAL_TCP_PORT)
  Aspera.assert(!(is_local_auto_port && (!@stop || !start)), type: Error) { 'Cannot set options `stop` or `start` to false with port zero' }
  # keep PID for optional shutdown
  @daemon_pid = nil
  @daemon_endpoint = nil
  daemon_endpoint = url
  Log.dump(:daemon_endpoint, daemon_endpoint)
  # retry loop
  begin
    # no address: local bind
    daemon_endpoint = "#{LOCAL_SOCKET_ADDR}#{daemon_endpoint}" if daemon_endpoint.match?(/^#{PORT_SEP}[0-9]+$/o)
    # Create stub (without credentials)
    @transfer_client = ::Transferd::Api::TransferService::Stub.new(daemon_endpoint, :this_channel_is_insecure)
    # Initiate actual connection
    get_info_response = @transfer_client.get_info(::Transferd::Api::InstanceInfoRequest.new)
    @daemon_endpoint = daemon_endpoint
    Log.log.debug { "Daemon info: #{get_info_response}" }
    Log.log.warn('Attached to existing daemon') unless @daemon_pid || !start || !@stop
    at_exit { shutdown }
  rescue GRPC::Unavailable => e
    # if transferd is external: do not start it, or other error
    raise if !start || !e.message.include?('failed to connect')
    # we already tried to start a daemon, but it failed
    Aspera.assert(@daemon_pid.nil?) { "Daemon started with PID #{@daemon_pid}, but connection failed to #{daemon_endpoint}" }
    Log.log.warn('no daemon present, starting daemon...') if !start
    # transferd only supports local ip and port
    daemon_uri = URI.parse("ipv4://#{daemon_endpoint}")
    Aspera.assert(daemon_uri.scheme.eql?('ipv4')) { "Invalid scheme daemon URI #{daemon_endpoint}" }
    # create a config file for daemon
    config = {
      address:      daemon_uri.host,
      port:         daemon_uri.port,
      fasp_runtime: {
        use_embedded: false,
        user_defined: {
          bin: Products::Transferd.sdk_directory,
          etc: Products::Transferd.sdk_directory
        }
      }
    }
    # config file and logs are created in same folder
    transferd_base_tmp = TempFileManager.instance.new_file_path_global('transferd')
    Log.log.debug { "transferd base tmp #{transferd_base_tmp}" }
    conf_file = "#{transferd_base_tmp}.conf"
    log_stdout = "#{transferd_base_tmp}.out"
    log_stderr = "#{transferd_base_tmp}.err"
    File.write(conf_file, config.to_json)
    @daemon_pid = Environment.secure_execute(
      Ascp::Installation.instance.path(:transferd), '--config', conf_file,
      mode: :background,
      out: log_stdout,
      err: log_stderr
    )
    begin
      # wait for process to initialize, max 2 seconds
      Timeout.timeout(2.0) do
        # this returns if process dies (within 2 seconds)
        _, status = Process.wait2(@daemon_pid)
        raise "Transfer daemon exited with status #{status.exitstatus}. Check files: #{log_stdout} and #{log_stderr}"
      end
    rescue Timeout::Error
      nil
    end
    Log.log.debug { "Daemon started with pid #{@daemon_pid}" }
    Process.detach(@daemon_pid) unless @stop
    at_exit { shutdown }
    # update port for next connection attempt (if auto high port was requested)
    daemon_endpoint = "#{LOCAL_SOCKET_ADDR}#{PORT_SEP}#{Products::Transferd.daemon_port_from_log(log_stdout)}" if is_local_auto_port
    # local daemon started, try again
    retry
  end
end

Instance Attribute Details

#daemon_endpoint ⇒ Object (readonly)

Actual endpoint the daemon is listening on (resolved after connect, e.g. when port 0 was used)



142
143
144
# File 'lib/aspera/agent/transferd.rb', line 142

def daemon_endpoint
  @daemon_endpoint
end

Class Method Details

.transfer_status(transfer_id, agent_params) ⇒ Hash

Re-query a previously started transfer from a running transferd daemon. agent_params must contain 'url' (the daemon endpoint, e.g. "127.0.0.1:12345").

Returns:

  • (Hash) —

    normalized status hash



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
# File 'lib/aspera/agent/transferd.rb', line 28

def transfer_status(transfer_id, agent_params)
  client = ::Transferd::Api::TransferService::Stub.new(agent_params['url'], :this_channel_is_insecure)
  result = {'bytes_transferred' => 0}
  client.monitor_transfers(::Transferd::Api::RegistrationRequest.new(transferId: [transfer_id])) do |response|
    result = case response.status
    when :COMPLETED
      {
        'status' => 'completed', 'ended_at' => Time.now.utc.iso8601,
                'bytes_transferred' => response.transferInfo.bytesTransferred
      }
    when :FAILED, :CANCELED
      {
        'status' => response.status == :FAILED ? 'failed' : 'cancelled',
                'ended_at' => Time.now.utc.iso8601, 'bytes_transferred' => 0,
                'error' => JSON.parse(response.message.to_s)['Description']
      }
    else
      {'status' => 'running', 'bytes_transferred' => response.transferInfo.bytesTransferred}
    end
    break
  end
  result
rescue StandardError => e
  {'status' => 'running', 'bytes_transferred' => 0, 'error' => e.message}
end

Instance Method Details

#shutdown ⇒ Object



200
201
202
# File 'lib/aspera/agent/transferd.rb', line 200

def shutdown
  stop_daemon if @stop
end

#start_transfer(transfer_spec, token_regenerator: nil) ⇒ Object

:reek:UnusedParameters token_regenerator



145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
# File 'lib/aspera/agent/transferd.rb', line 145

def start_transfer(transfer_spec, token_regenerator: nil)
  Transfer::Spec.fix_transferd_resume_policy(transfer_spec)
  # create a transfer request
  transfer_request = ::Transferd::Api::TransferRequest.new(
    transferType: ::Transferd::Api::TransferType::FILE_REGULAR, # transfer type (file/stream)
    config: ::Transferd::Api::TransferConfig.new, # transfer configuration
    transferSpec: transfer_spec.to_json
  ) # transfer definition
  # send start transfer request to the transfer manager daemon
  start_response = @transfer_client.start_transfer(transfer_request)
  Aspera.assert(!start_response.status.eql?(:FAILED), type: Transfer::Error) { start_response.error.description }
  Log.log.debug { "start transfer response #{start_response}" }
  @transfer_id = start_response.transferId
  Log.log.debug { "transfer started with id #{@transfer_id}" }
end

#wait_for_transfers_completion ⇒ Object



161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
# File 'lib/aspera/agent/transferd.rb', line 161

def wait_for_transfers_completion
  # set to true when we know the total size of the transfer
  session_started = false
  bytes_expected = nil
  # monitor transfer status
  @transfer_client.monitor_transfers(::Transferd::Api::RegistrationRequest.new(transferId: [@transfer_id])) do |response|
    Log.dump(:response, response.to_h)
    # Log.log.debug{"#{response.sessionInfo.preTransferBytes} #{response.transferInfo.bytesTransferred}"}
    case response.status
    when :RUNNING
      if !session_started
        notify_progress(:session_start, session_id: @transfer_id)
        session_started = true
      end
      if bytes_expected.nil? &&
          !response.sessionInfo.preTransferBytes.eql?(0)
        bytes_expected = response.sessionInfo.preTransferBytes
        notify_progress(:session_size, session_id: @transfer_id, info: bytes_expected)
      end
      notify_progress(:transfer, session_id: @transfer_id, info: response.transferInfo.bytesTransferred)
    when :COMPLETED
      notify_progress(:transfer, session_id: @transfer_id, info: bytes_expected) if bytes_expected
      notify_progress(:session_end, session_id: @transfer_id)
      notify_progress(:end)
      break
    when :FAILED, :CANCELED
      notify_progress(:session_end, session_id: @transfer_id)
      notify_progress(:end)
      raise Transfer::Error, JSON.parse(response.message)['Description']
    when :QUEUED, :UNKNOWN_STATUS, :PAUSED, :ORPHANED
      notify_progress(:sessions_init, info: response.status.to_s.downcase)
    else
      Log.log.error { "unknown status#{response.status}" }
    end
  end
  # TODO: return status
  return []
end