Class: Aspera::Cli::TransferAgent
- Inherits:
-
Object
- Object
- Aspera::Cli::TransferAgent
- Defined in:
- lib/aspera/cli/transfer_agent.rb
Overview
The Transfer agent is a common interface to start a transfer using one of the supported transfer agents. Provide CLI options to select one of the transfer agents (FASP/ascp client)
Constant Summary collapse
- FILE_LIST_FROM_ARGS =
@argsspecial value for --sources : read file list from arguments '@args'- CP4I_REMOTE_HOST_LB =
'N/A'
Instance Attribute Summary collapse
-
#transfer_options ⇒ Object
Returns the value of attribute transfer_options.
-
#user_transfer_spec ⇒ Object
Returns the value of attribute user_transfer_spec.
Class Method Summary collapse
-
.declare_options(options) ⇒ Object
Declare all transfer CLI options (metadata only - no handler binding yet).
Instance Method Summary collapse
-
#agent_instance ⇒ Object
analyze options and create new agent if not already created or set.
- #agent_instance=(instance) ⇒ Object
-
#destination_folder(direction) ⇒ String
Get destination folder.
- #httpgw_url_cb=(httpgw_url_proc) ⇒ Object
-
#initialize(context) ⇒ TransferAgent
constructor
A new instance of TransferAgent.
-
#list_to_paths(file_list) ⇒ Object
Transform the list of paths to a list of hash with source/dest.
-
#option_transfer(_option_sym, operation, value = nil) ⇒ Object
Composite option handler for :transfer String value: shorthand for agent type, stored as => value Hash value: merged into @transfer_options (may include 'agent' key).
-
#shutdown ⇒ Object
shut down if agent requires it.
-
#source_list ⇒ Array
List of source files.
-
#start(transfer_spec, rest_token: nil) ⇒ Object
Start a transfer and wait for completion, plugins shall use this method.
-
#ts_source_paths(default: nil) ⇒ Array?
This is how the list of files to be transferred is specified get paths suitable for transfer spec from command line computation is done only once, cache is kept in @transfer_paths.
Constructor Details
#initialize(context) ⇒ TransferAgent
Returns a new instance of TransferAgent.
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 |
# File 'lib/aspera/cli/transfer_agent.rb', line 52 def initialize(context) Aspera.assert_type(context, Context){'context'} Aspera.assert_type(context., Parser){'context.options'} @context = context # Command line can override transfer spec @user_transfer_spec = { 'create_dir' => true, 'resume_policy' => 'sparse_csum' } # options for transfer agent (agent type + agent-specific parameters) @transfer_options = {} # the currently selected transfer agent @agent = nil # source/destination pair, like "paths" of transfer spec @transfer_paths = nil # HTTPGW URL provided by webapp @httpgw_url_lambda = nil self.class.(@context.) @context..set_handler(:ts, object: self, method: :user_transfer_spec) @context..set_handler(:transfer, object: self, method: :option_transfer) @context..set_handler(:transfer_info, object: self, method: :transfer_options) @context.. @notification_cb = nil if !@context..get_option(:notify_to).nil? @notification_cb = ->(transfer_spec, global_status) do @context.mailer.send_email_template(email_template_default: DEFAULT_TRANSFER_NOTIFY_TEMPLATE, values: { subject: "#{Info::CMD_NAME} transfer: #{global_status}", status: global_status, ts: transfer_spec }) end end end |
Instance Attribute Details
#transfer_options ⇒ Object
Returns the value of attribute transfer_options.
86 87 88 |
# File 'lib/aspera/cli/transfer_agent.rb', line 86 def @transfer_options end |
#user_transfer_spec ⇒ Object
Returns the value of attribute user_transfer_spec.
86 87 88 |
# File 'lib/aspera/cli/transfer_agent.rb', line 86 def user_transfer_spec @user_transfer_spec end |
Class Method Details
.declare_options(options) ⇒ Object
Declare all transfer CLI options (metadata only - no handler binding yet).
41 42 43 44 45 46 47 48 |
# File 'lib/aspera/cli/transfer_agent.rb', line 41 def () .declare(:ts, description: 'Override transfer spec values', schema: Schema::Registry::TRANSFER_SPEC) .declare(:to_folder, description: 'Destination folder for transferred files') .declare(:sources, description: "How list of transferred files is provided (#{FILE_LIST_OPTIONS.join(',')})", default: FILE_LIST_FROM_ARGS) .declare(:src_type, description: 'Type of file list', allowed: %i[list pair], default: :list) .declare(:transfer, description: 'Transfer agent type, or agent parameters with optional agent key', allowed: [Hash, String], schema: Schema::Registry::TRANSFER_AGENT_OPTIONS) .declare(:transfer_info, description: 'Parameters for transfer agent', allowed: Hash, deprecation: 'use --transfer instead', schema: Schema::Registry::TRANSFER_AGENT_OPTIONS) end |
Instance Method Details
#agent_instance ⇒ Object
analyze options and create new agent if not already created or set
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 140 141 142 |
# File 'lib/aspera/cli/transfer_agent.rb', line 109 def agent_instance return @agent unless @agent.nil? # agent type: from composite option 'agent' key, default :direct raw_type = @transfer_options['agent'] || :direct agent_type = Parser.get_from_list(raw_type.to_s, 'transfer agent', Agent::Factory::ALL.keys) # set keys as symbols, strip internal keys not forwarded to the agent constructor = @transfer_options.except('agent', 'asynchronous').symbolize_keys [:progress] = @context. [:config_dir] = @context.main_folder # special cases case agent_type when :node if !.key?(:url) param_set_name = @context.presets.plugin_default_name(:node) Aspera.assert(!param_set_name.nil?, type: Cli::BadArgument){"No default node configured. Please specify #{Options.option_name_to_line(:transfer)}.url or #{Options.option_name_to_line(:transfer)}"} .merge!(@context.presets.by_name(param_set_name).symbolize_keys) end when :direct # by default do not display ascp native progress bar [:quiet] = true unless .key?(:quiet) [:check_ignore_cb] = ->(host, port){@context.http_config.ignore_cert?(host, port)} # JRuby [:trusted_certs] = @context.http_config.trusted_cert_locations unless .key?(:trusted_certs) when :httpgw unless .key?(:url) || @httpgw_url_lambda.nil? Log.log.debug('retrieving HTTPGW URL from webapp') [:url] = @httpgw_url_lambda.call end end # get agent instance self.agent_instance = Agent::Factory.instance.create(agent_type, ) Log.log.debug{"transfer agent is a #{@agent.class}"} return @agent end |
#agent_instance=(instance) ⇒ Object
104 105 106 |
# File 'lib/aspera/cli/transfer_agent.rb', line 104 def agent_instance=(instance) @agent = instance end |
#destination_folder(direction) ⇒ String
Get destination folder
147 148 149 150 151 152 153 154 155 156 157 158 159 160 |
# File 'lib/aspera/cli/transfer_agent.rb', line 147 def destination_folder(direction) dest_folder = @context..get_option(:to_folder) # do not expand path, if user wants to expand path: user @path: return dest_folder unless dest_folder.nil? dest_folder = @user_transfer_spec['destination_root'] return dest_folder unless dest_folder.nil? # default: / on remote, . on local case direction.to_s when Transfer::Spec::DIRECTION_SEND then dest_folder = '/' when Transfer::Spec::DIRECTION_RECEIVE then dest_folder = '.' else Aspera.error_unexpected_value(direction) end return dest_folder end |
#httpgw_url_cb=(httpgw_url_proc) ⇒ Object
170 171 172 173 |
# File 'lib/aspera/cli/transfer_agent.rb', line 170 def httpgw_url_cb=(httpgw_url_proc) Aspera.assert_type(httpgw_url_proc, Proc){'httpgw_url_cb'} @httpgw_url_lambda = httpgw_url_proc end |
#list_to_paths(file_list) ⇒ Object
Transform the list of paths to a list of hash with source/dest
177 178 179 180 181 182 183 184 185 186 187 188 189 |
# File 'lib/aspera/cli/transfer_agent.rb', line 177 def list_to_paths(file_list) source_type = @context..get_option(:src_type, mandatory: true) @transfer_paths = case source_type when :list # when providing a list, just specify source file_list.map{ |i| {'source' => i}} when :pair Aspera.assert(file_list.length.even?, type: Cli::BadArgument){"When using pair, provide an even number of paths: #{file_list.length}"} file_list.each_slice(2).map{ |s, d| {'source' => s, 'destination' => d}} else Aspera.error_unexpected_value(source_type) end end |
#option_transfer(_option_sym, operation, value = nil) ⇒ Object
Composite option handler for :transfer String value: shorthand for agent type, stored as => value Hash value: merged into @transfer_options (may include 'agent' key)
91 92 93 94 95 96 97 98 99 100 101 102 |
# File 'lib/aspera/cli/transfer_agent.rb', line 91 def option_transfer(_option_sym, operation, value = nil) Aspera.assert_values(operation, %i[set get]) case operation when :set value = {'agent' => value} if value.is_a?(String) Aspera.assert_type(value, Hash) @transfer_options = @transfer_options.deep_merge(value) when :get return @transfer_options end nil end |
#shutdown ⇒ Object
shut down if agent requires it
314 315 316 |
# File 'lib/aspera/cli/transfer_agent.rb', line 314 def shutdown @agent.shutdown if @agent.respond_to?(:shutdown) end |
#source_list ⇒ Array
Returns list of source files.
163 164 165 166 167 |
# File 'lib/aspera/cli/transfer_agent.rb', line 163 def source_list return ts_source_paths.map do |i| i['source'] end end |
#start(transfer_spec, rest_token: nil) ⇒ Object
Start a transfer and wait for completion, plugins shall use this method
234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 |
# File 'lib/aspera/cli/transfer_agent.rb', line 234 def start(transfer_spec, rest_token: nil) # check parameters Aspera.assert_type(transfer_spec, Hash){'transfer_spec'} raise "Wrong remote host: #{CP4I_REMOTE_HOST_LB}" if transfer_spec['remote_host'].eql?(CP4I_REMOTE_HOST_LB) # process :src option case transfer_spec['direction'] when Transfer::Spec::DIRECTION_RECEIVE # init default if required in any case @user_transfer_spec['destination_root'] ||= destination_folder(transfer_spec['direction']) when Transfer::Spec::DIRECTION_SEND if transfer_spec.dig('tags', Transfer::Spec::TAG_RESERVED, 'node', 'access_key') # gen4 @user_transfer_spec.delete('destination_root') if @user_transfer_spec.key?('destination_root_id') elsif transfer_spec.key?('token') # gen3 # in that case, destination is set in return by application (API/upload_setup) # but to_folder was used in initial API call @user_transfer_spec.delete('destination_root') else # init default if required @user_transfer_spec['destination_root'] ||= destination_folder(transfer_spec['direction']) end end # update command line paths, unless destination already has one @user_transfer_spec['paths'] = transfer_spec['paths'] || ts_source_paths # updated transfer spec with command line transfer_spec.deep_merge!(@user_transfer_spec) # resolve pseudo-parameter: target_rate -> target_rate_kbps (overrides target_rate_kbps if both are present) transfer_spec['target_rate_kbps'] = Transfer::Spec.rate_string_to_kbps(transfer_spec.delete('target_rate')) if transfer_spec.key?('target_rate') # recursively remove values that are nil (user wants to delete) transfer_spec.deep_do{ |hash, key, value, _unused| hash.delete(key) if value.nil?} # if TS from app has content_protection (e.g. F5), that means content is protected: ask password if not provided transfer_spec['content_protection_password'] = @context..prompt_user_input('content protection password', sensitive: true) if transfer_spec['content_protection'].eql?('decrypt') && !transfer_spec.key?('content_protection_password') # create transfer agent agent_instance.start_transfer(transfer_spec, token_regenerator: rest_token) # --- async mode --- if @transfer_options['asynchronous'] agent_type = (@transfer_options['agent'] || 'direct').to_s raise Cli::BadArgument, 'asynchronous mode is not supported for agent httpgw' if agent_type.eql?('httpgw') # For direct agent: the transfer lives in this Ruby process (threads). # We return the job_id from the agent itself; no persistent store needed. if agent_type.eql?('direct') job_id = agent_instance.instance_variable_get(:@sessions).last[:job_id] Log.log.info{"Async direct transfer started: job_id=#{job_id}"} return Transfer::Result.async(job_id: job_id) end # For remote-daemon agents (desktop, node, connect, transferd): persist transfer_id # so the status can be re-queried in a later process invocation. agent_params = @transfer_options.except('agent', 'asynchronous') case agent_type when 'desktop' agent_params['application_id'] = agent_instance.application_id when 'connect' agent_params['app_id'] = agent_instance.app_id when 'transferd' # store the resolved daemon endpoint (auto-port may have been assigned) agent_params['url'] = agent_instance.daemon_endpoint end job_id = SecureRandom.uuid async_store.write(job_id, { 'job_id' => job_id, 'agent_type' => agent_type, 'transfer_id' => agent_instance.instance_variable_get(:@transfer_id).to_s, 'agent_params' => agent_params, 'status' => 'running', 'bytes_transferred' => 0, 'started_at' => Time.now.utc.iso8601, 'ended_at' => nil, 'error' => nil }) Log.log.info{"Async transfer started: job_id=#{job_id}"} return Transfer::Result.async(job_id: job_id) end # --- synchronous mode (default) --- result = agent_instance.wait_for_completion @notification_cb&.call(transfer_spec, result) return result end |
#ts_source_paths(default: nil) ⇒ Array?
This is how the list of files to be transferred is specified get paths suitable for transfer spec from command line computation is done only once, cache is kept in @transfer_paths
196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 |
# File 'lib/aspera/cli/transfer_agent.rb', line 196 def ts_source_paths(default: nil) # return cache if set return @transfer_paths unless @transfer_paths.nil? # start with lower priority : get paths from transfer spec on command line @transfer_paths = @user_transfer_spec['paths'] if @user_transfer_spec.key?('paths') # is there a source list option ? sources = @context..get_option(:sources) @transfer_paths = case sources when FILE_LIST_FROM_ARGS Log.log.debug('getting file list as parameters') Aspera.assert_type(default, Array, NilClass) # get remaining arguments list = @context..get_next_argument('source file list', multiple: true, default: default) raise Cli::BadArgument, 'specify at least one file on command line or use ' \ "--sources=#{FILE_LIST_FROM_TRANSFER_SPEC} to use transfer spec" if !list.is_a?(Array) || list.empty? list_to_paths(list) when FILE_LIST_FROM_TRANSFER_SPEC Log.log.debug('assume list provided in transfer spec') special_case_direct_with_list = (@transfer_options['agent'] || :direct).to_sym.eql?(:direct) && Transfer::Parameters.ascp_args_file_list?(@transfer_options['ascp_args']) Aspera.assert(!@transfer_paths.nil? || special_case_direct_with_list, type: Cli::BadArgument){'transfer spec on command line must have sources'} # can be nil @transfer_paths when Array Log.log.debug('getting file list as extended value') Aspera.assert_array_all(sources, String, type: Cli::BadArgument){'sources must be a Array of String'} list_to_paths(sources) else Aspera.error_unexpected_value(sources){'sources'} end Log.log.debug{"paths=#{@transfer_paths}"} return @transfer_paths end |