| Index: third_party/gsutil/boto/boto/storage_uri.py
|
| diff --git a/third_party/gsutil/20110627/boto/boto/storage_uri.py b/third_party/gsutil/boto/boto/storage_uri.py
|
| similarity index 66%
|
| rename from third_party/gsutil/20110627/boto/boto/storage_uri.py
|
| rename to third_party/gsutil/boto/boto/storage_uri.py
|
| index e683b04a3de4c133591a6741ec8eb59f77be0d82..9661c9f733e2a69fb1af359b5cde7a866ba999e2 100755
|
| --- a/third_party/gsutil/20110627/boto/boto/storage_uri.py
|
| +++ b/third_party/gsutil/boto/boto/storage_uri.py
|
| @@ -1,4 +1,5 @@
|
| # Copyright 2010 Google Inc.
|
| +# Copyright (c) 2011, Nexenta Systems Inc.
|
| #
|
| # Permission is hereby granted, free of charge, to any person obtaining a
|
| # copy of this software and associated documentation files (the
|
| @@ -40,6 +41,13 @@ class StorageUri(object):
|
| # https_connection_factory).
|
| connection_args = None
|
|
|
| + # Map of provider scheme ('s3' or 'gs') to AWSAuthConnection object. We
|
| + # maintain a pool here in addition to the connection pool implemented
|
| + # in AWSAuthConnection because the latter re-creates its connection pool
|
| + # every time that class is instantiated (so the current pool is used to
|
| + # avoid re-instantiating AWSAuthConnection).
|
| + provider_pool = {}
|
| +
|
| def __init__(self):
|
| """Uncallable constructor on abstract base StorageUri class.
|
| """
|
| @@ -56,8 +64,11 @@ class StorageUri(object):
|
|
|
| def check_response(self, resp, level, uri):
|
| if resp is None:
|
| - raise InvalidUriError('Attempt to get %s for "%s" failed. This '
|
| - 'probably indicates the URI is invalid' %
|
| + raise InvalidUriError('Attempt to get %s for "%s" failed.\nThis '
|
| + 'can happen if the URI refers to a non-'
|
| + 'existent object or if you meant to\noperate '
|
| + 'on a directory (e.g., leaving off -R option '
|
| + 'on gsutil cp, mv, or ls of a\nbucket)' %
|
| (level, uri))
|
|
|
| def connect(self, access_key_id=None, secret_access_key=None, **kwargs):
|
| @@ -70,7 +81,6 @@ class StorageUri(object):
|
| @rtype: L{AWSAuthConnection<boto.gs.connection.AWSAuthConnection>}
|
| @return: A connection to storage service provider of the given URI.
|
| """
|
| -
|
| connection_args = dict(self.connection_args or ())
|
| # Use OrdinaryCallingFormat instead of boto-default
|
| # SubdomainCallingFormat because the latter changes the hostname
|
| @@ -81,18 +91,26 @@ class StorageUri(object):
|
| # the resumable upload/download tests.
|
| from boto.s3.connection import OrdinaryCallingFormat
|
| connection_args['calling_format'] = OrdinaryCallingFormat()
|
| + if (hasattr(self, 'suppress_consec_slashes') and
|
| + 'suppress_consec_slashes' not in connection_args):
|
| + connection_args['suppress_consec_slashes'] = (
|
| + self.suppress_consec_slashes)
|
| connection_args.update(kwargs)
|
| if not self.connection:
|
| - if self.scheme == 's3':
|
| + if self.scheme in self.provider_pool:
|
| + self.connection = self.provider_pool[self.scheme]
|
| + elif self.scheme == 's3':
|
| from boto.s3.connection import S3Connection
|
| self.connection = S3Connection(access_key_id,
|
| secret_access_key,
|
| **connection_args)
|
| + self.provider_pool[self.scheme] = self.connection
|
| elif self.scheme == 'gs':
|
| from boto.gs.connection import GSConnection
|
| self.connection = GSConnection(access_key_id,
|
| secret_access_key,
|
| **connection_args)
|
| + self.provider_pool[self.scheme] = self.connection
|
| elif self.scheme == 'file':
|
| from boto.file.connection import FileConnection
|
| self.connection = FileConnection(self)
|
| @@ -173,8 +191,10 @@ class BucketStorageUri(StorageUri):
|
| Callers should instantiate this class by calling boto.storage_uri().
|
| """
|
|
|
| + delim = '/'
|
| +
|
| def __init__(self, scheme, bucket_name=None, object_name=None,
|
| - debug=0, connection_args=None):
|
| + debug=0, connection_args=None, suppress_consec_slashes=True):
|
| """Instantiate a BucketStorageUri from scheme,bucket,object tuple.
|
|
|
| @type scheme: string
|
| @@ -189,6 +209,8 @@ class BucketStorageUri(StorageUri):
|
| @param connection_args: optional map containing args to be
|
| passed to {S3,GS}Connection constructor (e.g., to override
|
| https_connection_factory).
|
| + @param suppress_consec_slashes: If provided, controls whether
|
| + consecutive slashes will be suppressed in key paths.
|
|
|
| After instantiation the components are available in the following
|
| fields: uri, scheme, bucket_name, object_name.
|
| @@ -199,6 +221,7 @@ class BucketStorageUri(StorageUri):
|
| self.object_name = object_name
|
| if connection_args:
|
| self.connection_args = connection_args
|
| + self.suppress_consec_slashes = suppress_consec_slashes
|
| if self.bucket_name and self.object_name:
|
| self.uri = ('%s://%s/%s' % (self.scheme, self.bucket_name,
|
| self.object_name))
|
| @@ -218,10 +241,13 @@ class BucketStorageUri(StorageUri):
|
| if not self.bucket_name:
|
| raise InvalidUriError('clone_replace_name() on bucket-less URI %s' %
|
| self.uri)
|
| - return BucketStorageUri(self.scheme, self.bucket_name, new_name,
|
| - self.debug)
|
| + return BucketStorageUri(
|
| + self.scheme, bucket_name=self.bucket_name, object_name=new_name,
|
| + debug=self.debug,
|
| + suppress_consec_slashes=self.suppress_consec_slashes)
|
|
|
| def get_acl(self, validate=True, headers=None, version_id=None):
|
| + """returns a bucket's acl"""
|
| if not self.bucket_name:
|
| raise InvalidUriError('get_acl on bucket-less URI (%s)' % self.uri)
|
| bucket = self.get_bucket(validate, headers)
|
| @@ -231,6 +257,33 @@ class BucketStorageUri(StorageUri):
|
| self.check_response(acl, 'acl', self.uri)
|
| return acl
|
|
|
| + def get_def_acl(self, validate=True, headers=None):
|
| + """returns a bucket's default object acl"""
|
| + if not self.bucket_name:
|
| + raise InvalidUriError('get_acl on bucket-less URI (%s)' % self.uri)
|
| + bucket = self.get_bucket(validate, headers)
|
| + # This works for both bucket- and object- level ACLs (former passes
|
| + # key_name=None):
|
| + acl = bucket.get_def_acl(self.object_name, headers)
|
| + self.check_response(acl, 'acl', self.uri)
|
| + return acl
|
| +
|
| + def get_cors(self, validate=True, headers=None):
|
| + """returns a bucket's CORS XML"""
|
| + if not self.bucket_name:
|
| + raise InvalidUriError('get_cors on bucket-less URI (%s)' % self.uri)
|
| + bucket = self.get_bucket(validate, headers)
|
| + cors = bucket.get_cors(headers)
|
| + self.check_response(cors, 'cors', self.uri)
|
| + return cors
|
| +
|
| + def set_cors(self, cors, validate=True, headers=None):
|
| + """sets or updates a bucket's CORS XML"""
|
| + if not self.bucket_name:
|
| + raise InvalidUriError('set_cors on bucket-less URI (%s)' % self.uri)
|
| + bucket = self.get_bucket(validate, headers)
|
| + bucket.set_cors(cors.to_xml(), headers)
|
| +
|
| def get_location(self, validate=True, headers=None):
|
| if not self.bucket_name:
|
| raise InvalidUriError('get_location on bucket-less URI (%s)' %
|
| @@ -238,6 +291,15 @@ class BucketStorageUri(StorageUri):
|
| bucket = self.get_bucket(validate, headers)
|
| return bucket.get_location()
|
|
|
| + def get_subresource(self, subresource, validate=True, headers=None,
|
| + version_id=None):
|
| + if not self.bucket_name:
|
| + raise InvalidUriError(
|
| + 'get_subresource on bucket-less URI (%s)' % self.uri)
|
| + bucket = self.get_bucket(validate, headers)
|
| + return bucket.get_subresource(subresource, self.object_name, headers,
|
| + version_id)
|
| +
|
| def add_group_email_grant(self, permission, email_address, recursive=False,
|
| validate=True, headers=None):
|
| if self.scheme != 'gs':
|
| @@ -255,8 +317,8 @@ class BucketStorageUri(StorageUri):
|
| bucket.add_group_email_grant(permission, email_address, recursive,
|
| headers)
|
| else:
|
| - raise InvalidUriError('add_group_email_grant() on bucket-less URI %s' %
|
| - self.uri)
|
| + raise InvalidUriError('add_group_email_grant() on bucket-less URI '
|
| + '%s' % self.uri)
|
|
|
| def add_email_grant(self, permission, email_address, recursive=False,
|
| validate=True, headers=None):
|
| @@ -292,21 +354,49 @@ class BucketStorageUri(StorageUri):
|
| bucket = self.get_bucket(headers)
|
| return bucket.list_grants(headers)
|
|
|
| + def is_file_uri(self):
|
| + """Returns True if this URI names a file or directory."""
|
| + return False
|
| +
|
| + def is_cloud_uri(self):
|
| + """Returns True if this URI names a bucket or object."""
|
| + return True
|
| +
|
| def names_container(self):
|
| - """Returns True if this URI names a bucket (vs. an object).
|
| """
|
| - return not self.object_name
|
| + Returns True if this URI names a directory or bucket. Will return
|
| + False for bucket subdirs; providing bucket subdir semantics needs to
|
| + be done by the caller (like gsutil does).
|
| + """
|
| + return bool(not self.object_name)
|
|
|
| def names_singleton(self):
|
| - """Returns True if this URI names an object (vs. a bucket).
|
| - """
|
| - return self.object_name
|
| + """Returns True if this URI names a file or object."""
|
| + return bool(self.object_name)
|
|
|
| - def is_file_uri(self):
|
| + def names_directory(self):
|
| + """Returns True if this URI names a directory."""
|
| return False
|
|
|
| - def is_cloud_uri(self):
|
| - return True
|
| + def names_provider(self):
|
| + """Returns True if this URI names a provider."""
|
| + return bool(not self.bucket_name)
|
| +
|
| + def names_bucket(self):
|
| + """Returns True if this URI names a bucket."""
|
| + return self.names_container()
|
| +
|
| + def names_file(self):
|
| + """Returns True if this URI names a file."""
|
| + return False
|
| +
|
| + def names_object(self):
|
| + """Returns True if this URI names an object."""
|
| + return self.names_singleton()
|
| +
|
| + def is_stream(self):
|
| + """Returns True if this URI represents input/output stream."""
|
| + return False
|
|
|
| def create_bucket(self, headers=None, location='', policy=None):
|
| if self.bucket_name is None:
|
| @@ -334,14 +424,25 @@ class BucketStorageUri(StorageUri):
|
|
|
| def set_acl(self, acl_or_str, key_name='', validate=True, headers=None,
|
| version_id=None):
|
| + """sets or updates a bucket's acl"""
|
| if not self.bucket_name:
|
| raise InvalidUriError('set_acl on bucket-less URI (%s)' %
|
| self.uri)
|
| self.get_bucket(validate, headers).set_acl(acl_or_str, key_name,
|
| headers, version_id)
|
|
|
| + def set_def_acl(self, acl_or_str, key_name='', validate=True, headers=None,
|
| + version_id=None):
|
| + """sets or updates a bucket's default object acl"""
|
| + if not self.bucket_name:
|
| + raise InvalidUriError('set_acl on bucket-less URI (%s)' %
|
| + self.uri)
|
| + self.get_bucket(validate, headers).set_def_acl(acl_or_str, key_name,
|
| + headers)
|
| +
|
| def set_canned_acl(self, acl_str, validate=True, headers=None,
|
| version_id=None):
|
| + """sets or updates a bucket's acl to a predefined (canned) value"""
|
| if not self.object_name:
|
| raise InvalidUriError('set_canned_acl on object-less URI (%s)' %
|
| self.uri)
|
| @@ -349,6 +450,26 @@ class BucketStorageUri(StorageUri):
|
| self.check_response(key, 'key', self.uri)
|
| key.set_canned_acl(acl_str, headers, version_id)
|
|
|
| + def set_def_canned_acl(self, acl_str, validate=True, headers=None,
|
| + version_id=None):
|
| + """sets or updates a bucket's default object acl to a predefined
|
| + (canned) value"""
|
| + if not self.object_name:
|
| + raise InvalidUriError('set_canned_acl on object-less URI (%s)' %
|
| + self.uri)
|
| + key = self.get_key(validate, headers)
|
| + self.check_response(key, 'key', self.uri)
|
| + key.set_def_canned_acl(acl_str, headers, version_id)
|
| +
|
| + def set_subresource(self, subresource, value, validate=True, headers=None,
|
| + version_id=None):
|
| + if not self.bucket_name:
|
| + raise InvalidUriError(
|
| + 'set_subresource on bucket-less URI (%s)' % self.uri)
|
| + bucket = self.get_bucket(validate, headers)
|
| + bucket.set_subresource(subresource, value, self.object_name, headers,
|
| + version_id)
|
| +
|
| def set_contents_from_string(self, s, headers=None, replace=True,
|
| cb=None, num_cb=10, policy=None, md5=None,
|
| reduced_redundancy=False):
|
| @@ -356,6 +477,21 @@ class BucketStorageUri(StorageUri):
|
| key.set_contents_from_string(s, headers, replace, cb, num_cb, policy,
|
| md5, reduced_redundancy)
|
|
|
| + def enable_logging(self, target_bucket, target_prefix=None, validate=True,
|
| + headers=None, version_id=None):
|
| + if not self.bucket_name:
|
| + raise InvalidUriError(
|
| + 'disable_logging on bucket-less URI (%s)' % self.uri)
|
| + bucket = self.get_bucket(validate, headers)
|
| + bucket.enable_logging(target_bucket, target_prefix, headers=headers)
|
| +
|
| + def disable_logging(self, validate=True, headers=None, version_id=None):
|
| + if not self.bucket_name:
|
| + raise InvalidUriError(
|
| + 'disable_logging on bucket-less URI (%s)' % self.uri)
|
| + bucket = self.get_bucket(validate, headers)
|
| + bucket.disable_logging(headers=headers)
|
| +
|
|
|
|
|
| class FileStorageUri(StorageUri):
|
| @@ -366,7 +502,9 @@ class FileStorageUri(StorageUri):
|
| See file/README about how we map StorageUri operations onto a file system.
|
| """
|
|
|
| - def __init__(self, object_name, debug):
|
| + delim = os.sep
|
| +
|
| + def __init__(self, object_name, debug, is_stream=False):
|
| """Instantiate a FileStorageUri from a path name.
|
|
|
| @type object_name: string
|
| @@ -384,6 +522,7 @@ class FileStorageUri(StorageUri):
|
| self.object_name = object_name
|
| self.uri = 'file://' + object_name
|
| self.debug = debug
|
| + self.stream = is_stream
|
|
|
| def clone_replace_name(self, new_name):
|
| """Instantiate a FileStorageUri from the current FileStorageUri,
|
| @@ -392,20 +531,52 @@ class FileStorageUri(StorageUri):
|
| @type new_name: string
|
| @param new_name: new object name
|
| """
|
| - return FileStorageUri(new_name, self.debug)
|
| + return FileStorageUri(new_name, self.debug, self.stream)
|
| +
|
| + def is_file_uri(self):
|
| + """Returns True if this URI names a file or directory."""
|
| + return True
|
| +
|
| + def is_cloud_uri(self):
|
| + """Returns True if this URI names a bucket or object."""
|
| + return False
|
|
|
| def names_container(self):
|
| - """Returns True if this URI names a directory.
|
| - """
|
| - return os.path.isdir(self.object_name)
|
| + """Returns True if this URI names a directory or bucket."""
|
| + return self.names_directory()
|
|
|
| def names_singleton(self):
|
| - """Returns True if this URI names a file.
|
| - """
|
| - return os.path.isfile(self.object_name)
|
| + """Returns True if this URI names a file (or stream) or object."""
|
| + return not self.names_container()
|
|
|
| - def is_file_uri(self):
|
| - return True
|
| + def names_directory(self):
|
| + """Returns True if this URI names a directory."""
|
| + if self.stream:
|
| + return False
|
| + return os.path.isdir(self.object_name)
|
|
|
| - def is_cloud_uri(self):
|
| + def names_provider(self):
|
| + """Returns True if this URI names a provider."""
|
| + return False
|
| +
|
| + def names_bucket(self):
|
| + """Returns True if this URI names a bucket."""
|
| + return False
|
| +
|
| + def names_file(self):
|
| + """Returns True if this URI names a file."""
|
| + return self.names_singleton()
|
| +
|
| + def names_object(self):
|
| + """Returns True if this URI names an object."""
|
| return False
|
| +
|
| + def is_stream(self):
|
| + """Returns True if this URI represents input/output stream.
|
| + """
|
| + return bool(self.stream)
|
| +
|
| + def close(self):
|
| + """Closes the underlying file.
|
| + """
|
| + self.get_key().close()
|
|
|