This is an automated email from the ASF dual-hosted git repository.

miaoliyao pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/shardingsphere-on-cloud.git


The following commit(s) were added to refs/heads/main by this push:
     new 3e912c7  feat:wrapping the backup subcommand (#202)
3e912c7 is described below

commit 3e912c7a8a0e2dbe33014930f2c21807a87f6166
Author: lltgo <[email protected]>
AuthorDate: Tue Feb 14 14:53:13 2023 +0800

    feat:wrapping the backup subcommand (#202)
    
    * feat: wrapping the backup subcommand
    
    * feat: get the backup id by regex
    
    * chore: revert test
    
    * chore: update License
---
 pitr/agent/go.mod                                  |  1 +
 pitr/agent/go.sum                                  |  2 +
 pitr/agent/internal/pkg/opengauss.go               | 57 ++++++++++++++++++++++
 .../cmd_test.go => internal/pkg/opengauss_test.go} | 32 ++++++------
 pitr/agent/pkg/cmds/cmd.go                         | 44 ++++++++---------
 pitr/agent/pkg/cmds/cmd_test.go                    | 12 +++--
 6 files changed, 102 insertions(+), 46 deletions(-)

diff --git a/pitr/agent/go.mod b/pitr/agent/go.mod
index 8117a69..74ac0c0 100644
--- a/pitr/agent/go.mod
+++ b/pitr/agent/go.mod
@@ -11,6 +11,7 @@ require (
 
 require (
        github.com/andybalholm/brotli v1.0.4 // indirect
+       github.com/dlclark/regexp2 v1.8.0 // indirect
        github.com/go-logr/logr v1.2.3 // indirect
        github.com/google/go-cmp v0.5.9 // indirect
        github.com/google/uuid v1.3.0 // indirect
diff --git a/pitr/agent/go.sum b/pitr/agent/go.sum
index aaa8ca3..93702c8 100644
--- a/pitr/agent/go.sum
+++ b/pitr/agent/go.sum
@@ -5,6 +5,8 @@ github.com/creack/pty v1.1.9/go.mod 
h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ3
 github.com/davecgh/go-spew v1.1.0/go.mod 
h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
 github.com/davecgh/go-spew v1.1.1 
h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
 github.com/davecgh/go-spew v1.1.1/go.mod 
h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
+github.com/dlclark/regexp2 v1.8.0 
h1:rJD5HeGIT/2b5CDk63FVCwZA3qgYElfg+oQK7uH5pfE=
+github.com/dlclark/regexp2 v1.8.0/go.mod 
h1:DHkYz0B9wPfa6wondMfaivmHpzrQ3v9q8cnmRbL6yW8=
 github.com/go-logr/logr v1.2.3 h1:2DntVwHkVopvECVRSlL5PSo9eG+cAkDCuckLubN+rq0=
 github.com/go-logr/logr v1.2.3/go.mod 
h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A=
 github.com/gofiber/fiber/v2 v2.42.0 
h1:Fnp7ybWvS+sjNQsFvkhf4G8OhXswvB6Vee8hM/LyS+8=
diff --git a/pitr/agent/internal/pkg/opengauss.go 
b/pitr/agent/internal/pkg/opengauss.go
index f68e726..356770c 100644
--- a/pitr/agent/internal/pkg/opengauss.go
+++ b/pitr/agent/internal/pkg/opengauss.go
@@ -16,3 +16,60 @@
  */
 
 package pkg
+
+import (
+       "fmt"
+
+       "github.com/dlclark/regexp2"
+
+       "github.com/apache/shardingsphere-on-cloud/pitr/agent/pkg/cmds"
+)
+
+type openGauss struct {
+       shell string
+}
+
+const (
+       _backupFmt = "gs_probackup backup --backup-path=%s --instance=%s 
--backup-mode=%s --pgdata=%s 2>&1"
+)
+
+func (og *openGauss) AsyncBackup(backupPath, instanceName, backupMode, pgData 
string) (string, error) {
+       cmd := fmt.Sprintf(_backupFmt, backupPath, instanceName, backupMode, 
pgData)
+       outputs, err := cmds.Commands(og.shell, fmt.Sprintf(_backupFmt, 
backupPath, instanceName, backupMode, pgData))
+       if err != nil {
+               return "", fmt.Errorf("cmds.Commands[shell=%s,cmd=%s] return 
err=%w", og.shell, cmd, err)
+       }
+
+       for output := range outputs {
+               if output.Error != nil {
+                       return "", fmt.Errorf("output.Error[%w] is not nil", 
output.Error)
+               }
+
+               // get the backup id from the first line
+               bid, err := og.getBackupID(output.Message)
+               if err != nil {
+                       return "", fmt.Errorf("og.getBackupID[source=%s] return 
err=%w", output.Message, err)
+               }
+               //ignore other output
+               go og.ignore(outputs)
+               return bid, nil
+       }
+       return "", fmt.Errorf("unknow err")
+}
+
+func (og *openGauss) ignore(outputs chan *cmds.Output) {
+       defer func() {
+               _ = recover()
+       }()
+
+       for range outputs {
+               //ignore all
+       }
+       //outputs closed
+}
+
+func (og *openGauss) getBackupID(msg string) (string, error) {
+       re := regexp2.MustCompile("(?<=backup ID:\\s+)\\w+(?=,)", 0)
+       match, err := re.FindStringMatch(msg)
+       return match.String(), err
+}
diff --git a/pitr/agent/pkg/cmds/cmd_test.go 
b/pitr/agent/internal/pkg/opengauss_test.go
similarity index 74%
copy from pitr/agent/pkg/cmds/cmd_test.go
copy to pitr/agent/internal/pkg/opengauss_test.go
index e93eb77..0107bc8 100644
--- a/pitr/agent/pkg/cmds/cmd_test.go
+++ b/pitr/agent/internal/pkg/opengauss_test.go
@@ -15,32 +15,28 @@
 * limitations under the License.
  */
 
-package cmds
+package pkg
 
 import (
        "fmt"
        "testing"
+       "time"
 )
 
-const (
-       sh = "/bin/sh"
-)
-
-func TestCommand(t *testing.T) {
-       output, err := Commands(sh, "ping www.baidu.com")
+func TestOpenGauss_AsyncBackup(t *testing.T) {
+       og := &openGauss{
+               shell: "/bin/sh",
+       }
+       backupID, err := og.AsyncBackup(
+               "/home/omm/data",
+               "ins-default-0",
+               "full",
+               "/data/opengauss/3.1.1/data/single_node/",
+       )
        if err != nil {
                t.Fatal(err)
        }
+       fmt.Println(backupID)
 
-       for {
-               select {
-               case out, ok := <-output:
-                       if ok {
-                               fmt.Print(out.LineNo, "\t", out.Message)
-                       } else {
-                               return
-                       }
-               }
-       }
-
+       time.Sleep(time.Second * 10)
 }
diff --git a/pitr/agent/pkg/cmds/cmd.go b/pitr/agent/pkg/cmds/cmd.go
index 513172f..23c78c7 100644
--- a/pitr/agent/pkg/cmds/cmd.go
+++ b/pitr/agent/pkg/cmds/cmd.go
@@ -20,7 +20,6 @@ package cmds
 import (
        "bufio"
        "fmt"
-       "io"
        "os/exec"
 
        "github.com/apache/shardingsphere-on-cloud/pitr/agent/pkg/syncutils"
@@ -42,39 +41,34 @@ func Commands(name string, args ...string) (chan *Output, 
error) {
        if err != nil {
                return nil, fmt.Errorf("can not obtain stdout pipe for 
command[args=%+v]:%s", args, err)
        }
-       if err := cmd.Start(); err != nil {
+       if err = cmd.Start(); err != nil {
                return nil, fmt.Errorf("the command is err[args=%+v]:%s", args, 
err)
        }
 
-       reader := bufio.NewReader(stdout)
-
-       output := make(chan *Output, 10)
-       index := uint32(1)
-
+       var (
+               scanner = bufio.NewScanner(stdout)
+               output  = make(chan *Output)
+               index   = uint32(1)
+       )
        go func() {
-               if err := syncutils.NewRecoverFuncWithErrRet("", func() error {
-                       for {
-                               msg, err := reader.ReadString('\n')
-                               if io.EOF == err {
-                                       goto end
-                               } else if err != nil {
-                                       output <- &Output{
-                                               LineNo:  index,
-                                               Message: msg,
-                                               Error:   err,
-                                       }
-                                       goto end
-                               }
-
+               if err = syncutils.NewRecoverFuncWithErrRet("", func() error {
+                       for scanner.Scan() {
                                output <- &Output{
                                        LineNo:  index,
-                                       Message: msg,
+                                       Message: scanner.Text(),
+                                       Error:   err,
                                }
-
                                index++
                        }
-               end:
-                       if err := cmd.Wait(); err != nil {
+
+                       if err = scanner.Err(); err != nil {
+                               output <- &Output{
+                                       LineNo: index,
+                                       Error:  err,
+                               }
+                       }
+
+                       if err = cmd.Wait(); err != nil {
                                output <- &Output{
                                        Error: err,
                                }
diff --git a/pitr/agent/pkg/cmds/cmd_test.go b/pitr/agent/pkg/cmds/cmd_test.go
index e93eb77..62001ab 100644
--- a/pitr/agent/pkg/cmds/cmd_test.go
+++ b/pitr/agent/pkg/cmds/cmd_test.go
@@ -26,8 +26,11 @@ const (
        sh = "/bin/sh"
 )
 
+var backup = "gs_probackup backup -B /home/omm/data --instance=ins-default-0 
-b full -D /data/opengauss/3.1.1/data/single_node/  2>&1"
+var ping = "ping www.baidu.com"
+
 func TestCommand(t *testing.T) {
-       output, err := Commands(sh, "ping www.baidu.com")
+       output, err := Commands(sh, backup)
        if err != nil {
                t.Fatal(err)
        }
@@ -36,11 +39,14 @@ func TestCommand(t *testing.T) {
                select {
                case out, ok := <-output:
                        if ok {
-                               fmt.Print(out.LineNo, "\t", out.Message)
+                               if out.Error != nil {
+                                       fmt.Println(out.LineNo, "\t", 
out.Error.Error())
+                               } else {
+                                       fmt.Println(out.LineNo, "\t", 
out.Message)
+                               }
                        } else {
                                return
                        }
                }
        }
-
 }

Reply via email to