rename Block/ObjectStoreAdapter -> Block/ObjectStore

Signed-off-by: Steve Kriss <steve@heptio.com>
This commit is contained in:
Steve Kriss
2017-11-08 16:58:47 -08:00
parent 71bb702297
commit 21e2019540
12 changed files with 122 additions and 134 deletions
@@ -27,9 +27,7 @@ import (
"github.com/heptio/ark/pkg/cloudprovider"
)
var _ cloudprovider.BlockStorageAdapter = &blockStorageAdapter{}
type blockStorageAdapter struct {
type blockStore struct {
ec2 *ec2.EC2
}
@@ -46,7 +44,7 @@ func getSession(config *aws.Config) (*session.Session, error) {
return sess, nil
}
func NewBlockStorageAdapter(region string) (cloudprovider.BlockStorageAdapter, error) {
func NewBlockStore(region string) (cloudprovider.BlockStore, error) {
if region == "" {
return nil, errors.New("missing region in aws configuration in config file")
}
@@ -58,7 +56,7 @@ func NewBlockStorageAdapter(region string) (cloudprovider.BlockStorageAdapter, e
return nil, err
}
return &blockStorageAdapter{
return &blockStore{
ec2: ec2.New(sess),
}, nil
}
@@ -68,7 +66,7 @@ func NewBlockStorageAdapter(region string) (cloudprovider.BlockStorageAdapter, e
// from snapshot.
var iopsVolumeTypes = sets.NewString("io1")
func (op *blockStorageAdapter) CreateVolumeFromSnapshot(snapshotID, volumeType, volumeAZ string, iops *int64) (volumeID string, err error) {
func (op *blockStore) CreateVolumeFromSnapshot(snapshotID, volumeType, volumeAZ string, iops *int64) (volumeID string, err error) {
req := &ec2.CreateVolumeInput{
SnapshotId: &snapshotID,
AvailabilityZone: &volumeAZ,
@@ -87,7 +85,7 @@ func (op *blockStorageAdapter) CreateVolumeFromSnapshot(snapshotID, volumeType,
return *res.VolumeId, nil
}
func (op *blockStorageAdapter) GetVolumeInfo(volumeID, volumeAZ string) (string, *int64, error) {
func (op *blockStore) GetVolumeInfo(volumeID, volumeAZ string) (string, *int64, error) {
req := &ec2.DescribeVolumesInput{
VolumeIds: []*string{&volumeID},
}
@@ -119,7 +117,7 @@ func (op *blockStorageAdapter) GetVolumeInfo(volumeID, volumeAZ string) (string,
return volumeType, iops, nil
}
func (op *blockStorageAdapter) IsVolumeReady(volumeID, volumeAZ string) (ready bool, err error) {
func (op *blockStore) IsVolumeReady(volumeID, volumeAZ string) (ready bool, err error) {
req := &ec2.DescribeVolumesInput{
VolumeIds: []*string{&volumeID},
}
@@ -135,7 +133,7 @@ func (op *blockStorageAdapter) IsVolumeReady(volumeID, volumeAZ string) (ready b
return *res.Volumes[0].State == ec2.VolumeStateAvailable, nil
}
func (op *blockStorageAdapter) ListSnapshots(tagFilters map[string]string) ([]string, error) {
func (op *blockStore) ListSnapshots(tagFilters map[string]string) ([]string, error) {
req := &ec2.DescribeSnapshotsInput{}
for k, v := range tagFilters {
@@ -161,7 +159,7 @@ func (op *blockStorageAdapter) ListSnapshots(tagFilters map[string]string) ([]st
return ret, nil
}
func (op *blockStorageAdapter) CreateSnapshot(volumeID, volumeAZ string, tags map[string]string) (string, error) {
func (op *blockStore) CreateSnapshot(volumeID, volumeAZ string, tags map[string]string) (string, error) {
req := &ec2.CreateSnapshotInput{
VolumeId: &volumeID,
}
@@ -191,7 +189,7 @@ func (op *blockStorageAdapter) CreateSnapshot(volumeID, volumeAZ string, tags ma
return *res.SnapshotId, errors.WithStack(err)
}
func (op *blockStorageAdapter) DeleteSnapshot(snapshotID string) error {
func (op *blockStore) DeleteSnapshot(snapshotID string) error {
req := &ec2.DeleteSnapshotInput{
SnapshotId: &snapshotID,
}
@@ -29,15 +29,13 @@ import (
"github.com/heptio/ark/pkg/cloudprovider"
)
var _ cloudprovider.ObjectStorageAdapter = &objectStorageAdapter{}
type objectStorageAdapter struct {
type objectStore struct {
s3 *s3.S3
s3Uploader *s3manager.Uploader
kmsKeyID string
}
func NewObjectStorageAdapter(region, s3URL, kmsKeyID string, s3ForcePathStyle bool) (cloudprovider.ObjectStorageAdapter, error) {
func NewObjectStore(region, s3URL, kmsKeyID string, s3ForcePathStyle bool) (cloudprovider.ObjectStore, error) {
if region == "" {
return nil, errors.New("missing region in aws configuration in config file")
}
@@ -65,14 +63,14 @@ func NewObjectStorageAdapter(region, s3URL, kmsKeyID string, s3ForcePathStyle bo
return nil, err
}
return &objectStorageAdapter{
return &objectStore{
s3: s3.New(sess),
s3Uploader: s3manager.NewUploader(sess),
kmsKeyID: kmsKeyID,
}, nil
}
func (op *objectStorageAdapter) PutObject(bucket string, key string, body io.Reader) error {
func (op *objectStore) PutObject(bucket string, key string, body io.Reader) error {
req := &s3manager.UploadInput{
Bucket: &bucket,
Key: &key,
@@ -90,7 +88,7 @@ func (op *objectStorageAdapter) PutObject(bucket string, key string, body io.Rea
return errors.Wrapf(err, "error putting object %s", key)
}
func (op *objectStorageAdapter) GetObject(bucket string, key string) (io.ReadCloser, error) {
func (op *objectStore) GetObject(bucket string, key string) (io.ReadCloser, error) {
req := &s3.GetObjectInput{
Bucket: &bucket,
Key: &key,
@@ -104,7 +102,7 @@ func (op *objectStorageAdapter) GetObject(bucket string, key string) (io.ReadClo
return res.Body, nil
}
func (op *objectStorageAdapter) ListCommonPrefixes(bucket string, delimiter string) ([]string, error) {
func (op *objectStore) ListCommonPrefixes(bucket string, delimiter string) ([]string, error) {
req := &s3.ListObjectsV2Input{
Bucket: &bucket,
Delimiter: &delimiter,
@@ -125,7 +123,7 @@ func (op *objectStorageAdapter) ListCommonPrefixes(bucket string, delimiter stri
return ret, nil
}
func (op *objectStorageAdapter) ListObjects(bucket, prefix string) ([]string, error) {
func (op *objectStore) ListObjects(bucket, prefix string) ([]string, error) {
req := &s3.ListObjectsV2Input{
Bucket: &bucket,
Prefix: &prefix,
@@ -146,7 +144,7 @@ func (op *objectStorageAdapter) ListObjects(bucket, prefix string) ([]string, er
return ret, nil
}
func (op *objectStorageAdapter) DeleteObject(bucket string, key string) error {
func (op *objectStore) DeleteObject(bucket string, key string) error {
req := &s3.DeleteObjectInput{
Bucket: &bucket,
Key: &key,
@@ -157,7 +155,7 @@ func (op *objectStorageAdapter) DeleteObject(bucket string, key string) error {
return errors.Wrapf(err, "error deleting object %s", key)
}
func (op *objectStorageAdapter) CreateSignedURL(bucket, key string, ttl time.Duration) (string, error) {
func (op *objectStore) CreateSignedURL(bucket, key string, ttl time.Duration) (string, error) {
req, _ := op.s3.GetObjectRequest(&s3.GetObjectInput{
Bucket: aws.String(bucket),
Key: aws.String(key),
@@ -33,7 +33,7 @@ import (
"github.com/heptio/ark/pkg/cloudprovider"
)
type blockStorageAdapter struct {
type blockStore struct {
disks *disk.DisksClient
snaps *disk.SnapshotsClient
subscription string
@@ -42,8 +42,6 @@ type blockStorageAdapter struct {
apiTimeout time.Duration
}
var _ cloudprovider.BlockStorageAdapter = &blockStorageAdapter{}
const (
azureClientIDKey string = "AZURE_CLIENT_ID"
azureClientSecretKey string = "AZURE_CLIENT_SECRET"
@@ -72,7 +70,7 @@ func getConfig() map[string]string {
return cfg
}
func NewBlockStorageAdapter(location string, apiTimeout time.Duration) (cloudprovider.BlockStorageAdapter, error) {
func NewBlockStore(location string, apiTimeout time.Duration) (cloudprovider.BlockStore, error) {
if location == "" {
return nil, errors.New("missing location in azure configuration in config file")
}
@@ -120,7 +118,7 @@ func NewBlockStorageAdapter(location string, apiTimeout time.Duration) (cloudpro
return nil, errors.Errorf("location %q not found", location)
}
return &blockStorageAdapter{
return &blockStore{
disks: &disksClient,
snaps: &snapsClient,
subscription: cfg[azureSubscriptionIDKey],
@@ -130,7 +128,7 @@ func NewBlockStorageAdapter(location string, apiTimeout time.Duration) (cloudpro
}, nil
}
func (op *blockStorageAdapter) CreateVolumeFromSnapshot(snapshotID, volumeType, volumeAZ string, iops *int64) (string, error) {
func (op *blockStore) CreateVolumeFromSnapshot(snapshotID, volumeType, volumeAZ string, iops *int64) (string, error) {
fullSnapshotName := getFullSnapshotName(op.subscription, op.resourceGroup, snapshotID)
diskName := "restore-" + uuid.NewV4().String()
@@ -159,7 +157,7 @@ func (op *blockStorageAdapter) CreateVolumeFromSnapshot(snapshotID, volumeType,
return diskName, nil
}
func (op *blockStorageAdapter) GetVolumeInfo(volumeID, volumeAZ string) (string, *int64, error) {
func (op *blockStore) GetVolumeInfo(volumeID, volumeAZ string) (string, *int64, error) {
res, err := op.disks.Get(op.resourceGroup, volumeID)
if err != nil {
return "", nil, errors.WithStack(err)
@@ -168,7 +166,7 @@ func (op *blockStorageAdapter) GetVolumeInfo(volumeID, volumeAZ string) (string,
return string(res.AccountType), nil, nil
}
func (op *blockStorageAdapter) IsVolumeReady(volumeID, volumeAZ string) (ready bool, err error) {
func (op *blockStore) IsVolumeReady(volumeID, volumeAZ string) (ready bool, err error) {
res, err := op.disks.Get(op.resourceGroup, volumeID)
if err != nil {
return false, errors.WithStack(err)
@@ -181,7 +179,7 @@ func (op *blockStorageAdapter) IsVolumeReady(volumeID, volumeAZ string) (ready b
return *res.ProvisioningState == "Succeeded", nil
}
func (op *blockStorageAdapter) ListSnapshots(tagFilters map[string]string) ([]string, error) {
func (op *blockStore) ListSnapshots(tagFilters map[string]string) ([]string, error) {
res, err := op.snaps.ListByResourceGroup(op.resourceGroup)
if err != nil {
return nil, errors.WithStack(err)
@@ -215,7 +213,7 @@ Snapshot:
return ret, nil
}
func (op *blockStorageAdapter) CreateSnapshot(volumeID, volumeAZ string, tags map[string]string) (string, error) {
func (op *blockStore) CreateSnapshot(volumeID, volumeAZ string, tags map[string]string) (string, error) {
fullDiskName := getFullDiskName(op.subscription, op.resourceGroup, volumeID)
// snapshot names must be <= 80 characters long
var snapshotName string
@@ -258,7 +256,7 @@ func (op *blockStorageAdapter) CreateSnapshot(volumeID, volumeAZ string, tags ma
return snapshotName, nil
}
func (op *blockStorageAdapter) DeleteSnapshot(snapshotID string) error {
func (op *blockStore) DeleteSnapshot(snapshotID string) error {
ctx, cancel := context.WithTimeout(context.Background(), op.apiTimeout)
defer cancel()
@@ -29,13 +29,11 @@ import (
// ref. https://github.com/Azure-Samples/storage-blob-go-getting-started/blob/master/storageExample.go
type objectStorageAdapter struct {
type objectStore struct {
blobClient *storage.BlobStorageClient
}
var _ cloudprovider.ObjectStorageAdapter = &objectStorageAdapter{}
func NewObjectStorageAdapter() (cloudprovider.ObjectStorageAdapter, error) {
func NewObjectStore() (cloudprovider.ObjectStore, error) {
cfg := getConfig()
storageClient, err := storage.NewBasicClient(cfg[azureStorageAccountIDKey], cfg[azureStorageKeyKey])
@@ -45,12 +43,12 @@ func NewObjectStorageAdapter() (cloudprovider.ObjectStorageAdapter, error) {
blobClient := storageClient.GetBlobService()
return &objectStorageAdapter{
return &objectStore{
blobClient: &blobClient,
}, nil
}
func (op *objectStorageAdapter) PutObject(bucket string, key string, body io.Reader) error {
func (op *objectStore) PutObject(bucket string, key string, body io.Reader) error {
container, err := getContainerReference(op.blobClient, bucket)
if err != nil {
return err
@@ -64,7 +62,7 @@ func (op *objectStorageAdapter) PutObject(bucket string, key string, body io.Rea
return errors.WithStack(blob.CreateBlockBlobFromReader(body, nil))
}
func (op *objectStorageAdapter) GetObject(bucket string, key string) (io.ReadCloser, error) {
func (op *objectStore) GetObject(bucket string, key string) (io.ReadCloser, error) {
container, err := getContainerReference(op.blobClient, bucket)
if err != nil {
return nil, err
@@ -83,7 +81,7 @@ func (op *objectStorageAdapter) GetObject(bucket string, key string) (io.ReadClo
return res, nil
}
func (op *objectStorageAdapter) ListCommonPrefixes(bucket string, delimiter string) ([]string, error) {
func (op *objectStore) ListCommonPrefixes(bucket string, delimiter string) ([]string, error) {
container, err := getContainerReference(op.blobClient, bucket)
if err != nil {
return nil, err
@@ -108,7 +106,7 @@ func (op *objectStorageAdapter) ListCommonPrefixes(bucket string, delimiter stri
return ret, nil
}
func (op *objectStorageAdapter) ListObjects(bucket, prefix string) ([]string, error) {
func (op *objectStore) ListObjects(bucket, prefix string) ([]string, error) {
container, err := getContainerReference(op.blobClient, bucket)
if err != nil {
return nil, err
@@ -131,7 +129,7 @@ func (op *objectStorageAdapter) ListObjects(bucket, prefix string) ([]string, er
return ret, nil
}
func (op *objectStorageAdapter) DeleteObject(bucket string, key string) error {
func (op *objectStore) DeleteObject(bucket string, key string) error {
container, err := getContainerReference(op.blobClient, bucket)
if err != nil {
return err
@@ -147,7 +145,7 @@ func (op *objectStorageAdapter) DeleteObject(bucket string, key string) error {
const sasURIReadPermission = "r"
func (op *objectStorageAdapter) CreateSignedURL(bucket, key string, ttl time.Duration) (string, error) {
func (op *objectStore) CreateSignedURL(bucket, key string, ttl time.Duration) (string, error) {
container, err := getContainerReference(op.blobClient, bucket)
if err != nil {
return "", err
+23 -23
View File
@@ -99,35 +99,35 @@ func getRestoreResultsKey(backup, restore string) string {
}
type backupService struct {
objectStorage ObjectStorageAdapter
decoder runtime.Decoder
logger *logrus.Logger
objectStore ObjectStore
decoder runtime.Decoder
logger *logrus.Logger
}
var _ BackupService = &backupService{}
var _ BackupGetter = &backupService{}
// NewBackupService creates a backup service using the provided object storage adapter
func NewBackupService(objectStorage ObjectStorageAdapter, logger *logrus.Logger) BackupService {
// NewBackupService creates a backup service using the provided object store
func NewBackupService(objectStore ObjectStore, logger *logrus.Logger) BackupService {
return &backupService{
objectStorage: objectStorage,
decoder: scheme.Codecs.UniversalDecoder(api.SchemeGroupVersion),
logger: logger,
objectStore: objectStore,
decoder: scheme.Codecs.UniversalDecoder(api.SchemeGroupVersion),
logger: logger,
}
}
func (br *backupService) UploadBackup(bucket, backupName string, metadata, backup, log io.Reader) error {
// upload metadata file
metadataKey := getMetadataKey(backupName)
if err := br.objectStorage.PutObject(bucket, metadataKey, metadata); err != nil {
if err := br.objectStore.PutObject(bucket, metadataKey, metadata); err != nil {
// failure to upload metadata file is a hard-stop
return err
}
// upload tar file
if err := br.objectStorage.PutObject(bucket, getBackupContentsKey(backupName), backup); err != nil {
if err := br.objectStore.PutObject(bucket, getBackupContentsKey(backupName), backup); err != nil {
// try to delete the metadata file since the data upload failed
deleteErr := br.objectStorage.DeleteObject(bucket, metadataKey)
deleteErr := br.objectStore.DeleteObject(bucket, metadataKey)
return kerrors.NewAggregate([]error{err, deleteErr})
}
@@ -135,7 +135,7 @@ func (br *backupService) UploadBackup(bucket, backupName string, metadata, backu
// uploading log file is best-effort; if it fails, we log the error but call the overall upload a
// success
logKey := getBackupLogKey(backupName)
if err := br.objectStorage.PutObject(bucket, logKey, log); err != nil {
if err := br.objectStore.PutObject(bucket, logKey, log); err != nil {
br.logger.WithError(err).WithFields(logrus.Fields{
"bucket": bucket,
"key": logKey,
@@ -146,11 +146,11 @@ func (br *backupService) UploadBackup(bucket, backupName string, metadata, backu
}
func (br *backupService) DownloadBackup(bucket, backupName string) (io.ReadCloser, error) {
return br.objectStorage.GetObject(bucket, getBackupContentsKey(backupName))
return br.objectStore.GetObject(bucket, getBackupContentsKey(backupName))
}
func (br *backupService) GetAllBackups(bucket string) ([]*api.Backup, error) {
prefixes, err := br.objectStorage.ListCommonPrefixes(bucket, "/")
prefixes, err := br.objectStore.ListCommonPrefixes(bucket, "/")
if err != nil {
return nil, err
}
@@ -176,7 +176,7 @@ func (br *backupService) GetAllBackups(bucket string) ([]*api.Backup, error) {
func (br *backupService) GetBackup(bucket, name string) (*api.Backup, error) {
key := fmt.Sprintf(metadataFileFormatString, name)
res, err := br.objectStorage.GetObject(bucket, key)
res, err := br.objectStore.GetObject(bucket, key)
if err != nil {
return nil, err
}
@@ -201,7 +201,7 @@ func (br *backupService) GetBackup(bucket, name string) (*api.Backup, error) {
}
func (br *backupService) DeleteBackupDir(bucket, backupName string) error {
objects, err := br.objectStorage.ListObjects(bucket, backupName+"/")
objects, err := br.objectStore.ListObjects(bucket, backupName+"/")
if err != nil {
return err
}
@@ -212,7 +212,7 @@ func (br *backupService) DeleteBackupDir(bucket, backupName string) error {
"bucket": bucket,
"key": key,
}).Debug("Trying to delete object")
if err := br.objectStorage.DeleteObject(bucket, key); err != nil {
if err := br.objectStore.DeleteObject(bucket, key); err != nil {
errs = append(errs, err)
}
}
@@ -223,15 +223,15 @@ func (br *backupService) DeleteBackupDir(bucket, backupName string) error {
func (br *backupService) CreateSignedURL(target api.DownloadTarget, bucket string, ttl time.Duration) (string, error) {
switch target.Kind {
case api.DownloadTargetKindBackupContents:
return br.objectStorage.CreateSignedURL(bucket, getBackupContentsKey(target.Name), ttl)
return br.objectStore.CreateSignedURL(bucket, getBackupContentsKey(target.Name), ttl)
case api.DownloadTargetKindBackupLog:
return br.objectStorage.CreateSignedURL(bucket, getBackupLogKey(target.Name), ttl)
return br.objectStore.CreateSignedURL(bucket, getBackupLogKey(target.Name), ttl)
case api.DownloadTargetKindRestoreLog:
backup := extractBackupName(target.Name)
return br.objectStorage.CreateSignedURL(bucket, getRestoreLogKey(backup, target.Name), ttl)
return br.objectStore.CreateSignedURL(bucket, getRestoreLogKey(backup, target.Name), ttl)
case api.DownloadTargetKindRestoreResults:
backup := extractBackupName(target.Name)
return br.objectStorage.CreateSignedURL(bucket, getRestoreResultsKey(backup, target.Name), ttl)
return br.objectStore.CreateSignedURL(bucket, getRestoreResultsKey(backup, target.Name), ttl)
default:
return "", errors.Errorf("unsupported download target kind %q", target.Kind)
}
@@ -248,12 +248,12 @@ func extractBackupName(s string) string {
func (br *backupService) UploadRestoreLog(bucket, backup, restore string, log io.Reader) error {
key := getRestoreLogKey(backup, restore)
return br.objectStorage.PutObject(bucket, key, log)
return br.objectStore.PutObject(bucket, key, log)
}
func (br *backupService) UploadRestoreResults(bucket, backup, restore string, results io.Reader) error {
key := getRestoreResultsKey(backup, restore)
return br.objectStorage.PutObject(bucket, key, results)
return br.objectStore.PutObject(bucket, key, results)
}
// cachedBackupService wraps a real backup service with a cache for getting cloud backups.
+5 -5
View File
@@ -82,7 +82,7 @@ func TestUploadBackup(t *testing.T) {
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
var (
objStore = &testutil.ObjectStorageAdapter{}
objStore = &testutil.ObjectStore{}
bucket = "test-bucket"
backupName = "test-backup"
logger, _ = testlogger.NewNullLogger()
@@ -118,7 +118,7 @@ func TestUploadBackup(t *testing.T) {
func TestDownloadBackup(t *testing.T) {
var (
o = &testutil.ObjectStorageAdapter{}
o = &testutil.ObjectStore{}
bucket = "b"
backup = "bak"
logger, _ = testlogger.NewNullLogger()
@@ -158,7 +158,7 @@ func TestDeleteBackup(t *testing.T) {
bucket = "bucket"
backup = "bak"
objects = []string{"bak/ark-backup.json", "bak/bak.tar.gz", "bak/bak.log.gz"}
objStore = &testutil.ObjectStorageAdapter{}
objStore = &testutil.ObjectStore{}
logger, _ = testlogger.NewNullLogger()
)
@@ -230,7 +230,7 @@ func TestGetAllBackups(t *testing.T) {
t.Run(test.name, func(t *testing.T) {
var (
bucket = "bucket"
objStore = &testutil.ObjectStorageAdapter{}
objStore = &testutil.ObjectStore{}
logger, _ = testlogger.NewNullLogger()
)
@@ -327,7 +327,7 @@ func TestCreateSignedURL(t *testing.T) {
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
var (
objectStorage = &testutil.ObjectStorageAdapter{}
objectStorage = &testutil.ObjectStore{}
logger, _ = testlogger.NewNullLogger()
backupService = NewBackupService(objectStorage, logger)
)
@@ -31,14 +31,12 @@ import (
"github.com/heptio/ark/pkg/cloudprovider"
)
type blockStorageAdapter struct {
type blockStore struct {
gce *compute.Service
project string
}
var _ cloudprovider.BlockStorageAdapter = &blockStorageAdapter{}
func NewBlockStorageAdapter(project string) (cloudprovider.BlockStorageAdapter, error) {
func NewBlockStore(project string) (cloudprovider.BlockStore, error) {
if project == "" {
return nil, errors.New("missing project in gcp configuration in config file")
}
@@ -63,13 +61,13 @@ func NewBlockStorageAdapter(project string) (cloudprovider.BlockStorageAdapter,
return nil, errors.Errorf("error getting project %q", project)
}
return &blockStorageAdapter{
return &blockStore{
gce: gce,
project: project,
}, nil
}
func (op *blockStorageAdapter) CreateVolumeFromSnapshot(snapshotID, volumeType, volumeAZ string, iops *int64) (volumeID string, err error) {
func (op *blockStore) CreateVolumeFromSnapshot(snapshotID, volumeType, volumeAZ string, iops *int64) (volumeID string, err error) {
res, err := op.gce.Snapshots.Get(op.project, snapshotID).Do()
if err != nil {
return "", errors.WithStack(err)
@@ -88,7 +86,7 @@ func (op *blockStorageAdapter) CreateVolumeFromSnapshot(snapshotID, volumeType,
return disk.Name, nil
}
func (op *blockStorageAdapter) GetVolumeInfo(volumeID, volumeAZ string) (string, *int64, error) {
func (op *blockStore) GetVolumeInfo(volumeID, volumeAZ string) (string, *int64, error) {
res, err := op.gce.Disks.Get(op.project, volumeAZ, volumeID).Do()
if err != nil {
return "", nil, errors.WithStack(err)
@@ -97,7 +95,7 @@ func (op *blockStorageAdapter) GetVolumeInfo(volumeID, volumeAZ string) (string,
return res.Type, nil, nil
}
func (op *blockStorageAdapter) IsVolumeReady(volumeID, volumeAZ string) (ready bool, err error) {
func (op *blockStore) IsVolumeReady(volumeID, volumeAZ string) (ready bool, err error) {
disk, err := op.gce.Disks.Get(op.project, volumeAZ, volumeID).Do()
if err != nil {
return false, errors.WithStack(err)
@@ -107,7 +105,7 @@ func (op *blockStorageAdapter) IsVolumeReady(volumeID, volumeAZ string) (ready b
return disk.Status == "READY", nil
}
func (op *blockStorageAdapter) ListSnapshots(tagFilters map[string]string) ([]string, error) {
func (op *blockStore) ListSnapshots(tagFilters map[string]string) ([]string, error) {
useParentheses := len(tagFilters) > 1
subFilters := make([]string, 0, len(tagFilters))
@@ -134,7 +132,7 @@ func (op *blockStorageAdapter) ListSnapshots(tagFilters map[string]string) ([]st
return ret, nil
}
func (op *blockStorageAdapter) CreateSnapshot(volumeID, volumeAZ string, tags map[string]string) (string, error) {
func (op *blockStore) CreateSnapshot(volumeID, volumeAZ string, tags map[string]string) (string, error) {
// snapshot names must adhere to RFC1035 and be 1-63 characters
// long
var snapshotName string
@@ -180,7 +178,7 @@ func (op *blockStorageAdapter) CreateSnapshot(volumeID, volumeAZ string, tags ma
return gceSnap.Name, nil
}
func (op *blockStorageAdapter) DeleteSnapshot(snapshotID string) error {
func (op *blockStore) DeleteSnapshot(snapshotID string) error {
_, err := op.gce.Snapshots.Delete(op.project, snapshotID).Do()
return errors.WithStack(err)
@@ -31,15 +31,13 @@ import (
"github.com/heptio/ark/pkg/cloudprovider"
)
type objectStorageAdapter struct {
type objectStore struct {
gcs *storage.Service
googleAccessID string
privateKey []byte
}
var _ cloudprovider.ObjectStorageAdapter = &objectStorageAdapter{}
func NewObjectStorageAdapter(googleAccessID string, privateKey []byte) (cloudprovider.ObjectStorageAdapter, error) {
func NewObjectStore(googleAccessID string, privateKey []byte) (cloudprovider.ObjectStore, error) {
client, err := google.DefaultClient(oauth2.NoContext, storage.DevstorageReadWriteScope)
if err != nil {
return nil, errors.WithStack(err)
@@ -50,14 +48,14 @@ func NewObjectStorageAdapter(googleAccessID string, privateKey []byte) (cloudpro
return nil, errors.WithStack(err)
}
return &objectStorageAdapter{
return &objectStore{
gcs: gcs,
googleAccessID: googleAccessID,
privateKey: privateKey,
}, nil
}
func (op *objectStorageAdapter) PutObject(bucket string, key string, body io.Reader) error {
func (op *objectStore) PutObject(bucket string, key string, body io.Reader) error {
obj := &storage.Object{
Name: key,
}
@@ -67,7 +65,7 @@ func (op *objectStorageAdapter) PutObject(bucket string, key string, body io.Rea
return errors.WithStack(err)
}
func (op *objectStorageAdapter) GetObject(bucket string, key string) (io.ReadCloser, error) {
func (op *objectStore) GetObject(bucket string, key string) (io.ReadCloser, error) {
res, err := op.gcs.Objects.Get(bucket, key).Download()
if err != nil {
return nil, errors.WithStack(err)
@@ -76,7 +74,7 @@ func (op *objectStorageAdapter) GetObject(bucket string, key string) (io.ReadClo
return res.Body, nil
}
func (op *objectStorageAdapter) ListCommonPrefixes(bucket string, delimiter string) ([]string, error) {
func (op *objectStore) ListCommonPrefixes(bucket string, delimiter string) ([]string, error) {
res, err := op.gcs.Objects.List(bucket).Delimiter(delimiter).Do()
if err != nil {
return nil, errors.WithStack(err)
@@ -92,7 +90,7 @@ func (op *objectStorageAdapter) ListCommonPrefixes(bucket string, delimiter stri
return ret, nil
}
func (op *objectStorageAdapter) ListObjects(bucket, prefix string) ([]string, error) {
func (op *objectStore) ListObjects(bucket, prefix string) ([]string, error) {
res, err := op.gcs.Objects.List(bucket).Prefix(prefix).Do()
if err != nil {
return nil, errors.WithStack(err)
@@ -106,11 +104,11 @@ func (op *objectStorageAdapter) ListObjects(bucket, prefix string) ([]string, er
return ret, nil
}
func (op *objectStorageAdapter) DeleteObject(bucket string, key string) error {
func (op *objectStore) DeleteObject(bucket string, key string) error {
return errors.Wrapf(op.gcs.Objects.Delete(bucket, key).Do(), "error deleting object %s", key)
}
func (op *objectStorageAdapter) CreateSignedURL(bucket, key string, ttl time.Duration) (string, error) {
func (op *objectStore) CreateSignedURL(bucket, key string, ttl time.Duration) (string, error) {
if op.googleAccessID == "" {
return "", errors.New("unable to create a pre-signed URL - make sure GOOGLE_APPLICATION_CREDENTIALS points to a valid GCE service account file (missing email address)")
}
+10 -10
View File
@@ -56,20 +56,20 @@ const (
)
type snapshotService struct {
blockStorage BlockStorageAdapter
blockStore BlockStore
}
var _ SnapshotService = &snapshotService{}
// NewSnapshotService creates a snapshot service using the provided block storage adapter
func NewSnapshotService(blockStorage BlockStorageAdapter) SnapshotService {
// NewSnapshotService creates a snapshot service using the provided block store
func NewSnapshotService(blockStore BlockStore) SnapshotService {
return &snapshotService{
blockStorage: blockStorage,
blockStore: blockStore,
}
}
func (sr *snapshotService) CreateVolumeFromSnapshot(snapshotID string, volumeType string, volumeAZ string, iops *int64) (string, error) {
volumeID, err := sr.blockStorage.CreateVolumeFromSnapshot(snapshotID, volumeType, volumeAZ, iops)
volumeID, err := sr.blockStore.CreateVolumeFromSnapshot(snapshotID, volumeType, volumeAZ, iops)
if err != nil {
return "", err
}
@@ -85,7 +85,7 @@ func (sr *snapshotService) CreateVolumeFromSnapshot(snapshotID string, volumeTyp
case <-timeout.C:
return "", errors.Errorf("timeout reached waiting for volume %v to be ready", volumeID)
case <-ticker.C:
if ready, err := sr.blockStorage.IsVolumeReady(volumeID, volumeAZ); err == nil && ready {
if ready, err := sr.blockStore.IsVolumeReady(volumeID, volumeAZ); err == nil && ready {
return volumeID, nil
}
}
@@ -97,7 +97,7 @@ func (sr *snapshotService) GetAllSnapshots() ([]string, error) {
snapshotTagKey: snapshotTagVal,
}
res, err := sr.blockStorage.ListSnapshots(tags)
res, err := sr.blockStore.ListSnapshots(tags)
if err != nil {
return nil, err
}
@@ -110,13 +110,13 @@ func (sr *snapshotService) CreateSnapshot(volumeID, volumeAZ string) (string, er
snapshotTagKey: snapshotTagVal,
}
return sr.blockStorage.CreateSnapshot(volumeID, volumeAZ, tags)
return sr.blockStore.CreateSnapshot(volumeID, volumeAZ, tags)
}
func (sr *snapshotService) DeleteSnapshot(snapshotID string) error {
return sr.blockStorage.DeleteSnapshot(snapshotID)
return sr.blockStore.DeleteSnapshot(snapshotID)
}
func (sr *snapshotService) GetVolumeInfo(volumeID, volumeAZ string) (string, *int64, error) {
return sr.blockStorage.GetVolumeInfo(volumeID, volumeAZ)
return sr.blockStore.GetVolumeInfo(volumeID, volumeAZ)
}
+4 -4
View File
@@ -21,9 +21,9 @@ import (
"time"
)
// ObjectStorageAdapter exposes basic object-storage operations required
// ObjectStore exposes basic object-storage operations required
// by Ark.
type ObjectStorageAdapter interface {
type ObjectStore interface {
// PutObject creates a new object using the data in body within the specified
// object storage bucket with the given key.
PutObject(bucket string, key string, body io.Reader) error
@@ -48,9 +48,9 @@ type ObjectStorageAdapter interface {
CreateSignedURL(bucket, key string, ttl time.Duration) (string, error)
}
// BlockStorageAdapter exposes basic block-storage operations required
// BlockStore exposes basic block-storage operations required
// by Ark.
type BlockStorageAdapter interface {
type BlockStore interface {
// CreateVolumeFromSnapshot creates a new block volume, initialized from the provided snapshot,
// and with the specified type and IOPS (if using provisioned IOPS).
CreateVolumeFromSnapshot(snapshotID, volumeType, volumeAZ string, iops *int64) (volumeID string, err error)
+18 -18
View File
@@ -325,12 +325,12 @@ func (s *server) watchConfig(config *api.Config) {
func (s *server) initBackupService(config *api.Config) error {
s.logger.Info("Configuring cloud provider for backup service")
objectStorage, err := getObjectStorageProvider(config.BackupStorageProvider.CloudProviderConfig, "backupStorageProvider", s.logger)
objectStore, err := getObjectStore(config.BackupStorageProvider.CloudProviderConfig, "backupStorageProvider", s.logger)
if err != nil {
return err
}
s.backupService = cloudprovider.NewBackupService(objectStorage, s.logger)
s.backupService = cloudprovider.NewBackupService(objectStore, s.logger)
return nil
}
@@ -341,11 +341,11 @@ func (s *server) initSnapshotService(config *api.Config) error {
}
s.logger.Info("Configuring cloud provider for snapshot service")
blockStorage, err := getBlockStorageProvider(*config.PersistentVolumeProvider, "persistentVolumeProvider")
blockStore, err := getBlockStore(*config.PersistentVolumeProvider, "persistentVolumeProvider")
if err != nil {
return err
}
s.snapshotService = cloudprovider.NewSnapshotService(blockStorage)
s.snapshotService = cloudprovider.NewSnapshotService(blockStore)
return nil
}
@@ -373,10 +373,10 @@ func hasOneCloudProvider(cloudConfig api.CloudProviderConfig) bool {
return found
}
func getObjectStorageProvider(cloudConfig api.CloudProviderConfig, field string, logger *logrus.Logger) (cloudprovider.ObjectStorageAdapter, error) {
func getObjectStore(cloudConfig api.CloudProviderConfig, field string, logger *logrus.Logger) (cloudprovider.ObjectStore, error) {
var (
objectStorage cloudprovider.ObjectStorageAdapter
err error
objectStore cloudprovider.ObjectStore
err error
)
if !hasOneCloudProvider(cloudConfig) {
@@ -385,7 +385,7 @@ func getObjectStorageProvider(cloudConfig api.CloudProviderConfig, field string,
switch {
case cloudConfig.AWS != nil:
objectStorage, err = arkaws.NewObjectStorageAdapter(
objectStore, err = arkaws.NewObjectStore(
cloudConfig.AWS.Region,
cloudConfig.AWS.S3Url,
cloudConfig.AWS.KMSKeyID,
@@ -411,22 +411,22 @@ func getObjectStorageProvider(cloudConfig api.CloudProviderConfig, field string,
logger.Warning("GOOGLE_APPLICATION_CREDENTIALS is undefined; some features such as downloading log files will not work")
}
objectStorage, err = gcp.NewObjectStorageAdapter(email, privateKey)
objectStore, err = gcp.NewObjectStore(email, privateKey)
case cloudConfig.Azure != nil:
objectStorage, err = azure.NewObjectStorageAdapter()
objectStore, err = azure.NewObjectStore()
}
if err != nil {
return nil, err
}
return objectStorage, nil
return objectStore, nil
}
func getBlockStorageProvider(cloudConfig api.CloudProviderConfig, field string) (cloudprovider.BlockStorageAdapter, error) {
func getBlockStore(cloudConfig api.CloudProviderConfig, field string) (cloudprovider.BlockStore, error) {
var (
blockStorage cloudprovider.BlockStorageAdapter
err error
blockStore cloudprovider.BlockStore
err error
)
if !hasOneCloudProvider(cloudConfig) {
@@ -435,18 +435,18 @@ func getBlockStorageProvider(cloudConfig api.CloudProviderConfig, field string)
switch {
case cloudConfig.AWS != nil:
blockStorage, err = arkaws.NewBlockStorageAdapter(cloudConfig.AWS.Region)
blockStore, err = arkaws.NewBlockStore(cloudConfig.AWS.Region)
case cloudConfig.GCP != nil:
blockStorage, err = gcp.NewBlockStorageAdapter(cloudConfig.GCP.Project)
blockStore, err = gcp.NewBlockStore(cloudConfig.GCP.Project)
case cloudConfig.Azure != nil:
blockStorage, err = azure.NewBlockStorageAdapter(cloudConfig.Azure.Location, cloudConfig.Azure.APITimeout.Duration)
blockStore, err = azure.NewBlockStore(cloudConfig.Azure.Location, cloudConfig.Azure.APITimeout.Duration)
}
if err != nil {
return nil, err
}
return blockStorage, nil
return blockStore, nil
}
func durationMin(a, b time.Duration) time.Duration {
@@ -21,13 +21,13 @@ import io "io"
import mock "github.com/stretchr/testify/mock"
import time "time"
// ObjectStorageAdapter is an autogenerated mock type for the ObjectStorageAdapter type
type ObjectStorageAdapter struct {
// ObjectStore is an autogenerated mock type for the ObjectStore type
type ObjectStore struct {
mock.Mock
}
// CreateSignedURL provides a mock function with given fields: bucket, key, ttl
func (_m *ObjectStorageAdapter) CreateSignedURL(bucket string, key string, ttl time.Duration) (string, error) {
func (_m *ObjectStore) CreateSignedURL(bucket string, key string, ttl time.Duration) (string, error) {
ret := _m.Called(bucket, key, ttl)
var r0 string
@@ -48,7 +48,7 @@ func (_m *ObjectStorageAdapter) CreateSignedURL(bucket string, key string, ttl t
}
// DeleteObject provides a mock function with given fields: bucket, key
func (_m *ObjectStorageAdapter) DeleteObject(bucket string, key string) error {
func (_m *ObjectStore) DeleteObject(bucket string, key string) error {
ret := _m.Called(bucket, key)
var r0 error
@@ -62,7 +62,7 @@ func (_m *ObjectStorageAdapter) DeleteObject(bucket string, key string) error {
}
// GetObject provides a mock function with given fields: bucket, key
func (_m *ObjectStorageAdapter) GetObject(bucket string, key string) (io.ReadCloser, error) {
func (_m *ObjectStore) GetObject(bucket string, key string) (io.ReadCloser, error) {
ret := _m.Called(bucket, key)
var r0 io.ReadCloser
@@ -85,7 +85,7 @@ func (_m *ObjectStorageAdapter) GetObject(bucket string, key string) (io.ReadClo
}
// ListCommonPrefixes provides a mock function with given fields: bucket, delimiter
func (_m *ObjectStorageAdapter) ListCommonPrefixes(bucket string, delimiter string) ([]string, error) {
func (_m *ObjectStore) ListCommonPrefixes(bucket string, delimiter string) ([]string, error) {
ret := _m.Called(bucket, delimiter)
var r0 []string
@@ -108,7 +108,7 @@ func (_m *ObjectStorageAdapter) ListCommonPrefixes(bucket string, delimiter stri
}
// ListObjects provides a mock function with given fields: bucket, prefix
func (_m *ObjectStorageAdapter) ListObjects(bucket string, prefix string) ([]string, error) {
func (_m *ObjectStore) ListObjects(bucket string, prefix string) ([]string, error) {
ret := _m.Called(bucket, prefix)
var r0 []string
@@ -131,7 +131,7 @@ func (_m *ObjectStorageAdapter) ListObjects(bucket string, prefix string) ([]str
}
// PutObject provides a mock function with given fields: bucket, key, body
func (_m *ObjectStorageAdapter) PutObject(bucket string, key string, body io.Reader) error {
func (_m *ObjectStore) PutObject(bucket string, key string, body io.Reader) error {
ret := _m.Called(bucket, key, body)
var r0 error