"""COPIED FROM garcon-contrib 0.4.4rc0.""" from ftplib import FTP import os from boto.s3.connection import Bucket from boto.s3.connection import S3Connection from garcon import task import paramiko def _extract_bucket_path(url): """Extract the bucket name and path from the provided S3 url Args: url (str): S3 URL like 's3://bucket_name/folder1/folder2/' Return: tuple: bucket name and bucket path """ bucket_path = '' bucket_name = '' if url.startswith('s3://'): split_path = url.split('/', 3) bucket_name = split_path[2] if len(split_path) > 3: bucket_path = split_path[3] if not bucket_name: raise Exception('The S3 url \'{url}\' is not valid.'.format(url=url)) return (bucket_name, bucket_path) def _upload_from_local_to_ftp( sftp_creds, file_name, remote_ftp_dir, local_dir, pkey=None): """Upload files from local file system to FTP Args: sftp_creds (dict): credentials for connecting to ftp server file_name (str): name of the file to be uploaded remote_ftp_dir (str): destination directory on FTP local_dir (str): local directory where the file is located pkey (str): path to ssh private key """ host = sftp_creds['host'] port = int(sftp_creds['port']) username = sftp_creds['username'] remote_ftp_dir = remote_ftp_dir.rstrip('/') if pkey or port != 21: transport = paramiko.Transport((host, port)) if pkey: transport.connect( username=username, pkey=paramiko.RSAKey.from_private_key_file(pkey)) else: transport.connect( username=username, password=sftp_creds['password']) sftp = paramiko.SFTPClient.from_transport(transport) dir_list = remote_ftp_dir.strip('/').split('/') full_path = '' for current_dir in dir_list: try: full_path = '{previous}/{current}'.format( previous=full_path, current=current_dir) sftp.chdir(full_path) except (IOError, FileExistsError): sftp.mkdir(full_path) sftp.chdir(full_path) sftp.put( '{}/{}'.format(local_dir, file_name), '/{}/{}'.format(remote_ftp_dir, file_name)) sftp.close() transport.close() else: ftp = FTP(host) ftp.login(username, sftp_creds['password']) try: ftp.mkd(remote_ftp_dir) except Exception: pass ftp.cwd(remote_ftp_dir) with open('{}/{}'.format(local_dir, file_name), 'rb') as fp: ftp.storbinary( 'STOR {}'.format(file_name), fp, 1024) ftp.close() @task.decorate(timeout=7200) def copy_from_s3_to_ftp( activity, sftp_creds, file_names_list, remote_s3_dir_path, remote_ftp_dir_path, pkey): """Copy files from S3 to FTP Args: activity (ActivityWorker): The swf activity worker. sftp_creds (dict): sftp connection credentials dictionary file_names_list (list): list of file names remote_s3_dir_path (str): s3 path to the folder which contains the source file remote_ftp_dir_path (str): path to the destination sftp location pkey (Optional[str]): path to private key. If present, use key, otherwise use username and password """ remote_s3_dir_path = remote_s3_dir_path.rstrip('/') bucket_name, bucket_path = _extract_bucket_path(remote_s3_dir_path) bucket = Bucket(S3Connection(), bucket_name) for file_name in file_names_list: local_dir = './' remote_dir = '{remote_ftp_path}/'.format( remote_ftp_path=remote_ftp_dir_path) key = bucket.get_key('{bucket_path}/{file_name}'.format( bucket_path=bucket_path, file_name=file_name)) key.get_contents_to_filename('{}/{}'.format(local_dir, file_name)) _upload_from_local_to_ftp( sftp_creds, file_name, remote_dir, local_dir, pkey) # Remove file from local dir when done uploading os.remove('{}/{}'.format(local_dir, file_name))