From dff3a533a3ac9f22e0684ba56b3633196164acc6 Mon Sep 17 00:00:00 2001 From: nmahalle Date: Fri, 12 Feb 2021 00:05:30 -0800 Subject: [PATCH 1/8] New genric connector for cloud --- .../nzconnector/connector/AZConnector.go | 373 ++++++++++++++++++ bnr-utils/nzconnector/connector/Connector.go | 9 + .../nzconnector/connector/S3Connector.go | 56 +++ bnr-utils/nzconnector/main.go | 126 ++++++ 4 files changed, 564 insertions(+) create mode 100644 bnr-utils/nzconnector/connector/AZConnector.go create mode 100644 bnr-utils/nzconnector/connector/Connector.go create mode 100644 bnr-utils/nzconnector/connector/S3Connector.go create mode 100644 bnr-utils/nzconnector/main.go diff --git a/bnr-utils/nzconnector/connector/AZConnector.go b/bnr-utils/nzconnector/connector/AZConnector.go new file mode 100644 index 0000000..70fc8ad --- /dev/null +++ b/bnr-utils/nzconnector/connector/AZConnector.go @@ -0,0 +1,373 @@ +package Connector + +import ( + "fmt" + "strings" + "strconv" + "net/url" + "time" + "os" + "path/filepath" + "path" + "context" + "log" + "github.com/Azure/azure-storage-blob-go/azblob" +) + +type iaz interface { + azargs() + getServiceURL() (azblob.ServiceURL, error) + getContainerURL() (azblob.ContainerURL, error) + getBlockBlobURL(blobname string) (azblob.BlockBlobURL, error) + getBlobURL(blobname string) (azblob.BlobURL, error) + uploadFile(absfilepath string, relfilepath string, uniqueid string) (error) + downloadFile(outfilepath string, blobname string, streams uint, blockSize int64) error + +} + +type AZConnector struct { + azaccount string + azkey string + azcontainer string + streams uint + blocksize int64 +} + +type job struct { + uniqueid string + bkpdir string +} + +type uploadJob struct { + job + absfilepath string +} + +type jobResult struct { + *job + err error +} + +type downloadJob struct { + conn AZConnector + blobname string + outfilepath string +} + +type downloadJobResult struct { + blobname string + err error +} + +func (c *AZConnector) ParseConnectorArgs(args string) { + arguments := strings.Split(args, ";") + for _, arg := range arguments { + kv := strings.Split(arg, ":") + switch kv[0] { + case "STORAGE_ACCOUNT": + c.azaccount = kv[1] + case "KEY": + c.azkey = kv[1] + case "CONTAINER": + c.azcontainer = kv[1] + case "STREAMS": + u32, err := strconv.ParseUint(kv[1], 10, 32) + if (err == nil ) { + c.streams = uint(u32) + } + case "BLOCKSIZE": + u64, err := strconv.ParseInt(kv[1], 10, 64) + if (err == nil ) { + c.blocksize = int64(u64) + } + } + } +} + +func (j *uploadJob) upload(conn *AZConnector) error { + relfilepath, err := filepath.Rel(j.job.bkpdir, j.absfilepath) + if err != nil { + return fmt.Errorf("Unable to traverse %s, %s: %v", j.job.bkpdir, j.absfilepath, err) + } + + log.Println("Uploading file :", j.absfilepath) + return conn.uploadFile(j.absfilepath, relfilepath, j.job.uniqueid) +} + +func (t *AZConnector) Upload() { + t.azargs() + fmt.Println("Uploading with az connector") +} + +func (cn *AZConnector) getServiceURL() (azblob.ServiceURL, error) { + var serviceURL azblob.ServiceURL + us := fmt.Sprintf("https://%s.blob.core.windows.net/", cn.azaccount) + u, err := url.Parse(us) + if err != nil { + return serviceURL, fmt.Errorf("Unable to parse URL: %s : %v", us, err) + } + + credential, err := azblob.NewSharedKeyCredential(cn.azaccount, cn.azkey) + if err != nil { + return serviceURL, fmt.Errorf("Unable to create shared credentials: %v", err) + } + + p := azblob.NewPipeline(credential, azblob.PipelineOptions{ + Retry: azblob.RetryOptions{ + TryTimeout: 5 * time.Minute, + }, + }) + + serviceURL = azblob.NewServiceURL(*u, p) + return serviceURL, nil +} + +func (cn *AZConnector) getContainerURL() (azblob.ContainerURL, error) { + var containerURL azblob.ContainerURL + serviceURL, err := cn.getServiceURL() + if err == nil { + containerURL = serviceURL.NewContainerURL(cn.azcontainer) + } + return containerURL, err +} + +func (cn *AZConnector) getBlobURL(blobname string) (azblob.BlobURL, error) { + var blobURL azblob.BlobURL + containerURL, err := cn.getContainerURL() + if err == nil { + blobURL = containerURL.NewBlobURL(blobname) + } + return blobURL, err +} + +func (cn *AZConnector) getBlockBlobURL(blobname string) (azblob.BlockBlobURL, error) { + var blockBlobURL azblob.BlockBlobURL + containerURL, err := cn.getContainerURL() + if err == nil { + blockBlobURL = containerURL.NewBlockBlobURL(blobname) + } + return blockBlobURL, err +} + +func (cn *AZConnector) uploadFile(absfilepath string, relfilepath string, uniqueid string) (error){ + // Upload the file to a block blob + blockBlobURL, err := cn.getBlockBlobURL(uniqueid+"/"+relfilepath) + if err != nil { + return err + } + + file, err := os.Open(absfilepath) + if err != nil { + return fmt.Errorf("Error in opening backup file : %v", err) + } + + _, err = azblob.UploadFileToBlockBlob(context.Background(), file, blockBlobURL, + azblob.UploadToBlockBlobOptions{ + BlockSize: int64(cn.blocksize * 1024 * 1024), + Parallelism: uint16(cn.streams), + }) + + return err +} + + +func (cn *AZConnector) UploadBkp(bkpdir string, uniqueid string, backupdir string, paralleljobs int) (error){ + var err error + work := make(chan *uploadJob, paralleljobs) + result := make(chan *jobResult, paralleljobs) + done := make(chan bool) + + go func() { + for { + select { + case j, ok := <- work: + if ! ok { + // done + close(result) + return + } + err := j.upload(cn) + jr := jobResult{ job:&j.job, err:err } + result <- &jr + } + } + }() + + filesuploaded := 0 + go func() { + for { + select { + case r, ok := <- result: + if ! ok { + // work done + done <- true + return + } + if r.err != nil { + // stopping right here so that we + // don't keep on uploading when one has failed + log.Fatalf("%s: %v", r.job, r.err) + } + filesuploaded++ // this is fine, since this is single threaded increment + } + } + }() + + err = filepath.Walk(backupdir, + func(absfilepath string, info os.FileInfo, err error) error { + if info.IsDir() { + return nil + } + j := uploadJob{ job: job{uniqueid, bkpdir}, absfilepath: absfilepath } + work <- &j // this will hang until at least one of the prior uploads finish if other.paralleljobs + // are already running + return err + }) + close(work) + <- done + log.Println("Upload successful. Total files uploaded:", filesuploaded) + return err + } + +func (j *downloadJob) download(conn *AZConnector) error { + + log.Println("Downloading file :", j.blobname) + return conn.downloadFile(j.outfilepath, j.blobname, conn.streams, conn.blocksize) +} + +func (cn *AZConnector) downloadFile(outfilepath string, blobname string, streams uint, blockSize int64) error { + + filehandle, err := os.Create(outfilepath) + if err != nil { + return fmt.Errorf("Error in creating file inside backup dir: %v",err) + } + + defer filehandle.Close() + + blobURL, err := cn.getBlobURL(blobname) + if err != nil { + return err + } + + // Perform download + err = azblob.DownloadBlobToFile(context.Background(), blobURL, 0, 0, filehandle, + azblob.DownloadFromBlobOptions{ + BlockSize: int64(blockSize * 1024 * 1024), + RetryReaderOptionsPerBlock: azblob.RetryReaderOptions{MaxRetryRequests: 20}, + Parallelism: uint16(streams), + }) + if err != nil { + return fmt.Errorf("Error in downloading an Azure blob to a file: %v",err) + } + return err +} + + +func (cn *AZConnector) DownloadBkp(outdir string, uniqueid string, blobpath string, paralleljobs int) (error){ + var err error + work := make(chan *downloadJob, paralleljobs) + result := make(chan *downloadJobResult, paralleljobs) + done := make(chan bool) + + // start the workers + go func() { + for { + select { + case j, ok := <- work: + if ! ok { + // done + close(result) + return + } + err := j.download(cn) + jr := downloadJobResult{ blobname:j.blobname, err:err } + result <- &jr + } + } + }() + + filesdownloaded := 0 + go func() { + for { + select { + case r, ok := <- result: + if ! ok { + // work done + done <- true + return + } + if r.err != nil { + // stopping right here so that we + // don't keep on uploading when one has failed + log.Fatalf("%s: %v", r.blobname, r.err) + } + filesdownloaded++ // this is fine, since this is single threaded increment + } + } + }() + + + blobfound := 0 + + for marker := (azblob.Marker{}); marker.NotDone(); { + containerURL, err := cn.getContainerURL() + if err != nil { + return err + } + + // Get a result segment starting with the blob indicated by the current Marker. + listBlob, err := containerURL.ListBlobsFlatSegment(context.Background(), marker, azblob.ListBlobsSegmentOptions{}) + if err != nil { + return fmt.Errorf("Error in listing segment of blobs: %v",err) + } + + // ListBlobs returns the start of the next segment; you MUST use this to get + // the next segment (after processing the current result segment). + marker = listBlob.NextMarker + // Process the blobs returned in this result segment (if the segment is empty, the loop body won't execute) + for _, blobInfo := range listBlob.Segment.BlobItems { + if strings.HasPrefix(blobInfo.Name, blobpath) { + + // Set up file to download the blob to + dir, filename := filepath.Split(blobInfo.Name) + + relfilepath, err := filepath.Rel(uniqueid,dir) + if err != nil { + return fmt.Errorf("Error in fetching download relative path: %v",err) + } + + dumpdir := filepath.Join(outdir, relfilepath) + err = os.MkdirAll(dumpdir, 0777) + if err != nil { + return fmt.Errorf("Error in creating backup directory structure: %v",err) + } + + outfilepath := path.Join(dumpdir, filename) + j := downloadJob{ conn:*cn, outfilepath:outfilepath, blobname: blobInfo.Name } + work <- &j + + blobfound-- + } + } + + if blobfound > 0 { + log.Println("No matching blob found. Please check if DB name, hostname, uniqueid or containername is correct") + return fmt.Errorf("No matching blob found.") + } + if blobfound == 0 { + blobfound++ + } + } + close(work) + <- done + log.Println("Total files downloaded:", filesdownloaded) + return err +} + + +func (t AZConnector) azargs() { + fmt.Println("STORAGE_ACCOUNT : ", t.azaccount) + fmt.Println("STORAGE_ACCOUNT : ", t.blocksize) + fmt.Println("STORAGE_ACCOUNT : ", t.streams) +} + diff --git a/bnr-utils/nzconnector/connector/Connector.go b/bnr-utils/nzconnector/connector/Connector.go new file mode 100644 index 0000000..f2c7623 --- /dev/null +++ b/bnr-utils/nzconnector/connector/Connector.go @@ -0,0 +1,9 @@ +package Connector + +type IConnector interface { + ParseConnectorArgs(string) + UploadBkp(string, string, string, int) (error) + Upload() + DownloadBkp(string, string, string, int) (error) +} + diff --git a/bnr-utils/nzconnector/connector/S3Connector.go b/bnr-utils/nzconnector/connector/S3Connector.go new file mode 100644 index 0000000..53baeb9 --- /dev/null +++ b/bnr-utils/nzconnector/connector/S3Connector.go @@ -0,0 +1,56 @@ +package Connector + +import ( + "fmt" + "strings" +) + +type is3 interface { + s3args() +} + +type S3connector struct { + access_key_id string + secret_access_key string + default_region string + bucket_url string + endpoint string + uniqueid string +} + +func (c *S3connector) Upload() { + c.s3args() + fmt.Println("Uploading with s3 connector") +} + +func (c *S3connector) UploadBkp(bkpdir string, uniqueid string, backupdir string, paralleljobs int) (error){ + fmt.Println("Uploading") + return nil +} +func (cn *S3connector) DownloadBkp(outdir string, uniqueid string, blobpath string, paralleljobs int) (error){ + fmt.Println("Downloading") + return nil +} +func (c *S3connector) ParseConnectorArgs(args string) { + arguments := strings.Split(args, ":") + for _, arg := range arguments { + kv := strings.Split(arg, "=") + switch kv[0] { + case "ACCESS_KEY_ID": + c.access_key_id = kv[1] + case "SECRET_ACCESS_KEY": + c.secret_access_key = kv[1] + case "DEFAULT_REGION": + c.default_region = kv[1] + case "BUCKET_URL": + c.bucket_url = kv[1] + case "ENDPOINT": + c.endpoint = kv[1] + } + } +} + +func (c S3connector) s3args() { + fmt.Println("ACCESS_KEY_ID : ", c.access_key_id) +} + diff --git a/bnr-utils/nzconnector/main.go b/bnr-utils/nzconnector/main.go new file mode 100644 index 0000000..a58b4e0 --- /dev/null +++ b/bnr-utils/nzconnector/main.go @@ -0,0 +1,126 @@ +package main + +import ( + "flag" + "nzconnector/connector" + "fmt" + "path/filepath" + "log" + "os" + "time" + "path" + "strings" +) + +type BackupInfo struct { + dbname string + dir string + npshost string + backupset string +} + +type ConnectorInfo struct { + connector string + connectorArgs string +} + +type OtherArgs struct { + uniqueid string + logfiledir string + upload *bool + download *bool + paralleljobs int +} + + +func parseArgs(backupinfo *BackupInfo, connectorInfo *ConnectorInfo, otherargs *OtherArgs) { + flag.StringVar(&backupinfo.dbname, "db", "", "Database name") + flag.StringVar(&backupinfo.dir, "dir", "", "Full path to the directory in which the backup already exists or should be downloaded") + flag.StringVar(&backupinfo.npshost, "npshost", "", "Name of the NPS host as it appears in the backups") + flag.StringVar(&backupinfo.backupset, "backupset", "", "Name of the backupset to be uploaded/downloaded") + + flag.StringVar(&connectorInfo.connector, "connector", "", "Destination cloud store") + flag.StringVar(&connectorInfo.connectorArgs, "connectorArgs", "", "Arguments for cloud store") + + flag.StringVar(&otherargs.uniqueid,"uniqueid", "", "Azure blob storage container") + flag.StringVar(&otherargs.logfiledir,"logfiledir", "/tmp", "Logfile directory for this utility. Default is /tmp dir") + otherargs.upload = flag.Bool("upload", false, "Upload to cloud") + otherargs.download = flag.Bool("download", false, "Download from cloud") + flag.IntVar(&otherargs.paralleljobs,"paralleljobs",6,"Number of parallel files to upload/download") +} + +func handleErrors(err error) { + if err != nil { + log.Fatalln(err) + } +} + +func parseConnectorArgs(e Connector.IConnector, args string) { + if e != nil { + e.ParseConnectorArgs(args) + } +} + +func GetConnector(connectorType string) Connector.IConnector { + switch connectorType { + case "s3": + return &Connector.S3connector{} + case "az": + return &Connector.AZConnector{} + default: + return nil + } +} + +func main() { + var backupinfo BackupInfo + var connectorInfo ConnectorInfo + var otherargs OtherArgs + // parse input args + parseArgs(&backupinfo, &connectorInfo, &otherargs) + flag.Parse() + + connector := GetConnector(connectorInfo.connector) + parseConnectorArgs(connector, connectorInfo.connectorArgs) + + // log file configuration setup + logfilename := fmt.Sprintf("nz_azConnector_%d_%s.log", os.Getppid(), time.Now().Format("2006-01-02")) + logfilepath := path.Join(otherargs.logfiledir, logfilename) + filehandle, err := os.OpenFile(logfilepath, os.O_WRONLY|os.O_CREATE|os.O_APPEND, 0644) + if err != nil { + fmt.Errorf("Error in opening logfile: %v",err) + } + log.SetOutput(filehandle) + prefixStr := fmt.Sprintf("%s ", time.Now().UTC().Format("2006-01-02 15:04:05 EST")) + fmt.Sprintf("%-7s", "[INFO]") + log.SetFlags(0) + log.SetPrefix(prefixStr) + + dirlist := strings.Split(backupinfo.dir," ") + log.Println("Backup/Restore directory :",dirlist) + log.Println("DB name :", backupinfo.dbname) + log.Println("Nps hostname :", backupinfo.npshost) + log.Println("BackupsetID :", backupinfo.backupset) + log.Println("UniqueID :", otherargs.uniqueid) + log.Println("Number of files to upload/download in parallel :", otherargs.paralleljobs) + + for _, bkpdir := range dirlist { + if (*otherargs.upload) { + + // now do the upload + log.Println("Uploading backup data to azure cloud from backup dir", bkpdir) + backupdir := filepath.Join(bkpdir, "Netezza", backupinfo.npshost, backupinfo.dbname, backupinfo.backupset) + _, err = os.Stat(backupdir) + handleErrors(err) + err = connector.UploadBkp(bkpdir, otherargs.uniqueid, backupdir, otherargs.paralleljobs) + handleErrors(err) + log.Println("Upload successful") + } + if (*otherargs.download) { + log.Println("Downloading backup data from azure cloud to restore dir", bkpdir) + blobpath := filepath.Join(otherargs.uniqueid, "Netezza",backupinfo.npshost, backupinfo.dbname, backupinfo.backupset) + err = connector.DownloadBkp(bkpdir, otherargs.uniqueid, blobpath, otherargs.paralleljobs) + handleErrors(err) + log.Println("Download successful") + } + } +} From c8601b57e987f8e75563acc18c9803a611ab9238 Mon Sep 17 00:00:00 2001 From: nmahalle Date: Thu, 18 Feb 2021 23:29:10 -0800 Subject: [PATCH 2/8] Added Upload/Download code for S3 --- .../nzconnector/connector/AZConnector.go | 24 +- bnr-utils/nzconnector/connector/Connector.go | 70 +++- .../nzconnector/connector/S3Connector.go | 302 +++++++++++++++++- bnr-utils/nzconnector/factory/Factory.go | 17 + bnr-utils/nzconnector/main.go | 122 ++----- 5 files changed, 409 insertions(+), 126 deletions(-) create mode 100644 bnr-utils/nzconnector/factory/Factory.go diff --git a/bnr-utils/nzconnector/connector/AZConnector.go b/bnr-utils/nzconnector/connector/AZConnector.go index 70fc8ad..2673fb8 100644 --- a/bnr-utils/nzconnector/connector/AZConnector.go +++ b/bnr-utils/nzconnector/connector/AZConnector.go @@ -95,7 +95,6 @@ func (j *uploadJob) upload(conn *AZConnector) error { } func (t *AZConnector) Upload() { - t.azargs() fmt.Println("Uploading with az connector") } @@ -171,10 +170,15 @@ func (cn *AZConnector) uploadFile(absfilepath string, relfilepath string, unique } -func (cn *AZConnector) UploadBkp(bkpdir string, uniqueid string, backupdir string, paralleljobs int) (error){ +func (cn *AZConnector) UploadBkp(bkpdir string, otherargs *OtherArgs, backupinfo *BackupInfo ) (error){ var err error - work := make(chan *uploadJob, paralleljobs) - result := make(chan *jobResult, paralleljobs) + backupdir := filepath.Join(bkpdir, "Netezza", backupinfo.npshost, backupinfo.dbname, backupinfo.backupset) + _, err = os.Stat(backupdir) + if err != nil { + return fmt.Errorf("Error in Creating backupdir : %v", err) + } + work := make(chan *uploadJob, otherargs.paralleljobs) + result := make(chan *jobResult, otherargs.paralleljobs) done := make(chan bool) go func() { @@ -218,7 +222,7 @@ func (cn *AZConnector) UploadBkp(bkpdir string, uniqueid string, backupdir strin if info.IsDir() { return nil } - j := uploadJob{ job: job{uniqueid, bkpdir}, absfilepath: absfilepath } + j := uploadJob{ job: job{otherargs.uniqueid, bkpdir}, absfilepath: absfilepath } work <- &j // this will hang until at least one of the prior uploads finish if other.paralleljobs // are already running return err @@ -263,12 +267,14 @@ func (cn *AZConnector) downloadFile(outfilepath string, blobname string, streams } -func (cn *AZConnector) DownloadBkp(outdir string, uniqueid string, blobpath string, paralleljobs int) (error){ +func (cn *AZConnector) DownloadBkp(outdir string, otherargs *OtherArgs, backupinfo *BackupInfo) (error){ var err error - work := make(chan *downloadJob, paralleljobs) - result := make(chan *downloadJobResult, paralleljobs) + work := make(chan *downloadJob, otherargs.paralleljobs) + result := make(chan *downloadJobResult, otherargs.paralleljobs) done := make(chan bool) + blobpath := filepath.Join(otherargs.uniqueid, "Netezza",backupinfo.npshost, backupinfo.dbname, backupinfo.backupset) + // start the workers go func() { for { @@ -331,7 +337,7 @@ func (cn *AZConnector) DownloadBkp(outdir string, uniqueid string, blobpath stri // Set up file to download the blob to dir, filename := filepath.Split(blobInfo.Name) - relfilepath, err := filepath.Rel(uniqueid,dir) + relfilepath, err := filepath.Rel(otherargs.uniqueid,dir) if err != nil { return fmt.Errorf("Error in fetching download relative path: %v",err) } diff --git a/bnr-utils/nzconnector/connector/Connector.go b/bnr-utils/nzconnector/connector/Connector.go index f2c7623..955a9c1 100644 --- a/bnr-utils/nzconnector/connector/Connector.go +++ b/bnr-utils/nzconnector/connector/Connector.go @@ -1,9 +1,73 @@ package Connector +import ( + "flag" + "log" + "fmt" + "time" + "os" + "path" +) + type IConnector interface { ParseConnectorArgs(string) - UploadBkp(string, string, string, int) (error) - Upload() - DownloadBkp(string, string, string, int) (error) + UploadBkp(string, *OtherArgs, *BackupInfo) (error) + DownloadBkp(string, *OtherArgs, *BackupInfo) (error) +} + +type BackupInfo struct { + dbname string + Dir string + npshost string + backupset string +} + +type ConnectorInfo struct { + Connector string + ConnectorArgs string +} + +type OtherArgs struct { + uniqueid string + logfiledir string + Upload *bool + Download *bool + paralleljobs int +} + +func ParseArgs(backupinfo *BackupInfo, connectorInfo *ConnectorInfo, otherargs *OtherArgs) { + flag.StringVar(&backupinfo.dbname, "db", "", "Database name") + flag.StringVar(&backupinfo.Dir, "dir", "", "Full path to the directory in which the backup already exists or should be downloaded") + flag.StringVar(&backupinfo.npshost, "npshost", "", "Name of the NPS host as it appears in the backups") + flag.StringVar(&backupinfo.backupset, "backupset", "", "Name of the backupset to be uploaded/downloaded") + + flag.StringVar(&connectorInfo.Connector, "connector", "", "Destination cloud store") + flag.StringVar(&connectorInfo.ConnectorArgs, "connectorArgs", "", "Arguments for cloud store") + + flag.StringVar(&otherargs.uniqueid,"uniqueid", "", "Azure blob storage container") + flag.StringVar(&otherargs.logfiledir,"logfiledir", "/tmp", "Logfile directory for this utility. Default is /tmp dir") + otherargs.Upload = flag.Bool("upload", false, "Upload to cloud") + otherargs.Download = flag.Bool("download", false, "Download from cloud") + flag.IntVar(&otherargs.paralleljobs,"paralleljobs",6,"Number of parallel files to upload/download") } +func SetUpLogFile(backupinfo *BackupInfo, connectorInfo *ConnectorInfo, otherargs *OtherArgs) { + // log file configuration setup + logfilename := fmt.Sprintf("nz_%sConnector_%d_%s.log", connectorInfo.Connector, os.Getppid(), time.Now().Format("2006-01-02-150405")) + logfilepath := path.Join(otherargs.logfiledir, logfilename) + filehandle, err := os.OpenFile(logfilepath, os.O_WRONLY|os.O_CREATE|os.O_APPEND, 0644) + if err != nil { + fmt.Errorf("Error in opening logfile: %v",err) + } + log.SetOutput(filehandle) + prefixStr := fmt.Sprintf("%s ", time.Now().UTC().Format("2006-01-02 15:04:05 EST")) + fmt.Sprintf("%-7s", "[INFO]") + log.SetFlags(0) + log.SetPrefix(prefixStr) + + log.Println("Backup/Restore directory :",backupinfo.Dir) + log.Println("DB name :", backupinfo.dbname) + log.Println("Nps hostname :", backupinfo.npshost) + log.Println("BackupsetID :", backupinfo.backupset) + log.Println("UniqueID :", otherargs.uniqueid) + log.Println("Number of files to upload/download in parallel :", otherargs.paralleljobs) +} diff --git a/bnr-utils/nzconnector/connector/S3Connector.go b/bnr-utils/nzconnector/connector/S3Connector.go index 53baeb9..19f375e 100644 --- a/bnr-utils/nzconnector/connector/S3Connector.go +++ b/bnr-utils/nzconnector/connector/S3Connector.go @@ -3,10 +3,24 @@ package Connector import ( "fmt" "strings" + "strconv" + "os" + "path/filepath" + "path" + "log" + "github.com/aws/aws-sdk-go/aws" + "github.com/aws/aws-sdk-go/aws/session" + "github.com/aws/aws-sdk-go/service/s3" + "github.com/aws/aws-sdk-go/service/s3/s3manager" + "github.com/aws/aws-sdk-go/aws/credentials" ) type is3 interface { - s3args() + getUploader()(*s3manager.Uploader, error) + getDownloader()(*s3manager.Downloader, error) + getSession() (sess *session.Session) + uploadFile(absfilepath string, relfilepath string, uniqueid string) (error) + downloadFile(outfilepath string, key string, conn *s3manager.Downloader) (error) } type S3connector struct { @@ -15,26 +29,131 @@ type S3connector struct { default_region string bucket_url string endpoint string - uniqueid string + streams int + blocksize int64 } -func (c *S3connector) Upload() { - c.s3args() - fmt.Println("Uploading with s3 connector") +type jobS3 struct { + uniqueid string + bkpdir string } -func (c *S3connector) UploadBkp(bkpdir string, uniqueid string, backupdir string, paralleljobs int) (error){ - fmt.Println("Uploading") - return nil +type uploadJobS3 struct { + jobS3 + absfilepath string } -func (cn *S3connector) DownloadBkp(outdir string, uniqueid string, blobpath string, paralleljobs int) (error){ - fmt.Println("Downloading") - return nil + +type jobResultS3 struct { + *jobS3 + err error +} + +type downloadJobS3 struct { + conn *s3manager.Downloader + key string + outfilepath string +} + +type downloadJobResultS3 struct { + key string + err error } + + +func (j *uploadJobS3) uploadS3(conn *S3connector) error { + relfilepath, err := filepath.Rel(j.jobS3.bkpdir, j.absfilepath) + if err != nil { + return fmt.Errorf("Unable to traverse %s, %s: %v", j.jobS3.bkpdir, j.absfilepath, err) + } + + log.Println("Uploading file :", j.absfilepath) + return conn.uploadFileS3(j.absfilepath, relfilepath, j.uniqueid) +} + +func (c *S3connector) uploadFileS3(absfilepath string, relfilepath string, uniqueid string) (error){ + // Upload the file to a block blob + uploader := c.getUploader() + + file, err := os.Open(absfilepath) + if err != nil { + return fmt.Errorf("Error in opening backup file : %v", err) + } + + result,err := uploader.Upload(&s3manager.UploadInput{ + Bucket: &c.bucket_url, + Key: aws.String(filepath.Join(uniqueid, relfilepath)), + Body: file, + }) + + if (result == nil) { + return fmt.Errorf("Error in uploading file : %v", err) + } + return err +} + +func (j *downloadJobS3) download(c *S3connector) error { + log.Println("Downloading file :", j.key) + return c.downloadFile(j.outfilepath, j.key, j.conn ) +} + +func (c *S3connector) downloadFile(outfilepath string, key string, conn *s3manager.Downloader) error { + filehandle, err := os.Create(outfilepath) + if err != nil { + return fmt.Errorf("Error in creating file inside backup dir: %v",err) + } + + defer filehandle.Close() + + numBytes, err := conn.Download(filehandle, + &s3.GetObjectInput{ + Bucket: aws.String(c.bucket_url), + Key: aws.String(key), + }) + + if((numBytes < 0) && (err != nil)){ + return fmt.Errorf("Error in downloading File: %v",err) + } + return err +} + +func (c *S3connector) getSession() (sess *session.Session) { + sess, err := session.NewSession(&aws.Config{ + Region: aws.String(c.default_region), + Credentials: credentials.NewStaticCredentials(c.access_key_id,c.secret_access_key,""), + CredentialsChainVerboseErrors: aws.Bool(true) }) + if (err != nil) { + log.Fatalln("Session failed:", err) + } + + return sess +} + +func (c *S3connector) getUploader() (*s3manager.Uploader) { + + sessn := c.getSession() + uploader := s3manager.NewUploader(sessn,func(u *s3manager.Uploader) { + u.PartSize = c.blocksize * 1024 * 1024 // 64MB per part + u.Concurrency = c.streams + }) + + return uploader +} + +func (c *S3connector) getDownloader() (*s3manager.Downloader, *session.Session) { + + sessn := c.getSession() + downloader := s3manager.NewDownloader(sessn,func(d *s3manager.Downloader) { + d.PartSize = c.blocksize * 1024 * 1024 // 64MB per part + d.Concurrency = c.streams + }) + + return downloader,sessn +} + func (c *S3connector) ParseConnectorArgs(args string) { - arguments := strings.Split(args, ":") + arguments := strings.Split(args, ";") for _, arg := range arguments { - kv := strings.Split(arg, "=") + kv := strings.Split(arg, ":") switch kv[0] { case "ACCESS_KEY_ID": c.access_key_id = kv[1] @@ -46,11 +165,164 @@ func (c *S3connector) ParseConnectorArgs(args string) { c.bucket_url = kv[1] case "ENDPOINT": c.endpoint = kv[1] + case "STREAMS": + i, err := strconv.Atoi(kv[1]) + if (err == nil ) { + c.streams = i + } + case "BLOCKSIZE": + u64, err := strconv.ParseInt(kv[1], 10, 64) + if (err == nil ) { + c.blocksize = int64(u64) + } + } } } -func (c S3connector) s3args() { - fmt.Println("ACCESS_KEY_ID : ", c.access_key_id) +func (c *S3connector) UploadBkp(bkpdir string, otherargs *OtherArgs, backupinfo *BackupInfo ) (error){ + var err error + log.Println("Uploading Using S3 Connector") + _, err = os.Stat(bkpdir) + if err != nil { + return fmt.Errorf("Error Directory not present : %v", err) + } + + work := make(chan *uploadJobS3, otherargs.paralleljobs) + result := make(chan *jobResultS3, otherargs.paralleljobs) + done := make(chan bool) + + go func() { + for { + select { + case j, ok := <- work: + if ! ok { + // done + close(result) + return + } + err := j.uploadS3(c) + jr := jobResultS3{jobS3:&j.jobS3, err:err } + result <- &jr + } + } + }() + + filesuploaded := 0 + go func() { + for { + select { + case r, ok := <- result: + if ! ok { + // work done + done <- true + return + } + if r.err != nil { + // stopping right here so that we + // don't keep on uploading when one has failed + log.Fatalf("%s: %v", r.jobS3, r.err) + } + filesuploaded++ // this is fine, since this is single threaded increment + } + } + }() + + log.Println("Uploading Dir :", bkpdir) + err = filepath.Walk(bkpdir, + func(absfilepath string, info os.FileInfo, err error) error { + if info.IsDir() { + return nil + } + j := uploadJobS3{ jobS3: jobS3{otherargs.uniqueid, bkpdir}, absfilepath: absfilepath } + work <- &j // this will hang until at least one of the prior uploads finish if other.paralleljobs + // are already running + return err + }) + + close(work) + <- done + log.Println("Upload using S3 connector successful. Total files uploaded:", filesuploaded) + return err +} + +func (cn *S3connector) DownloadBkp(outdir string, otherargs *OtherArgs, backupinfo *BackupInfo) (error){ + log.Println("Downloading Using S3 Connector") + + var err error + work := make(chan *downloadJobS3, otherargs.paralleljobs) + result := make(chan *downloadJobResultS3, otherargs.paralleljobs) + done := make(chan bool) + + + // start the workers + go func() { + for { + select { + case j, ok := <- work: + if ! ok { + // done + close(result) + return + } + err := j.download(cn) + jr := downloadJobResultS3{ key:j.key, err:err } + result <- &jr + } + } + }() + filesdownloaded := 0 + go func() { + for { + select { + case r, ok := <- result: + if ! ok { + // work done + done <- true + return + } + if r.err != nil { + // stopping right here so that we + // don't keep on uploading when one has failed + log.Fatalf("%s: %v", r.key, r.err) + } + filesdownloaded++ // this is fine, since this is single threaded increment + } + } + }() + + down,sess := cn.getDownloader() + client := s3.New(sess) + params := &s3.ListObjectsInput{Bucket: &cn.bucket_url, Prefix: &otherargs.uniqueid} + client.ListObjectsPages(params, func(page *s3.ListObjectsOutput, more bool) (bool) { + for _, obj := range page.Contents { + key := *obj.Key + // Create the directories in the path + + dir, filename := filepath.Split(key) + relfilepath, err := filepath.Rel(otherargs.uniqueid,dir) + + if err != nil { + log.Fatalf("Error in fetching download relative path: %v",err) + } + + file := filepath.Join(outdir, relfilepath) + err = os.MkdirAll(file, 0777) + if err != nil { + log.Fatalf("Error in creating backup directory structure: %v",err) + } + + outfilepath := path.Join(file, filename) + j := downloadJobS3{ conn:down, key:key, outfilepath:outfilepath } + work <- &j + + } + return true + }) + + close(work) + <- done + log.Println("Total files downloaded using S3 Connector:", filesdownloaded) + return err } diff --git a/bnr-utils/nzconnector/factory/Factory.go b/bnr-utils/nzconnector/factory/Factory.go new file mode 100644 index 0000000..af0582b --- /dev/null +++ b/bnr-utils/nzconnector/factory/Factory.go @@ -0,0 +1,17 @@ +package Factory + +import ( + "nzconnector/connector" +) + +func GetConnector(connectorType string) Connector.IConnector { + switch connectorType { + case "s3": + return &Connector.S3connector{} + case "az": + return &Connector.AZConnector{} + default: + return nil + } +} + diff --git a/bnr-utils/nzconnector/main.go b/bnr-utils/nzconnector/main.go index a58b4e0..15fcdf8 100644 --- a/bnr-utils/nzconnector/main.go +++ b/bnr-utils/nzconnector/main.go @@ -3,123 +3,47 @@ package main import ( "flag" "nzconnector/connector" - "fmt" - "path/filepath" + "nzconnector/factory" "log" - "os" - "time" - "path" "strings" ) -type BackupInfo struct { - dbname string - dir string - npshost string - backupset string -} - -type ConnectorInfo struct { - connector string - connectorArgs string -} - -type OtherArgs struct { - uniqueid string - logfiledir string - upload *bool - download *bool - paralleljobs int -} - - -func parseArgs(backupinfo *BackupInfo, connectorInfo *ConnectorInfo, otherargs *OtherArgs) { - flag.StringVar(&backupinfo.dbname, "db", "", "Database name") - flag.StringVar(&backupinfo.dir, "dir", "", "Full path to the directory in which the backup already exists or should be downloaded") - flag.StringVar(&backupinfo.npshost, "npshost", "", "Name of the NPS host as it appears in the backups") - flag.StringVar(&backupinfo.backupset, "backupset", "", "Name of the backupset to be uploaded/downloaded") - - flag.StringVar(&connectorInfo.connector, "connector", "", "Destination cloud store") - flag.StringVar(&connectorInfo.connectorArgs, "connectorArgs", "", "Arguments for cloud store") - - flag.StringVar(&otherargs.uniqueid,"uniqueid", "", "Azure blob storage container") - flag.StringVar(&otherargs.logfiledir,"logfiledir", "/tmp", "Logfile directory for this utility. Default is /tmp dir") - otherargs.upload = flag.Bool("upload", false, "Upload to cloud") - otherargs.download = flag.Bool("download", false, "Download from cloud") - flag.IntVar(&otherargs.paralleljobs,"paralleljobs",6,"Number of parallel files to upload/download") -} - -func handleErrors(err error) { - if err != nil { - log.Fatalln(err) - } -} - func parseConnectorArgs(e Connector.IConnector, args string) { if e != nil { e.ParseConnectorArgs(args) } } -func GetConnector(connectorType string) Connector.IConnector { - switch connectorType { - case "s3": - return &Connector.S3connector{} - case "az": - return &Connector.AZConnector{} - default: - return nil - } -} func main() { - var backupinfo BackupInfo - var connectorInfo ConnectorInfo - var otherargs OtherArgs - // parse input args - parseArgs(&backupinfo, &connectorInfo, &otherargs) + var backupinfo Connector.BackupInfo + var connectorInfo Connector.ConnectorInfo + var otherargs Connector.OtherArgs + var err error +// parse input args + Connector.ParseArgs(&backupinfo, &connectorInfo, &otherargs) flag.Parse() + Connector.SetUpLogFile(&backupinfo, &connectorInfo, &otherargs) + connector := Factory.GetConnector(connectorInfo.Connector) + parseConnectorArgs(connector, connectorInfo.ConnectorArgs) - connector := GetConnector(connectorInfo.connector) - parseConnectorArgs(connector, connectorInfo.connectorArgs) - - // log file configuration setup - logfilename := fmt.Sprintf("nz_azConnector_%d_%s.log", os.Getppid(), time.Now().Format("2006-01-02")) - logfilepath := path.Join(otherargs.logfiledir, logfilename) - filehandle, err := os.OpenFile(logfilepath, os.O_WRONLY|os.O_CREATE|os.O_APPEND, 0644) - if err != nil { - fmt.Errorf("Error in opening logfile: %v",err) - } - log.SetOutput(filehandle) - prefixStr := fmt.Sprintf("%s ", time.Now().UTC().Format("2006-01-02 15:04:05 EST")) + fmt.Sprintf("%-7s", "[INFO]") - log.SetFlags(0) - log.SetPrefix(prefixStr) - - dirlist := strings.Split(backupinfo.dir," ") - log.Println("Backup/Restore directory :",dirlist) - log.Println("DB name :", backupinfo.dbname) - log.Println("Nps hostname :", backupinfo.npshost) - log.Println("BackupsetID :", backupinfo.backupset) - log.Println("UniqueID :", otherargs.uniqueid) - log.Println("Number of files to upload/download in parallel :", otherargs.paralleljobs) - + dirlist := strings.Split(backupinfo.Dir," ") for _, bkpdir := range dirlist { - if (*otherargs.upload) { - + if (*otherargs.Upload) { // now do the upload - log.Println("Uploading backup data to azure cloud from backup dir", bkpdir) - backupdir := filepath.Join(bkpdir, "Netezza", backupinfo.npshost, backupinfo.dbname, backupinfo.backupset) - _, err = os.Stat(backupdir) - handleErrors(err) - err = connector.UploadBkp(bkpdir, otherargs.uniqueid, backupdir, otherargs.paralleljobs) - handleErrors(err) + log.Println("Uploading backup data to cloud from backup dir", bkpdir) + err = connector.UploadBkp(bkpdir, &otherargs, &backupinfo) + if (err != nil) { + log.Fatalln(err) + } log.Println("Upload successful") } - if (*otherargs.download) { - log.Println("Downloading backup data from azure cloud to restore dir", bkpdir) - blobpath := filepath.Join(otherargs.uniqueid, "Netezza",backupinfo.npshost, backupinfo.dbname, backupinfo.backupset) - err = connector.DownloadBkp(bkpdir, otherargs.uniqueid, blobpath, otherargs.paralleljobs) - handleErrors(err) + if (*otherargs.Download) { + log.Println("Downloading backup data from cloud to restore dir", bkpdir) + err = connector.DownloadBkp(bkpdir, &otherargs, &backupinfo) + if (err != nil) { + log.Fatalln(err) + } log.Println("Download successful") } } From dac1547908e193c95df6283ea96afd2ffbf1d7bb Mon Sep 17 00:00:00 2001 From: nmahalle Date: Fri, 19 Feb 2021 02:48:26 -0800 Subject: [PATCH 3/8] Added code to dowload database directory --- .../nzconnector/connector/S3Connector.go | 32 ++++++++++--------- 1 file changed, 17 insertions(+), 15 deletions(-) diff --git a/bnr-utils/nzconnector/connector/S3Connector.go b/bnr-utils/nzconnector/connector/S3Connector.go index 19f375e..ed3fd8a 100644 --- a/bnr-utils/nzconnector/connector/S3Connector.go +++ b/bnr-utils/nzconnector/connector/S3Connector.go @@ -254,6 +254,7 @@ func (cn *S3connector) DownloadBkp(outdir string, otherargs *OtherArgs, backupin result := make(chan *downloadJobResultS3, otherargs.paralleljobs) done := make(chan bool) + bkpath := filepath.Join(otherargs.uniqueid, "Netezza",backupinfo.npshost, backupinfo.dbname, backupinfo.backupset) // start the workers go func() { @@ -297,25 +298,26 @@ func (cn *S3connector) DownloadBkp(outdir string, otherargs *OtherArgs, backupin client.ListObjectsPages(params, func(page *s3.ListObjectsOutput, more bool) (bool) { for _, obj := range page.Contents { key := *obj.Key - // Create the directories in the path + if strings.HasPrefix(key, bkpath){ + // Create the directories in the path - dir, filename := filepath.Split(key) - relfilepath, err := filepath.Rel(otherargs.uniqueid,dir) + dir, filename := filepath.Split(key) + relfilepath, err := filepath.Rel(otherargs.uniqueid,dir) - if err != nil { - log.Fatalf("Error in fetching download relative path: %v",err) - } - - file := filepath.Join(outdir, relfilepath) - err = os.MkdirAll(file, 0777) - if err != nil { - log.Fatalf("Error in creating backup directory structure: %v",err) - } + if err != nil { + log.Fatalf("Error in fetching download relative path: %v",err) + } - outfilepath := path.Join(file, filename) - j := downloadJobS3{ conn:down, key:key, outfilepath:outfilepath } - work <- &j + file := filepath.Join(outdir, relfilepath) + err = os.MkdirAll(file, 0777) + if err != nil { + log.Fatalf("Error in creating backup directory structure: %v",err) + } + outfilepath := path.Join(file, filename) + j := downloadJobS3{ conn:down, key:key, outfilepath:outfilepath } + work <- &j + } } return true }) From c5ffc54803a18876e1eba9b7a0945542e60f701b Mon Sep 17 00:00:00 2001 From: nmahalle Date: Thu, 25 Feb 2021 06:31:47 -0800 Subject: [PATCH 4/8] changes to support cloud backup download --- .../nzconnector/connector/AZConnector.go | 148 +++++++++--------- bnr-utils/nzconnector/connector/Connector.go | 67 +++++++- .../nzconnector/connector/S3Connector.go | 138 ++++++++-------- bnr-utils/nzconnector/main.go | 48 +++--- 4 files changed, 242 insertions(+), 159 deletions(-) diff --git a/bnr-utils/nzconnector/connector/AZConnector.go b/bnr-utils/nzconnector/connector/AZConnector.go index 2673fb8..f601791 100644 --- a/bnr-utils/nzconnector/connector/AZConnector.go +++ b/bnr-utils/nzconnector/connector/AZConnector.go @@ -15,7 +15,6 @@ import ( ) type iaz interface { - azargs() getServiceURL() (azblob.ServiceURL, error) getContainerURL() (azblob.ContainerURL, error) getBlockBlobURL(blobname string) (azblob.BlockBlobURL, error) @@ -84,7 +83,7 @@ func (c *AZConnector) ParseConnectorArgs(args string) { } } -func (j *uploadJob) upload(conn *AZConnector) error { +func (j *uploadJob) uploadAZ(conn *AZConnector) error { relfilepath, err := filepath.Rel(j.job.bkpdir, j.absfilepath) if err != nil { return fmt.Errorf("Unable to traverse %s, %s: %v", j.job.bkpdir, j.absfilepath, err) @@ -94,10 +93,6 @@ func (j *uploadJob) upload(conn *AZConnector) error { return conn.uploadFile(j.absfilepath, relfilepath, j.job.uniqueid) } -func (t *AZConnector) Upload() { - fmt.Println("Uploading with az connector") -} - func (cn *AZConnector) getServiceURL() (azblob.ServiceURL, error) { var serviceURL azblob.ServiceURL us := fmt.Sprintf("https://%s.blob.core.windows.net/", cn.azaccount) @@ -170,70 +165,75 @@ func (cn *AZConnector) uploadFile(absfilepath string, relfilepath string, unique } -func (cn *AZConnector) UploadBkp(bkpdir string, otherargs *OtherArgs, backupinfo *BackupInfo ) (error){ +func (cn *AZConnector) Upload( otherargs *OtherArgs, backupinfo *BackupInfo ) (error){ var err error - backupdir := filepath.Join(bkpdir, "Netezza", backupinfo.npshost, backupinfo.dbname, backupinfo.backupset) - _, err = os.Stat(backupdir) - if err != nil { - return fmt.Errorf("Error in Creating backupdir : %v", err) - } - work := make(chan *uploadJob, otherargs.paralleljobs) - result := make(chan *jobResult, otherargs.paralleljobs) - done := make(chan bool) - - go func() { - for { - select { - case j, ok := <- work: - if ! ok { - // done - close(result) - return - } - err := j.upload(cn) - jr := jobResult{ job:&j.job, err:err } - result <- &jr - } + log.Println("Uploading Using AZ Connector") + dirlist := strings.Split(backupinfo.Dir," ") + for _, bkpdir := range dirlist { + backupdir := filepath.Join(bkpdir, "Netezza", backupinfo.npshost, backupinfo.dbname, backupinfo.backupset) + _, err = os.Stat(backupdir) + if err != nil { + return fmt.Errorf("Error in Creating backupdir : %v", err) } - }() - - filesuploaded := 0 - go func() { - for { - select { - case r, ok := <- result: - if ! ok { - // work done - done <- true - return + work := make(chan *uploadJob, otherargs.paralleljobs) + result := make(chan *jobResult, otherargs.paralleljobs) + done := make(chan bool) + + go func() { + for { + select { + case j, ok := <- work: + if ! ok { + // done + close(result) + return + } + err := j.uploadAZ(cn) + jr := jobResult{ job:&j.job, err:err } + result <- &jr } - if r.err != nil { - // stopping right here so that we - // don't keep on uploading when one has failed - log.Fatalf("%s: %v", r.job, r.err) + } + }() + + filesuploaded := 0 + go func() { + for { + select { + case r, ok := <- result: + if ! ok { + // work done + done <- true + return + } + if r.err != nil { + // stopping right here so that we + // don't keep on uploading when one has failed + log.Fatalf("%s: %v", r.job, r.err) + } + filesuploaded++ // this is fine, since this is single threaded increment } - filesuploaded++ // this is fine, since this is single threaded increment } - } - }() + }() - err = filepath.Walk(backupdir, - func(absfilepath string, info os.FileInfo, err error) error { - if info.IsDir() { - return nil - } - j := uploadJob{ job: job{otherargs.uniqueid, bkpdir}, absfilepath: absfilepath } - work <- &j // this will hang until at least one of the prior uploads finish if other.paralleljobs - // are already running - return err - }) + err = filepath.Walk(backupdir, + func(absfilepath string, info os.FileInfo, err error) error { + if info.IsDir() { + return nil + } + j := uploadJob{ job: job{otherargs.uniqueid, bkpdir}, absfilepath: absfilepath } + work <- &j // this will hang until at least one of the prior uploads finish if other.paralleljobs + // are already running + return err + }) close(work) <- done - log.Println("Upload successful. Total files uploaded:", filesuploaded) - return err + log.Println("Upload successful for Backup Dir :", bkpdir) + log.Println("Total files uploaded:", filesuploaded) } + return err +} -func (j *downloadJob) download(conn *AZConnector) error { +func (j *downloadJob) downloadAZ(conn *AZConnector) error { log.Println("Downloading file :", j.blobname) return conn.downloadFile(j.outfilepath, j.blobname, conn.streams, conn.blocksize) @@ -267,8 +267,12 @@ func (cn *AZConnector) downloadFile(outfilepath string, blobname string, streams } -func (cn *AZConnector) DownloadBkp(outdir string, otherargs *OtherArgs, backupinfo *BackupInfo) (error){ +func (cn *AZConnector) Download(otherargs *OtherArgs, backupinfo *BackupInfo) (error){ var err error + log.Println("Downloading Using AZ Connector") + outdir := backupinfo.Dir + arrayLoc:= []string{} + arrayContents:= []string{} work := make(chan *downloadJob, otherargs.paralleljobs) result := make(chan *downloadJobResult, otherargs.paralleljobs) done := make(chan bool) @@ -285,7 +289,7 @@ func (cn *AZConnector) DownloadBkp(outdir string, otherargs *OtherArgs, backupin close(result) return } - err := j.download(cn) + err := j.downloadAZ(cn) jr := downloadJobResult{ blobname:j.blobname, err:err } result <- &jr } @@ -349,6 +353,13 @@ func (cn *AZConnector) DownloadBkp(outdir string, otherargs *OtherArgs, backupin } outfilepath := path.Join(dumpdir, filename) + if (strings.HasSuffix(outfilepath,"locations.txt")){ + arrayLoc = append(arrayLoc,outfilepath) + } + if (strings.HasSuffix(outfilepath,"contents.txt")){ + arrayContents = append(arrayContents,outfilepath) + } + j := downloadJob{ conn:*cn, outfilepath:outfilepath, blobname: blobInfo.Name } work <- &j @@ -366,14 +377,11 @@ func (cn *AZConnector) DownloadBkp(outdir string, otherargs *OtherArgs, backupin } close(work) <- done + log.Println("File Downloaded to dir :", outdir) log.Println("Total files downloaded:", filesdownloaded) + if (*otherargs.cloudBackup) { + updateLocation(arrayLoc,outdir) + updateContents(arrayContents) + } return err } - - -func (t AZConnector) azargs() { - fmt.Println("STORAGE_ACCOUNT : ", t.azaccount) - fmt.Println("STORAGE_ACCOUNT : ", t.blocksize) - fmt.Println("STORAGE_ACCOUNT : ", t.streams) -} - diff --git a/bnr-utils/nzconnector/connector/Connector.go b/bnr-utils/nzconnector/connector/Connector.go index 955a9c1..c06213e 100644 --- a/bnr-utils/nzconnector/connector/Connector.go +++ b/bnr-utils/nzconnector/connector/Connector.go @@ -5,14 +5,16 @@ import ( "log" "fmt" "time" + "io/ioutil" "os" "path" + "strings" ) type IConnector interface { ParseConnectorArgs(string) - UploadBkp(string, *OtherArgs, *BackupInfo) (error) - DownloadBkp(string, *OtherArgs, *BackupInfo) (error) + Upload(*OtherArgs, *BackupInfo) (error) + Download( *OtherArgs, *BackupInfo) (error) } type BackupInfo struct { @@ -30,9 +32,11 @@ type ConnectorInfo struct { type OtherArgs struct { uniqueid string logfiledir string - Upload *bool - Download *bool + upload *bool + download *bool paralleljobs int + Operation int + cloudBackup *bool } func ParseArgs(backupinfo *BackupInfo, connectorInfo *ConnectorInfo, otherargs *OtherArgs) { @@ -46,8 +50,9 @@ func ParseArgs(backupinfo *BackupInfo, connectorInfo *ConnectorInfo, otherargs * flag.StringVar(&otherargs.uniqueid,"uniqueid", "", "Azure blob storage container") flag.StringVar(&otherargs.logfiledir,"logfiledir", "/tmp", "Logfile directory for this utility. Default is /tmp dir") - otherargs.Upload = flag.Bool("upload", false, "Upload to cloud") - otherargs.Download = flag.Bool("download", false, "Download from cloud") + otherargs.upload = flag.Bool("upload", false, "Upload to cloud") + otherargs.download = flag.Bool("download", false, "Download from cloud") + otherargs.cloudBackup = flag.Bool("cloudBackup", false, "Download backup taken on cloud") flag.IntVar(&otherargs.paralleljobs,"paralleljobs",6,"Number of parallel files to upload/download") } @@ -70,4 +75,52 @@ func SetUpLogFile(backupinfo *BackupInfo, connectorInfo *ConnectorInfo, otherarg log.Println("BackupsetID :", backupinfo.backupset) log.Println("UniqueID :", otherargs.uniqueid) log.Println("Number of files to upload/download in parallel :", otherargs.paralleljobs) -} +} + +func SetOperation(otherargs *OtherArgs) { + if (*otherargs.upload){ + otherargs.Operation = 0 + } + if (*otherargs.download){ + otherargs.Operation = 1 + } +} + +func updateLocation(arrLoc []string,outdir string){ + for _,locFile := range arrLoc{ + f, err := os.OpenFile(locFile, + os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644) + if err != nil { + log.Fatalln(err) + } + defer f.Close() + textAppend := "1,1,1," + outdir + if _, err := f.WriteString(textAppend); err != nil { + log.Fatalln(err) + } + } +} + +func updateContents(arrContents []string){ + for _,contentFile := range arrContents{ + input, err := ioutil.ReadFile(contentFile) + if err != nil { + log.Fatalln(err) + } + + lines := strings.Split(string(input), "\n") + lines = lines[:len(lines)-1] + var textline []string + for _, line := range lines { + r := []rune(line) + str := string(r[:len(r)-1]) + "1" + textline = append(textline,str) + } + output := strings.Join(textline, "\n") + err = ioutil.WriteFile(contentFile, []byte(output), 0644) + if err != nil { + log.Fatalln(err) + } + } +} + diff --git a/bnr-utils/nzconnector/connector/S3Connector.go b/bnr-utils/nzconnector/connector/S3Connector.go index ed3fd8a..ff6246d 100644 --- a/bnr-utils/nzconnector/connector/S3Connector.go +++ b/bnr-utils/nzconnector/connector/S3Connector.go @@ -91,7 +91,7 @@ func (c *S3connector) uploadFileS3(absfilepath string, relfilepath string, uniqu return err } -func (j *downloadJobS3) download(c *S3connector) error { +func (j *downloadJobS3) downloadS3(c *S3connector) error { log.Println("Downloading file :", j.key) return c.downloadFile(j.outfilepath, j.key, j.conn ) } @@ -113,6 +113,7 @@ func (c *S3connector) downloadFile(outfilepath string, key string, conn *s3manag if((numBytes < 0) && (err != nil)){ return fmt.Errorf("Error in downloading File: %v",err) } + return err } @@ -180,82 +181,88 @@ func (c *S3connector) ParseConnectorArgs(args string) { } } -func (c *S3connector) UploadBkp(bkpdir string, otherargs *OtherArgs, backupinfo *BackupInfo ) (error){ +func (c *S3connector) Upload( otherargs *OtherArgs, backupinfo *BackupInfo ) (error){ var err error log.Println("Uploading Using S3 Connector") - _, err = os.Stat(bkpdir) - if err != nil { - return fmt.Errorf("Error Directory not present : %v", err) - } - - work := make(chan *uploadJobS3, otherargs.paralleljobs) - result := make(chan *jobResultS3, otherargs.paralleljobs) - done := make(chan bool) - go func() { - for { - select { - case j, ok := <- work: - if ! ok { - // done - close(result) - return - } - err := j.uploadS3(c) - jr := jobResultS3{jobS3:&j.jobS3, err:err } - result <- &jr - } + dirlist := strings.Split(backupinfo.Dir," ") + for _, bkpdir := range dirlist { + backupdir := filepath.Join(bkpdir, "Netezza", backupinfo.npshost, backupinfo.dbname, backupinfo.backupset) + _, err = os.Stat(backupdir) + if err != nil { + return fmt.Errorf("Error Directory not present : %v", err) } - }() - filesuploaded := 0 - go func() { - for { - select { - case r, ok := <- result: + work := make(chan *uploadJobS3, otherargs.paralleljobs) + result := make(chan *jobResultS3, otherargs.paralleljobs) + done := make(chan bool) + + go func() { + for { + select { + case j, ok := <- work: if ! ok { - // work done - done <- true + // done + close(result) return + } + err := j.uploadS3(c) + jr := jobResultS3{jobS3:&j.jobS3, err:err } + result <- &jr } - if r.err != nil { - // stopping right here so that we - // don't keep on uploading when one has failed - log.Fatalf("%s: %v", r.jobS3, r.err) + } + }() + + filesuploaded := 0 + go func() { + for { + select { + case r, ok := <- result: + if ! ok { + // work done + done <- true + return + } + if r.err != nil { + // stopping right here so that we + // don't keep on uploading when one has failed + log.Fatalf("%s: %v", r.jobS3, r.err) + } + filesuploaded++ // this is fine, since this is single threaded increment } - filesuploaded++ // this is fine, since this is single threaded increment } - } - }() + }() - log.Println("Uploading Dir :", bkpdir) - err = filepath.Walk(bkpdir, - func(absfilepath string, info os.FileInfo, err error) error { - if info.IsDir() { - return nil - } - j := uploadJobS3{ jobS3: jobS3{otherargs.uniqueid, bkpdir}, absfilepath: absfilepath } - work <- &j // this will hang until at least one of the prior uploads finish if other.paralleljobs - // are already running - return err - }) - - close(work) - <- done - log.Println("Upload using S3 connector successful. Total files uploaded:", filesuploaded) + err = filepath.Walk(bkpdir, + func(absfilepath string, info os.FileInfo, err error) error { + if info.IsDir() { + return nil + } + j := uploadJobS3{ jobS3: jobS3{otherargs.uniqueid, bkpdir}, absfilepath: absfilepath } + work <- &j // this will hang until at least one of the prior uploads finish if other.paralleljobs + // are already running + return err + }) + + close(work) + <- done + log.Println("Upload using S3 connector successful for directory :", bkpdir) + log.Println("Total files uploaded:", filesuploaded) + } return err } -func (cn *S3connector) DownloadBkp(outdir string, otherargs *OtherArgs, backupinfo *BackupInfo) (error){ +func (cn *S3connector) Download(otherargs *OtherArgs, backupinfo *BackupInfo) (error){ log.Println("Downloading Using S3 Connector") - var err error + outdir := backupinfo.Dir + arrayLoc:= []string{} + arrayContents:= []string{} work := make(chan *downloadJobS3, otherargs.paralleljobs) result := make(chan *downloadJobResultS3, otherargs.paralleljobs) done := make(chan bool) bkpath := filepath.Join(otherargs.uniqueid, "Netezza",backupinfo.npshost, backupinfo.dbname, backupinfo.backupset) - // start the workers go func() { for { @@ -266,7 +273,7 @@ func (cn *S3connector) DownloadBkp(outdir string, otherargs *OtherArgs, backupin close(result) return } - err := j.download(cn) + err := j.downloadS3(cn) jr := downloadJobResultS3{ key:j.key, err:err } result <- &jr } @@ -315,6 +322,12 @@ func (cn *S3connector) DownloadBkp(outdir string, otherargs *OtherArgs, backupin } outfilepath := path.Join(file, filename) + if (strings.HasSuffix(outfilepath,"locations.txt")){ + arrayLoc = append(arrayLoc,outfilepath) + } + if (strings.HasSuffix(outfilepath,"contents.txt")){ + arrayContents = append(arrayContents,outfilepath) + } j := downloadJobS3{ conn:down, key:key, outfilepath:outfilepath } work <- &j } @@ -322,9 +335,12 @@ func (cn *S3connector) DownloadBkp(outdir string, otherargs *OtherArgs, backupin return true }) - close(work) - <- done - log.Println("Total files downloaded using S3 Connector:", filesdownloaded) + close(work) + <- done + log.Println("Total files downloaded using S3 Connector:", filesdownloaded) + if (*otherargs.cloudBackup) { + updateLocation(arrayLoc,outdir) + updateContents(arrayContents) + } return err } - diff --git a/bnr-utils/nzconnector/main.go b/bnr-utils/nzconnector/main.go index 15fcdf8..65dbd48 100644 --- a/bnr-utils/nzconnector/main.go +++ b/bnr-utils/nzconnector/main.go @@ -2,10 +2,14 @@ package main import ( "flag" + "log" "nzconnector/connector" "nzconnector/factory" - "log" - "strings" +) + +const( + upload = iota + download ) func parseConnectorArgs(e Connector.IConnector, args string) { @@ -27,24 +31,26 @@ func main() { connector := Factory.GetConnector(connectorInfo.Connector) parseConnectorArgs(connector, connectorInfo.ConnectorArgs) - dirlist := strings.Split(backupinfo.Dir," ") - for _, bkpdir := range dirlist { - if (*otherargs.Upload) { - // now do the upload - log.Println("Uploading backup data to cloud from backup dir", bkpdir) - err = connector.UploadBkp(bkpdir, &otherargs, &backupinfo) - if (err != nil) { - log.Fatalln(err) - } - log.Println("Upload successful") - } - if (*otherargs.Download) { - log.Println("Downloading backup data from cloud to restore dir", bkpdir) - err = connector.DownloadBkp(bkpdir, &otherargs, &backupinfo) - if (err != nil) { - log.Fatalln(err) - } - log.Println("Download successful") + Connector.SetOperation(&otherargs) + + switch otherargs.Operation { + case upload: + // now do the upload + log.Println("Uploading backup data to cloud for conector : ",connectorInfo.Connector) + err = connector.Upload(&otherargs, &backupinfo) + if (err != nil) { + log.Fatalln(err) + } + log.Println("Upload successful") + case download: + log.Println("Downloading backup data from cloud for connector : ", connectorInfo.Connector) + err = connector.Download(&otherargs, &backupinfo) + if (err != nil) { + log.Fatalln(err) + } + log.Println("Download successful") + default: + log.Fatalln("Invalid Operation, Supported Upload/Download") } - } + } From 90a6421137ecbb7330b1af9137ec8574f4038376 Mon Sep 17 00:00:00 2001 From: nmahalle Date: Tue, 9 Mar 2021 03:42:20 -0800 Subject: [PATCH 5/8] changed utility named to nz_cp_backup --- bnr-utils/{nzconnector => nz_cp_backup}/connector/AZConnector.go | 0 bnr-utils/{nzconnector => nz_cp_backup}/connector/Connector.go | 0 bnr-utils/{nzconnector => nz_cp_backup}/connector/S3Connector.go | 0 bnr-utils/{nzconnector => nz_cp_backup}/factory/Factory.go | 0 bnr-utils/{nzconnector => nz_cp_backup}/main.go | 0 5 files changed, 0 insertions(+), 0 deletions(-) rename bnr-utils/{nzconnector => nz_cp_backup}/connector/AZConnector.go (100%) rename bnr-utils/{nzconnector => nz_cp_backup}/connector/Connector.go (100%) rename bnr-utils/{nzconnector => nz_cp_backup}/connector/S3Connector.go (100%) rename bnr-utils/{nzconnector => nz_cp_backup}/factory/Factory.go (100%) rename bnr-utils/{nzconnector => nz_cp_backup}/main.go (100%) diff --git a/bnr-utils/nzconnector/connector/AZConnector.go b/bnr-utils/nz_cp_backup/connector/AZConnector.go similarity index 100% rename from bnr-utils/nzconnector/connector/AZConnector.go rename to bnr-utils/nz_cp_backup/connector/AZConnector.go diff --git a/bnr-utils/nzconnector/connector/Connector.go b/bnr-utils/nz_cp_backup/connector/Connector.go similarity index 100% rename from bnr-utils/nzconnector/connector/Connector.go rename to bnr-utils/nz_cp_backup/connector/Connector.go diff --git a/bnr-utils/nzconnector/connector/S3Connector.go b/bnr-utils/nz_cp_backup/connector/S3Connector.go similarity index 100% rename from bnr-utils/nzconnector/connector/S3Connector.go rename to bnr-utils/nz_cp_backup/connector/S3Connector.go diff --git a/bnr-utils/nzconnector/factory/Factory.go b/bnr-utils/nz_cp_backup/factory/Factory.go similarity index 100% rename from bnr-utils/nzconnector/factory/Factory.go rename to bnr-utils/nz_cp_backup/factory/Factory.go diff --git a/bnr-utils/nzconnector/main.go b/bnr-utils/nz_cp_backup/main.go similarity index 100% rename from bnr-utils/nzconnector/main.go rename to bnr-utils/nz_cp_backup/main.go From 9653db01f620356161488db84c73a7271805e5e7 Mon Sep 17 00:00:00 2001 From: nmahalle Date: Thu, 11 Mar 2021 07:13:06 -0800 Subject: [PATCH 6/8] Added -increment to add level of granularity of upload/download --- bnr-utils/nz_cp_backup/connector/AZConnector.go | 4 ++-- bnr-utils/nz_cp_backup/connector/Connector.go | 2 ++ bnr-utils/nz_cp_backup/connector/S3Connector.go | 12 +++++++----- 3 files changed, 11 insertions(+), 7 deletions(-) diff --git a/bnr-utils/nz_cp_backup/connector/AZConnector.go b/bnr-utils/nz_cp_backup/connector/AZConnector.go index f601791..a0855d6 100644 --- a/bnr-utils/nz_cp_backup/connector/AZConnector.go +++ b/bnr-utils/nz_cp_backup/connector/AZConnector.go @@ -170,7 +170,7 @@ func (cn *AZConnector) Upload( otherargs *OtherArgs, backupinfo *BackupInfo ) (e log.Println("Uploading Using AZ Connector") dirlist := strings.Split(backupinfo.Dir," ") for _, bkpdir := range dirlist { - backupdir := filepath.Join(bkpdir, "Netezza", backupinfo.npshost, backupinfo.dbname, backupinfo.backupset) + backupdir := filepath.Join(bkpdir, "Netezza", backupinfo.npshost, backupinfo.dbname, backupinfo.backupset, backupinfo.increment) _, err = os.Stat(backupdir) if err != nil { return fmt.Errorf("Error in Creating backupdir : %v", err) @@ -277,7 +277,7 @@ func (cn *AZConnector) Download(otherargs *OtherArgs, backupinfo *BackupInfo) (e result := make(chan *downloadJobResult, otherargs.paralleljobs) done := make(chan bool) - blobpath := filepath.Join(otherargs.uniqueid, "Netezza",backupinfo.npshost, backupinfo.dbname, backupinfo.backupset) + blobpath := filepath.Join(otherargs.uniqueid, "Netezza",backupinfo.npshost, backupinfo.dbname, backupinfo.backupset, backupinfo.increment) // start the workers go func() { diff --git a/bnr-utils/nz_cp_backup/connector/Connector.go b/bnr-utils/nz_cp_backup/connector/Connector.go index c06213e..53f3738 100644 --- a/bnr-utils/nz_cp_backup/connector/Connector.go +++ b/bnr-utils/nz_cp_backup/connector/Connector.go @@ -22,6 +22,7 @@ type BackupInfo struct { Dir string npshost string backupset string + increment string } type ConnectorInfo struct { @@ -44,6 +45,7 @@ func ParseArgs(backupinfo *BackupInfo, connectorInfo *ConnectorInfo, otherargs * flag.StringVar(&backupinfo.Dir, "dir", "", "Full path to the directory in which the backup already exists or should be downloaded") flag.StringVar(&backupinfo.npshost, "npshost", "", "Name of the NPS host as it appears in the backups") flag.StringVar(&backupinfo.backupset, "backupset", "", "Name of the backupset to be uploaded/downloaded") + flag.StringVar(&backupinfo.increment, "increment", "", "Increment Number to be uploaded/downloaded") flag.StringVar(&connectorInfo.Connector, "connector", "", "Destination cloud store") flag.StringVar(&connectorInfo.ConnectorArgs, "connectorArgs", "", "Arguments for cloud store") diff --git a/bnr-utils/nz_cp_backup/connector/S3Connector.go b/bnr-utils/nz_cp_backup/connector/S3Connector.go index ff6246d..8ed831a 100644 --- a/bnr-utils/nz_cp_backup/connector/S3Connector.go +++ b/bnr-utils/nz_cp_backup/connector/S3Connector.go @@ -187,7 +187,7 @@ func (c *S3connector) Upload( otherargs *OtherArgs, backupinfo *BackupInfo ) (er dirlist := strings.Split(backupinfo.Dir," ") for _, bkpdir := range dirlist { - backupdir := filepath.Join(bkpdir, "Netezza", backupinfo.npshost, backupinfo.dbname, backupinfo.backupset) + backupdir := filepath.Join(bkpdir, "Netezza", backupinfo.npshost, backupinfo.dbname, backupinfo.backupset, backupinfo.increment) _, err = os.Stat(backupdir) if err != nil { return fmt.Errorf("Error Directory not present : %v", err) @@ -238,9 +238,11 @@ func (c *S3connector) Upload( otherargs *OtherArgs, backupinfo *BackupInfo ) (er if info.IsDir() { return nil } - j := uploadJobS3{ jobS3: jobS3{otherargs.uniqueid, bkpdir}, absfilepath: absfilepath } - work <- &j // this will hang until at least one of the prior uploads finish if other.paralleljobs - // are already running + if ( strings.HasPrefix(absfilepath, backupdir) ) { + j := uploadJobS3{ jobS3: jobS3{otherargs.uniqueid, bkpdir}, absfilepath: absfilepath } + work <- &j // this will hang until at least one of the prior uploads finish if other.paralleljobs + // are already running + } return err }) @@ -262,7 +264,7 @@ func (cn *S3connector) Download(otherargs *OtherArgs, backupinfo *BackupInfo) (e result := make(chan *downloadJobResultS3, otherargs.paralleljobs) done := make(chan bool) - bkpath := filepath.Join(otherargs.uniqueid, "Netezza",backupinfo.npshost, backupinfo.dbname, backupinfo.backupset) + bkpath := filepath.Join(otherargs.uniqueid, "Netezza",backupinfo.npshost, backupinfo.dbname, backupinfo.backupset, backupinfo.increment) // start the workers go func() { for { From 1b6786a7cd3b032d559b0de0f0d8369fbcc79b98 Mon Sep 17 00:00:00 2001 From: nmahalle Date: Thu, 11 Mar 2021 07:15:51 -0800 Subject: [PATCH 7/8] Changed utility name to nzsyncbackup --- bnr-utils/{nz_cp_backup => nzsyncbackup}/connector/AZConnector.go | 0 bnr-utils/{nz_cp_backup => nzsyncbackup}/connector/Connector.go | 0 bnr-utils/{nz_cp_backup => nzsyncbackup}/connector/S3Connector.go | 0 bnr-utils/{nz_cp_backup => nzsyncbackup}/factory/Factory.go | 0 bnr-utils/{nz_cp_backup => nzsyncbackup}/main.go | 0 5 files changed, 0 insertions(+), 0 deletions(-) rename bnr-utils/{nz_cp_backup => nzsyncbackup}/connector/AZConnector.go (100%) rename bnr-utils/{nz_cp_backup => nzsyncbackup}/connector/Connector.go (100%) rename bnr-utils/{nz_cp_backup => nzsyncbackup}/connector/S3Connector.go (100%) rename bnr-utils/{nz_cp_backup => nzsyncbackup}/factory/Factory.go (100%) rename bnr-utils/{nz_cp_backup => nzsyncbackup}/main.go (100%) diff --git a/bnr-utils/nz_cp_backup/connector/AZConnector.go b/bnr-utils/nzsyncbackup/connector/AZConnector.go similarity index 100% rename from bnr-utils/nz_cp_backup/connector/AZConnector.go rename to bnr-utils/nzsyncbackup/connector/AZConnector.go diff --git a/bnr-utils/nz_cp_backup/connector/Connector.go b/bnr-utils/nzsyncbackup/connector/Connector.go similarity index 100% rename from bnr-utils/nz_cp_backup/connector/Connector.go rename to bnr-utils/nzsyncbackup/connector/Connector.go diff --git a/bnr-utils/nz_cp_backup/connector/S3Connector.go b/bnr-utils/nzsyncbackup/connector/S3Connector.go similarity index 100% rename from bnr-utils/nz_cp_backup/connector/S3Connector.go rename to bnr-utils/nzsyncbackup/connector/S3Connector.go diff --git a/bnr-utils/nz_cp_backup/factory/Factory.go b/bnr-utils/nzsyncbackup/factory/Factory.go similarity index 100% rename from bnr-utils/nz_cp_backup/factory/Factory.go rename to bnr-utils/nzsyncbackup/factory/Factory.go diff --git a/bnr-utils/nz_cp_backup/main.go b/bnr-utils/nzsyncbackup/main.go similarity index 100% rename from bnr-utils/nz_cp_backup/main.go rename to bnr-utils/nzsyncbackup/main.go From fac943f5ba7a65bc1daec981972c0281a4dd8a27 Mon Sep 17 00:00:00 2001 From: nmahalle Date: Mon, 12 Apr 2021 09:29:01 -0700 Subject: [PATCH 8/8] Changed Parsing,Logging, removed -cloudBackup parameter --- .../nzsyncbackup/connector/AZConnector.go | 91 +++++------- bnr-utils/nzsyncbackup/connector/Connector.go | 135 ++++++++++++------ .../nzsyncbackup/connector/S3Connector.go | 107 ++++++-------- bnr-utils/nzsyncbackup/factory/Factory.go | 12 +- bnr-utils/nzsyncbackup/main.go | 32 ++--- 5 files changed, 183 insertions(+), 194 deletions(-) diff --git a/bnr-utils/nzsyncbackup/connector/AZConnector.go b/bnr-utils/nzsyncbackup/connector/AZConnector.go index a0855d6..baf8042 100644 --- a/bnr-utils/nzsyncbackup/connector/AZConnector.go +++ b/bnr-utils/nzsyncbackup/connector/AZConnector.go @@ -3,7 +3,6 @@ package Connector import ( "fmt" "strings" - "strconv" "net/url" "time" "os" @@ -25,11 +24,11 @@ type iaz interface { } type AZConnector struct { - azaccount string - azkey string - azcontainer string - streams uint - blocksize int64 + Azaccount string + Azkey string + Azcontainer string + Streams uint + Blocksize int64 } type job struct { @@ -58,31 +57,6 @@ type downloadJobResult struct { err error } -func (c *AZConnector) ParseConnectorArgs(args string) { - arguments := strings.Split(args, ";") - for _, arg := range arguments { - kv := strings.Split(arg, ":") - switch kv[0] { - case "STORAGE_ACCOUNT": - c.azaccount = kv[1] - case "KEY": - c.azkey = kv[1] - case "CONTAINER": - c.azcontainer = kv[1] - case "STREAMS": - u32, err := strconv.ParseUint(kv[1], 10, 32) - if (err == nil ) { - c.streams = uint(u32) - } - case "BLOCKSIZE": - u64, err := strconv.ParseInt(kv[1], 10, 64) - if (err == nil ) { - c.blocksize = int64(u64) - } - } - } -} - func (j *uploadJob) uploadAZ(conn *AZConnector) error { relfilepath, err := filepath.Rel(j.job.bkpdir, j.absfilepath) if err != nil { @@ -95,13 +69,13 @@ func (j *uploadJob) uploadAZ(conn *AZConnector) error { func (cn *AZConnector) getServiceURL() (azblob.ServiceURL, error) { var serviceURL azblob.ServiceURL - us := fmt.Sprintf("https://%s.blob.core.windows.net/", cn.azaccount) + us := fmt.Sprintf("https://%s.blob.core.windows.net/", cn.Azaccount) u, err := url.Parse(us) if err != nil { return serviceURL, fmt.Errorf("Unable to parse URL: %s : %v", us, err) } - credential, err := azblob.NewSharedKeyCredential(cn.azaccount, cn.azkey) + credential, err := azblob.NewSharedKeyCredential(cn.Azaccount, cn.Azkey) if err != nil { return serviceURL, fmt.Errorf("Unable to create shared credentials: %v", err) } @@ -120,7 +94,7 @@ func (cn *AZConnector) getContainerURL() (azblob.ContainerURL, error) { var containerURL azblob.ContainerURL serviceURL, err := cn.getServiceURL() if err == nil { - containerURL = serviceURL.NewContainerURL(cn.azcontainer) + containerURL = serviceURL.NewContainerURL(cn.Azcontainer) } return containerURL, err } @@ -157,8 +131,8 @@ func (cn *AZConnector) uploadFile(absfilepath string, relfilepath string, unique _, err = azblob.UploadFileToBlockBlob(context.Background(), file, blockBlobURL, azblob.UploadToBlockBlobOptions{ - BlockSize: int64(cn.blocksize * 1024 * 1024), - Parallelism: uint16(cn.streams), + BlockSize: int64(cn.Blocksize * 1024 * 1024), + Parallelism: uint16(cn.Streams), }) return err @@ -167,13 +141,13 @@ func (cn *AZConnector) uploadFile(absfilepath string, relfilepath string, unique func (cn *AZConnector) Upload( otherargs *OtherArgs, backupinfo *BackupInfo ) (error){ var err error - log.Println("Uploading Using AZ Connector") dirlist := strings.Split(backupinfo.Dir," ") for _, bkpdir := range dirlist { + log.Println("Uploading backup data to azure cloud from backup dir", bkpdir) backupdir := filepath.Join(bkpdir, "Netezza", backupinfo.npshost, backupinfo.dbname, backupinfo.backupset, backupinfo.increment) _, err = os.Stat(backupdir) if err != nil { - return fmt.Errorf("Error in Creating backupdir : %v", err) + return fmt.Errorf("Cannot access directory '%s': %v. Please check if DB name, hostname are correct.", backupdir, err) } work := make(chan *uploadJob, otherargs.paralleljobs) result := make(chan *jobResult, otherargs.paralleljobs) @@ -208,7 +182,8 @@ func (cn *AZConnector) Upload( otherargs *OtherArgs, backupinfo *BackupInfo ) (e if r.err != nil { // stopping right here so that we // don't keep on uploading when one has failed - log.Fatalf("%s: %v", r.job, r.err) + log.Println("Error while uploading file. Ensure azure storage account name, azure key and container name are correct. If error persists contact IBM support team.", *r.job) + log.Fatalf("Azure storage account:%s accessing container:%s failed with error: %v", cn.Azaccount, cn.Azcontainer, r.err) } filesuploaded++ // this is fine, since this is single threaded increment } @@ -227,6 +202,9 @@ func (cn *AZConnector) Upload( otherargs *OtherArgs, backupinfo *BackupInfo ) (e }) close(work) <- done + if (err != nil) { + return(fmt.Errorf("Error reading directory: %s: %v. Please check if DB name, hostname are correct.", backupdir, err)) + } log.Println("Upload successful for Backup Dir :", bkpdir) log.Println("Total files uploaded:", filesuploaded) } @@ -236,7 +214,7 @@ func (cn *AZConnector) Upload( otherargs *OtherArgs, backupinfo *BackupInfo ) (e func (j *downloadJob) downloadAZ(conn *AZConnector) error { log.Println("Downloading file :", j.blobname) - return conn.downloadFile(j.outfilepath, j.blobname, conn.streams, conn.blocksize) + return conn.downloadFile(j.outfilepath, j.blobname, conn.Streams, conn.Blocksize) } func (cn *AZConnector) downloadFile(outfilepath string, blobname string, streams uint, blockSize int64) error { @@ -269,10 +247,10 @@ func (cn *AZConnector) downloadFile(outfilepath string, blobname string, streams func (cn *AZConnector) Download(otherargs *OtherArgs, backupinfo *BackupInfo) (error){ var err error - log.Println("Downloading Using AZ Connector") + log.Println("Downloading backup data from azure cloud to backup dir", backupinfo.Dir) outdir := backupinfo.Dir - arrayLoc:= []string{} - arrayContents:= []string{} + locations:= []string{} + contents:= []string{} work := make(chan *downloadJob, otherargs.paralleljobs) result := make(chan *downloadJobResult, otherargs.paralleljobs) done := make(chan bool) @@ -328,7 +306,7 @@ func (cn *AZConnector) Download(otherargs *OtherArgs, backupinfo *BackupInfo) (e // Get a result segment starting with the blob indicated by the current Marker. listBlob, err := containerURL.ListBlobsFlatSegment(context.Background(), marker, azblob.ListBlobsSegmentOptions{}) if err != nil { - return fmt.Errorf("Error in listing segment of blobs: %v",err) + return fmt.Errorf("Unable to list segment of blobs with storage account:%s and container:%s. Ensure azure storage account and container are correct.\n Error details: %v",cn.Azaccount, cn.Azcontainer, err) } // ListBlobs returns the start of the next segment; you MUST use this to get @@ -351,14 +329,15 @@ func (cn *AZConnector) Download(otherargs *OtherArgs, backupinfo *BackupInfo) (e if err != nil { return fmt.Errorf("Error in creating backup directory structure: %v",err) } + + switch filename { + case "locations.txt": + locations = append(locations, path.Join(dumpdir, filename)) + case "contents.txt": + contents = append(contents, path.Join(dumpdir, filename)) + } outfilepath := path.Join(dumpdir, filename) - if (strings.HasSuffix(outfilepath,"locations.txt")){ - arrayLoc = append(arrayLoc,outfilepath) - } - if (strings.HasSuffix(outfilepath,"contents.txt")){ - arrayContents = append(arrayContents,outfilepath) - } j := downloadJob{ conn:*cn, outfilepath:outfilepath, blobname: blobInfo.Name } work <- &j @@ -367,21 +346,15 @@ func (cn *AZConnector) Download(otherargs *OtherArgs, backupinfo *BackupInfo) (e } } - if blobfound > 0 { - log.Println("No matching blob found. Please check if DB name, hostname, uniqueid or containername is correct") - return fmt.Errorf("No matching blob found.") - } if blobfound == 0 { - blobfound++ + return fmt.Errorf("No matching blob found. Please check if DB name, hostname, uniqueid or containername are correct. If error persists contact IBM support team. Azaccount:%s AzContainer:%s Blobpath:%s, Uniqueid:%s", cn.Azaccount, cn.Azcontainer, blobpath, otherargs.uniqueid) } } close(work) <- done log.Println("File Downloaded to dir :", outdir) log.Println("Total files downloaded:", filesdownloaded) - if (*otherargs.cloudBackup) { - updateLocation(arrayLoc,outdir) - updateContents(arrayContents) - } + updateLocation(locations,outdir) + updateContents(contents) return err } diff --git a/bnr-utils/nzsyncbackup/connector/Connector.go b/bnr-utils/nzsyncbackup/connector/Connector.go index 53f3738..0b70dba 100644 --- a/bnr-utils/nzsyncbackup/connector/Connector.go +++ b/bnr-utils/nzsyncbackup/connector/Connector.go @@ -1,18 +1,18 @@ package Connector import ( - "flag" "log" "fmt" "time" + "io" "io/ioutil" "os" "path" "strings" + "github.com/spf13/cobra" ) type IConnector interface { - ParseConnectorArgs(string) Upload(*OtherArgs, *BackupInfo) (error) Download( *OtherArgs, *BackupInfo) (error) } @@ -25,48 +25,88 @@ type BackupInfo struct { increment string } -type ConnectorInfo struct { - Connector string - ConnectorArgs string -} - type OtherArgs struct { uniqueid string logfiledir string - upload *bool - download *bool paralleljobs int Operation int - cloudBackup *bool + upload *bool + download *bool + Connector string } -func ParseArgs(backupinfo *BackupInfo, connectorInfo *ConnectorInfo, otherargs *OtherArgs) { - flag.StringVar(&backupinfo.dbname, "db", "", "Database name") - flag.StringVar(&backupinfo.Dir, "dir", "", "Full path to the directory in which the backup already exists or should be downloaded") - flag.StringVar(&backupinfo.npshost, "npshost", "", "Name of the NPS host as it appears in the backups") - flag.StringVar(&backupinfo.backupset, "backupset", "", "Name of the backupset to be uploaded/downloaded") - flag.StringVar(&backupinfo.increment, "increment", "", "Increment Number to be uploaded/downloaded") - - flag.StringVar(&connectorInfo.Connector, "connector", "", "Destination cloud store") - flag.StringVar(&connectorInfo.ConnectorArgs, "connectorArgs", "", "Arguments for cloud store") - - flag.StringVar(&otherargs.uniqueid,"uniqueid", "", "Azure blob storage container") - flag.StringVar(&otherargs.logfiledir,"logfiledir", "/tmp", "Logfile directory for this utility. Default is /tmp dir") - otherargs.upload = flag.Bool("upload", false, "Upload to cloud") - otherargs.download = flag.Bool("download", false, "Download from cloud") - otherargs.cloudBackup = flag.Bool("cloudBackup", false, "Download backup taken on cloud") - flag.IntVar(&otherargs.paralleljobs,"paralleljobs",6,"Number of parallel files to upload/download") +func ParseArgs(backupinfo *BackupInfo, otherargs *OtherArgs, azconnect *AZConnector, s3connect *S3connector) { + var cmdAws = &cobra.Command{ + Use : "aws", + Short: "Upload/Download Backup to/from AWS Cloud", + Run: func(cmd *cobra.Command, args []string) { + otherargs.Connector = "aws" + }, + } + var cmdAzure = &cobra.Command{ + Use : "azure", + Short: "Upload/Download Backup to/From Azure Cloud", + Run: func(cmd *cobra.Command, args []string) { + otherargs.Connector = "azure" + }, + } + + + var rootCmd = &cobra.Command{Use: "nzsyncbackup"} + + + rootCmd.PersistentFlags().StringVar(&backupinfo.dbname, "db", "", "Database name") + rootCmd.PersistentFlags().StringVar(&backupinfo.Dir, "dir", "", "Full path to the directory in which the backup already exists or should be downloaded") + rootCmd.PersistentFlags().StringVar(&backupinfo.npshost, "npshost", "", "Name of the NPS host as it appears in the backups") + rootCmd.PersistentFlags().StringVar(&backupinfo.backupset, "backupset", "", "Name of the backupset to be uploaded/downloaded") + rootCmd.PersistentFlags().StringVar(&backupinfo.increment, "increment", "", "Increment Number to be uploaded/downloaded") + rootCmd.PersistentFlags().StringVar(&otherargs.uniqueid,"uniqueid", "", "Azure blob storage container") + rootCmd.PersistentFlags().StringVar(&otherargs.logfiledir,"logfiledir", "/tmp", "Logfile directory for this utility. Default is /tmp dir") + otherargs.upload = rootCmd.PersistentFlags().Bool("upload", false, "Upload Backup to Cloud") + otherargs.download = rootCmd.PersistentFlags().Bool("download", false, "Download backup from cloud") + rootCmd.PersistentFlags().IntVar(&otherargs.paralleljobs,"paralleljobs",6,"Number of parallel files to upload/download") + + + + + cmdAws.Flags().StringVar(&s3connect.Access_key_id, "access-key", "", "The access key for the object store [AWS_ACCESS_KEY_ID] (required)") + cmdAws.Flags().StringVar(&s3connect.Bucket_url, "bucket-url", "", "The bucket url to store backups to (required)") + cmdAws.Flags().StringVar(&s3connect.Default_region, "region", "", "The region of the object store bucket (required)") + cmdAws.Flags().StringVar(&s3connect.Secret_access_key, "secret-key", "", "The secret key for the object store [AWS_SECRET_ACCESS_KEY] (required)") + cmdAws.Flags().StringVar(&s3connect.Endpoint, "endpoint", "", "The endpoint for object to be store on IBM cloud") + cmdAws.Flags().IntVar(&s3connect.Streams, "streams", 16, "Number of blocks to upload/download in parallel default 16") + cmdAws.Flags().Int64Var(&s3connect.Blocksize, "blocksize", 100, "Block size in MB to upload/download file") + + cmdAzure.Flags().StringVar(&azconnect.Azkey, "account-key", "", "The Azure Blob account key (required)") + cmdAzure.Flags().StringVar(&azconnect.Azaccount, "account-name", "", "The Azure Blob account name (required)") + cmdAzure.Flags().UintVar(&azconnect.Streams, "streams", 16, "Number of blocks to upload/download in parallel default 16") + cmdAzure.Flags().Int64Var(&azconnect.Blocksize, "blocksize", 100, "Block size in MB to upload/download file") + cmdAzure.Flags().StringVar(&azconnect.Azcontainer, "container", "", "The Azure Blob container name (required)") + cmdAws.MarkFlagRequired("region") + cmdAws.MarkFlagRequired("access-key") + cmdAws.MarkFlagRequired("bucket-name") + cmdAws.MarkFlagRequired("secret-key") + cmdAzure.MarkFlagRequired("account-key") + cmdAzure.MarkFlagRequired("account-name") + cmdAzure.MarkFlagRequired("container") + rootCmd.AddCommand(cmdAws, cmdAzure) + + if err := rootCmd.Execute(); err != nil { + log.Println(err) + os.Exit(1) + } } -func SetUpLogFile(backupinfo *BackupInfo, connectorInfo *ConnectorInfo, otherargs *OtherArgs) { +func SetUpLogFile(backupinfo *BackupInfo, otherargs *OtherArgs) { // log file configuration setup - logfilename := fmt.Sprintf("nz_%sConnector_%d_%s.log", connectorInfo.Connector, os.Getppid(), time.Now().Format("2006-01-02-150405")) + logfilename := fmt.Sprintf("nz_%sConnector_%d_%s.log", otherargs.Connector, os.Getppid(), time.Now().Format("2006-01-02-150405")) logfilepath := path.Join(otherargs.logfiledir, logfilename) filehandle, err := os.OpenFile(logfilepath, os.O_WRONLY|os.O_CREATE|os.O_APPEND, 0644) if err != nil { fmt.Errorf("Error in opening logfile: %v",err) } - log.SetOutput(filehandle) + w := io.MultiWriter(os.Stdout, filehandle) + log.SetOutput(w) prefixStr := fmt.Sprintf("%s ", time.Now().UTC().Format("2006-01-02 15:04:05 EST")) + fmt.Sprintf("%-7s", "[INFO]") log.SetFlags(0) log.SetPrefix(prefixStr) @@ -90,15 +130,22 @@ func SetOperation(otherargs *OtherArgs) { func updateLocation(arrLoc []string,outdir string){ for _,locFile := range arrLoc{ - f, err := os.OpenFile(locFile, - os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644) + input, err := ioutil.ReadFile(locFile) if err != nil { - log.Fatalln(err) + log.Fatalf("Unable to open %s to read: %v\n", locFile, err) } - defer f.Close() - textAppend := "1,1,1," + outdir - if _, err := f.WriteString(textAppend); err != nil { - log.Fatalln(err) + lines := strings.Split(string(input), "\n") + if (len(lines) == 2 && !strings.HasSuffix(lines[len(lines) -2] , outdir)) { + f, err := os.OpenFile(locFile, + os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644) + if err != nil { + log.Fatalf("Unable to open %s for update: %v\n", locFile, err) + } + defer f.Close() + textAppend := "1,1,1," + outdir + "\n" + if _, err := f.WriteString(textAppend); err != nil { + log.Fatalf("Unable to update %s: %v\n", locFile, err) + } } } } @@ -107,21 +154,23 @@ func updateContents(arrContents []string){ for _,contentFile := range arrContents{ input, err := ioutil.ReadFile(contentFile) if err != nil { - log.Fatalln(err) + log.Fatalln("Unable to open %s to read: %v\n",contentFile,err) } lines := strings.Split(string(input), "\n") - lines = lines[:len(lines)-1] var textline []string - for _, line := range lines { - r := []rune(line) - str := string(r[:len(r)-1]) + "1" - textline = append(textline,str) + for i := 0 ; i < len(lines) ; i++ { + line := lines[i] + token := strings.Split(line, ",") + if ( token[len(token)-1] == "0" ) { + token[len(token)-1] = "1" + } + textline = append(textline, strings.Join(token, ",")) } output := strings.Join(textline, "\n") err = ioutil.WriteFile(contentFile, []byte(output), 0644) if err != nil { - log.Fatalln(err) + log.Fatalln("Unable to update %s: %v\n",contentFile,err) } } } diff --git a/bnr-utils/nzsyncbackup/connector/S3Connector.go b/bnr-utils/nzsyncbackup/connector/S3Connector.go index 8ed831a..59d64b3 100644 --- a/bnr-utils/nzsyncbackup/connector/S3Connector.go +++ b/bnr-utils/nzsyncbackup/connector/S3Connector.go @@ -3,7 +3,6 @@ package Connector import ( "fmt" "strings" - "strconv" "os" "path/filepath" "path" @@ -24,13 +23,13 @@ type is3 interface { } type S3connector struct { - access_key_id string - secret_access_key string - default_region string - bucket_url string - endpoint string - streams int - blocksize int64 + Access_key_id string + Secret_access_key string + Default_region string + Bucket_url string + Endpoint string + Streams int + Blocksize int64 } type jobS3 struct { @@ -80,7 +79,7 @@ func (c *S3connector) uploadFileS3(absfilepath string, relfilepath string, uniqu } result,err := uploader.Upload(&s3manager.UploadInput{ - Bucket: &c.bucket_url, + Bucket: &c.Bucket_url, Key: aws.String(filepath.Join(uniqueid, relfilepath)), Body: file, }) @@ -106,7 +105,7 @@ func (c *S3connector) downloadFile(outfilepath string, key string, conn *s3manag numBytes, err := conn.Download(filehandle, &s3.GetObjectInput{ - Bucket: aws.String(c.bucket_url), + Bucket: aws.String(c.Bucket_url), Key: aws.String(key), }) @@ -119,8 +118,8 @@ func (c *S3connector) downloadFile(outfilepath string, key string, conn *s3manag func (c *S3connector) getSession() (sess *session.Session) { sess, err := session.NewSession(&aws.Config{ - Region: aws.String(c.default_region), - Credentials: credentials.NewStaticCredentials(c.access_key_id,c.secret_access_key,""), + Region: aws.String(c.Default_region), + Credentials: credentials.NewStaticCredentials(c.Access_key_id,c.Secret_access_key,""), CredentialsChainVerboseErrors: aws.Bool(true) }) if (err != nil) { log.Fatalln("Session failed:", err) @@ -133,8 +132,8 @@ func (c *S3connector) getUploader() (*s3manager.Uploader) { sessn := c.getSession() uploader := s3manager.NewUploader(sessn,func(u *s3manager.Uploader) { - u.PartSize = c.blocksize * 1024 * 1024 // 64MB per part - u.Concurrency = c.streams + u.PartSize = c.Blocksize * 1024 * 1024 // 64MB per part + u.Concurrency = c.Streams }) return uploader @@ -144,53 +143,24 @@ func (c *S3connector) getDownloader() (*s3manager.Downloader, *session.Session) sessn := c.getSession() downloader := s3manager.NewDownloader(sessn,func(d *s3manager.Downloader) { - d.PartSize = c.blocksize * 1024 * 1024 // 64MB per part - d.Concurrency = c.streams + d.PartSize = c.Blocksize * 1024 * 1024 // 64MB per part + d.Concurrency = c.Streams }) return downloader,sessn } -func (c *S3connector) ParseConnectorArgs(args string) { - arguments := strings.Split(args, ";") - for _, arg := range arguments { - kv := strings.Split(arg, ":") - switch kv[0] { - case "ACCESS_KEY_ID": - c.access_key_id = kv[1] - case "SECRET_ACCESS_KEY": - c.secret_access_key = kv[1] - case "DEFAULT_REGION": - c.default_region = kv[1] - case "BUCKET_URL": - c.bucket_url = kv[1] - case "ENDPOINT": - c.endpoint = kv[1] - case "STREAMS": - i, err := strconv.Atoi(kv[1]) - if (err == nil ) { - c.streams = i - } - case "BLOCKSIZE": - u64, err := strconv.ParseInt(kv[1], 10, 64) - if (err == nil ) { - c.blocksize = int64(u64) - } - - } - } -} - func (c *S3connector) Upload( otherargs *OtherArgs, backupinfo *BackupInfo ) (error){ var err error log.Println("Uploading Using S3 Connector") dirlist := strings.Split(backupinfo.Dir," ") for _, bkpdir := range dirlist { + log.Println("Uploading backup data to aws s3 cloud from backup dir", bkpdir) backupdir := filepath.Join(bkpdir, "Netezza", backupinfo.npshost, backupinfo.dbname, backupinfo.backupset, backupinfo.increment) _, err = os.Stat(backupdir) if err != nil { - return fmt.Errorf("Error Directory not present : %v", err) + return fmt.Errorf("Cannot access directory '%s': %v. Please check if DB name, hostname are correct.", backupdir, err) } work := make(chan *uploadJobS3, otherargs.paralleljobs) @@ -226,7 +196,8 @@ func (c *S3connector) Upload( otherargs *OtherArgs, backupinfo *BackupInfo ) (er if r.err != nil { // stopping right here so that we // don't keep on uploading when one has failed - log.Fatalf("%s: %v", r.jobS3, r.err) + log.Println("Error while uploading file. Ensure aws s3 access-key-id, secret-access-key, bucket_url are correct. If error persists contact IBM support team.", *r.jobS3) + log.Fatalf("Failed to access AWS bucket: %s with error: %v", c.Bucket_url, r.err) } filesuploaded++ // this is fine, since this is single threaded increment } @@ -248,6 +219,9 @@ func (c *S3connector) Upload( otherargs *OtherArgs, backupinfo *BackupInfo ) (er close(work) <- done + if (err != nil) { + return(fmt.Errorf("Error reading directory: %s: %v. Please check if DB name, hostname are correct.", backupdir, err)) + } log.Println("Upload using S3 connector successful for directory :", bkpdir) log.Println("Total files uploaded:", filesuploaded) } @@ -255,11 +229,11 @@ func (c *S3connector) Upload( otherargs *OtherArgs, backupinfo *BackupInfo ) (er } func (cn *S3connector) Download(otherargs *OtherArgs, backupinfo *BackupInfo) (error){ - log.Println("Downloading Using S3 Connector") + log.Println("Downloading backup data from aws cloud to backup dir", backupinfo.Dir) var err error outdir := backupinfo.Dir - arrayLoc:= []string{} - arrayContents:= []string{} + locations:= []string{} + contents:= []string{} work := make(chan *downloadJobS3, otherargs.paralleljobs) result := make(chan *downloadJobResultS3, otherargs.paralleljobs) done := make(chan bool) @@ -303,7 +277,7 @@ func (cn *S3connector) Download(otherargs *OtherArgs, backupinfo *BackupInfo) (e down,sess := cn.getDownloader() client := s3.New(sess) - params := &s3.ListObjectsInput{Bucket: &cn.bucket_url, Prefix: &otherargs.uniqueid} + params := &s3.ListObjectsInput{Bucket: &cn.Bucket_url, Prefix: &otherargs.uniqueid} client.ListObjectsPages(params, func(page *s3.ListObjectsOutput, more bool) (bool) { for _, obj := range page.Contents { key := *obj.Key @@ -317,19 +291,20 @@ func (cn *S3connector) Download(otherargs *OtherArgs, backupinfo *BackupInfo) (e log.Fatalf("Error in fetching download relative path: %v",err) } - file := filepath.Join(outdir, relfilepath) - err = os.MkdirAll(file, 0777) + dumpdir := filepath.Join(outdir, relfilepath) + err = os.MkdirAll(dumpdir, 0777) if err != nil { log.Fatalf("Error in creating backup directory structure: %v",err) } - outfilepath := path.Join(file, filename) - if (strings.HasSuffix(outfilepath,"locations.txt")){ - arrayLoc = append(arrayLoc,outfilepath) - } - if (strings.HasSuffix(outfilepath,"contents.txt")){ - arrayContents = append(arrayContents,outfilepath) + switch filename { + case "locations.txt": + locations = append(locations, path.Join(dumpdir, filename)) + case "contents.txt": + contents = append(contents, path.Join(dumpdir, filename)) } + + outfilepath := path.Join(dumpdir, filename) j := downloadJobS3{ conn:down, key:key, outfilepath:outfilepath } work <- &j } @@ -337,12 +312,10 @@ func (cn *S3connector) Download(otherargs *OtherArgs, backupinfo *BackupInfo) (e return true }) - close(work) - <- done - log.Println("Total files downloaded using S3 Connector:", filesdownloaded) - if (*otherargs.cloudBackup) { - updateLocation(arrayLoc,outdir) - updateContents(arrayContents) - } + close(work) + <- done + log.Println("Total files downloaded using S3 Connector:", filesdownloaded) + updateLocation(locations,outdir) + updateContents(contents) return err } diff --git a/bnr-utils/nzsyncbackup/factory/Factory.go b/bnr-utils/nzsyncbackup/factory/Factory.go index af0582b..7eaf89d 100644 --- a/bnr-utils/nzsyncbackup/factory/Factory.go +++ b/bnr-utils/nzsyncbackup/factory/Factory.go @@ -1,15 +1,15 @@ package Factory import ( - "nzconnector/connector" + "nzsyncbackup/connector" ) -func GetConnector(connectorType string) Connector.IConnector { +func GetConnector(connectorType string, azconnect *Connector.AZConnector, s3connect *Connector.S3connector) Connector.IConnector { switch connectorType { - case "s3": - return &Connector.S3connector{} - case "az": - return &Connector.AZConnector{} + case "aws": + return s3connect + case "azure": + return azconnect default: return nil } diff --git a/bnr-utils/nzsyncbackup/main.go b/bnr-utils/nzsyncbackup/main.go index 65dbd48..43dd962 100644 --- a/bnr-utils/nzsyncbackup/main.go +++ b/bnr-utils/nzsyncbackup/main.go @@ -1,10 +1,9 @@ package main import ( - "flag" "log" - "nzconnector/connector" - "nzconnector/factory" + "nzsyncbackup/connector" + "nzsyncbackup/factory" ) const( @@ -12,38 +11,33 @@ const( download ) -func parseConnectorArgs(e Connector.IConnector, args string) { - if e != nil { - e.ParseConnectorArgs(args) - } -} - func main() { var backupinfo Connector.BackupInfo - var connectorInfo Connector.ConnectorInfo var otherargs Connector.OtherArgs + var azconnect Connector.AZConnector + var s3connect Connector.S3connector var err error // parse input args - Connector.ParseArgs(&backupinfo, &connectorInfo, &otherargs) - flag.Parse() - Connector.SetUpLogFile(&backupinfo, &connectorInfo, &otherargs) - connector := Factory.GetConnector(connectorInfo.Connector) - parseConnectorArgs(connector, connectorInfo.ConnectorArgs) + Connector.ParseArgs(&backupinfo, &otherargs, &azconnect, &s3connect) + connector := Factory.GetConnector(otherargs.Connector, &azconnect, &s3connect) + if ( connector != nil ) { Connector.SetOperation(&otherargs) - + log.Println("OPeration is ", otherargs.Operation) switch otherargs.Operation { case upload: + Connector.SetUpLogFile(&backupinfo, &otherargs) // now do the upload - log.Println("Uploading backup data to cloud for conector : ",connectorInfo.Connector) + log.Println("Uploading backup data to cloud for conector : ",otherargs.Connector) err = connector.Upload(&otherargs, &backupinfo) if (err != nil) { log.Fatalln(err) } log.Println("Upload successful") case download: - log.Println("Downloading backup data from cloud for connector : ", connectorInfo.Connector) + Connector.SetUpLogFile(&backupinfo, &otherargs) + log.Println("Downloading backup data from cloud for connector : ", otherargs.Connector) err = connector.Download(&otherargs, &backupinfo) if (err != nil) { log.Fatalln(err) @@ -52,5 +46,5 @@ func main() { default: log.Fatalln("Invalid Operation, Supported Upload/Download") } - + } }