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

woblerr pushed a commit to branch upgrade-go-1.25
in repository https://gitbox.apache.org/repos/asf/cloudberry-backup.git

commit 2861f472b71a9949bed5e91d8b9f98fa48ad24d4
Author: woblerr <[email protected]>
AuthorDate: Fri Jun 26 13:09:08 2026 +0300

    Fix lock e2e test transaction handling.
    The deadlock and signal handler e2e tests still held locks with raw 
BEGIN/COMMIT SQL on DBConn and reused those connections for pg_locks polling.
    With pgx v5 this can release the raw transaction lock or leave goroutines 
polling on closed connections, causing missing blocked-lock counts and 
DBConn.Get panics.
    Use explicit transactions, dedicated lock and polling connections, and 
synchronize goroutines before closing test connections.
---
 end_to_end/locks_test.go          | 119 +++++++++++++++++++++++++++--------
 end_to_end/signal_handler_test.go | 129 +++++++++++++++-----------------------
 2 files changed, 142 insertions(+), 106 deletions(-)

diff --git a/end_to_end/locks_test.go b/end_to_end/locks_test.go
index b05b1d17..046a5730 100644
--- a/end_to_end/locks_test.go
+++ b/end_to_end/locks_test.go
@@ -28,7 +28,11 @@ var _ = Describe("Deadlock handling", func() {
                }
                // Acquire AccessExclusiveLock on public.foo to block gpbackup 
when it attempts
                // to grab AccessShareLocks before its metadata dump section.
-               backupConn.MustExec("BEGIN; LOCK TABLE public.foo IN ACCESS 
EXCLUSIVE MODE")
+               initialLockConn := testutils.SetupTestDbConn("testdb")
+               defer initialLockConn.Close()
+               initialLockConn.MustBegin()
+               defer func() { _ = initialLockConn.Rollback() }()
+               initialLockConn.MustExec("LOCK TABLE public.foo IN ACCESS 
EXCLUSIVE MODE")
 
                // Execute gpbackup with --jobs 10 since there are 10 tables to 
back up
                args := []string{
@@ -37,13 +41,25 @@ var _ = Describe("Deadlock handling", func() {
                        "--jobs", "10",
                        "--verbose"}
                cmd := exec.Command(gpbackupPath, args...)
+
+               var wg sync.WaitGroup
+               releaseTriggerLock := make(chan struct{})
+               var releaseTriggerLockOnce sync.Once
+               defer func() {
+                       releaseTriggerLockOnce.Do(func() { 
close(releaseTriggerLock) })
+                       wg.Wait()
+               }()
+
                // Concurrently wait for gpbackup to block when it requests an 
AccessShareLock on public.foo. Once
                // that happens, acquire an AccessExclusiveLock on 
pg_catalog.pg_trigger to block gpbackup during its
                // trigger metadata dump. Then release the initial 
AccessExclusiveLock on public.foo (from the
                // beginning of the test) to unblock gpbackup and let gpbackup 
move forward to the trigger metadata dump.
-               anotherConn := testutils.SetupTestDbConn("testdb")
-               defer anotherConn.Close()
+               wg.Add(1)
                go func() {
+                       defer wg.Done()
+                       triggerLockConn := testutils.SetupTestDbConn("testdb")
+                       defer triggerLockConn.Close()
+
                        // Query to see if gpbackup's AccessShareLock request 
on public.foo is blocked
                        checkLockQuery := `SELECT count(*) FROM pg_locks l, 
pg_class c, pg_namespace n WHERE l.relation = c.oid AND n.oid = c.relnamespace 
AND n.nspname = 'public' AND c.relname = 'foo' AND l.granted = 'f' AND l.mode = 
'AccessShareLock'`
 
@@ -51,7 +67,7 @@ var _ = Describe("Deadlock handling", func() {
                        var gpbackupBlockedLockCount int
                        iterations := 100
                        for iterations > 0 {
-                               _ = anotherConn.Get(&gpbackupBlockedLockCount, 
checkLockQuery)
+                               _ = 
triggerLockConn.Get(&gpbackupBlockedLockCount, checkLockQuery)
                                if gpbackupBlockedLockCount < 1 {
                                        time.Sleep(100 * time.Millisecond)
                                        iterations--
@@ -64,8 +80,12 @@ var _ = Describe("Deadlock handling", func() {
                        // during the trigger metadata dump so that the test 
can queue a bunch of
                        // AccessExclusiveLock requests against the test 
tables. Afterwards, release the
                        // AccessExclusiveLock on public.foo to let gpbackup go 
to the trigger metadata dump.
-                       anotherConn.MustExec(`BEGIN; LOCK TABLE 
pg_catalog.pg_trigger IN ACCESS EXCLUSIVE MODE`)
-                       backupConn.MustExec("COMMIT")
+                       triggerLockConn.MustBegin()
+                       defer func() { _ = triggerLockConn.Rollback() }()
+                       triggerLockConn.MustExec(`LOCK TABLE 
pg_catalog.pg_trigger IN ACCESS EXCLUSIVE MODE`)
+                       initialLockConn.MustCommit()
+                       <-releaseTriggerLock
+                       triggerLockConn.MustCommit()
                }()
 
                // Concurrently wait for gpbackup to block on the trigger 
metadata dump section. Once we
@@ -74,7 +94,9 @@ var _ = Describe("Deadlock handling", func() {
                dataTables := []string{`public."FOObar"`, "public.foo", 
"public.holds", "public.sales", "public.bigtable",
                        "schema2.ao1", "schema2.ao2", "schema2.foo2", 
"schema2.foo3", "schema2.returns"}
                for _, dataTable := range dataTables {
+                       wg.Add(1)
                        go func(dataTable string) {
+                               defer wg.Done()
                                accessExclusiveLockConn := 
testutils.SetupTestDbConn("testdb")
                                defer accessExclusiveLockConn.Close()
 
@@ -96,23 +118,32 @@ var _ = Describe("Deadlock handling", func() {
 
                                // Queue an AccessExclusiveLock request on a 
test table which will later
                                // result in a detected deadlock during the 
gpbackup data dump section.
-                               
accessExclusiveLockConn.MustExec(fmt.Sprintf(`BEGIN; LOCK TABLE %s IN ACCESS 
EXCLUSIVE MODE; COMMIT`, dataTable))
+                               accessExclusiveLockConn.MustBegin()
+                               defer func() { _ = 
accessExclusiveLockConn.Rollback() }()
+                               
accessExclusiveLockConn.MustExec(fmt.Sprintf(`LOCK TABLE %s IN ACCESS EXCLUSIVE 
MODE`, dataTable))
+                               accessExclusiveLockConn.MustCommit()
                        }(dataTable)
                }
 
                // Concurrently wait for all AccessExclusiveLock requests on 
all 10 test tables to block.
                // Once that happens, release the AccessExclusiveLock on 
pg_catalog.pg_trigger to unblock
                // gpbackup and let gpbackup move forward to the data dump 
section.
-               var accessExclBlockedLockCount int
+               accessExclBlockedLockCountChan := make(chan int, 1)
+               wg.Add(1)
                go func() {
+                       defer wg.Done()
+                       accessExclBlockedLockConn := 
testutils.SetupTestDbConn("testdb")
+                       defer accessExclBlockedLockConn.Close()
+
                        // Query to check for ungranted AccessExclusiveLock 
requests on our test tables
                        checkLockQuery := `SELECT count(*) FROM pg_locks WHERE 
granted = 'f' AND mode = 'AccessExclusiveLock'`
 
                        // Wait up to 10 seconds
+                       var accessExclBlockedLockCount int
                        iterations := 100
                        for iterations > 0 {
-                               _ = backupConn.Get(&accessExclBlockedLockCount, 
checkLockQuery)
-                               if accessExclBlockedLockCount < 10 {
+                               _ = 
accessExclBlockedLockConn.Get(&accessExclBlockedLockCount, checkLockQuery)
+                               if accessExclBlockedLockCount < len(dataTables) 
{
                                        time.Sleep(100 * time.Millisecond)
                                        iterations--
                                } else {
@@ -121,15 +152,18 @@ var _ = Describe("Deadlock handling", func() {
                        }
 
                        // Unblock gpbackup by releasing AccessExclusiveLock on 
pg_catalog.pg_trigger
-                       anotherConn.MustExec("COMMIT")
+                       accessExclBlockedLockCountChan <- 
accessExclBlockedLockCount
+                       releaseTriggerLockOnce.Do(func() { 
close(releaseTriggerLock) })
                }()
 
                // gpbackup has finished
-               output, _ := cmd.CombinedOutput()
+               output, err := cmd.CombinedOutput()
                stdout := string(output)
+               accessExclBlockedLockCount := <-accessExclBlockedLockCountChan
+               Expect(err).ToNot(HaveOccurred(), "%s", stdout)
 
                // Check that 10 deadlock traps were placed during the test
-               Expect(accessExclBlockedLockCount).To(Equal(10))
+               Expect(accessExclBlockedLockCount).To(Equal(len(dataTables)))
                // No non-main worker should have been able to run COPY due to 
deadlock detection
                for i := 1; i < 10; i++ {
                        expectedLockString := fmt.Sprintf("[DEBUG]:-Worker %d: 
LOCK TABLE ", i)
@@ -156,7 +190,11 @@ var _ = Describe("Deadlock handling", func() {
                }
                // Acquire AccessExclusiveLock on public.foo to block gpbackup 
when it attempts
                // to grab AccessShareLocks before its metadata dump section.
-               backupConn.MustExec("BEGIN; LOCK TABLE public.foo IN ACCESS 
EXCLUSIVE MODE")
+               initialLockConn := testutils.SetupTestDbConn("testdb")
+               defer initialLockConn.Close()
+               initialLockConn.MustBegin()
+               defer func() { _ = initialLockConn.Rollback() }()
+               initialLockConn.MustExec("LOCK TABLE public.foo IN ACCESS 
EXCLUSIVE MODE")
 
                // Execute gpbackup with --copy-queue-size 2
                args := []string{
@@ -167,13 +205,24 @@ var _ = Describe("Deadlock handling", func() {
                        "--verbose"}
                cmd := exec.Command(gpbackupPath, args...)
 
+               var wg sync.WaitGroup
+               releaseTriggerLock := make(chan struct{})
+               var releaseTriggerLockOnce sync.Once
+               defer func() {
+                       releaseTriggerLockOnce.Do(func() { 
close(releaseTriggerLock) })
+                       wg.Wait()
+               }()
+
                // Concurrently wait for gpbackup to block when it requests an 
AccessShareLock on public.foo. Once
                // that happens, acquire an AccessExclusiveLock on 
pg_catalog.pg_trigger to block gpbackup during its
                // trigger metadata dump. Then release the initial 
AccessExclusiveLock on public.foo (from the
                // beginning of the test) to unblock gpbackup and let gpbackup 
move forward to the trigger metadata dump.
-               anotherConn := testutils.SetupTestDbConn("testdb")
-               defer anotherConn.Close()
+               wg.Add(1)
                go func() {
+                       defer wg.Done()
+                       triggerLockConn := testutils.SetupTestDbConn("testdb")
+                       defer triggerLockConn.Close()
+
                        // Query to see if gpbackup's AccessShareLock request 
on public.foo is blocked
                        checkLockQuery := `SELECT count(*) FROM pg_locks l, 
pg_class c, pg_namespace n WHERE l.relation = c.oid AND n.oid = c.relnamespace 
AND n.nspname = 'public' AND c.relname = 'foo' AND l.granted = 'f' AND l.mode = 
'AccessShareLock'`
 
@@ -181,7 +230,7 @@ var _ = Describe("Deadlock handling", func() {
                        var gpbackupBlockedLockCount int
                        iterations := 100
                        for iterations > 0 {
-                               _ = anotherConn.Get(&gpbackupBlockedLockCount, 
checkLockQuery)
+                               _ = 
triggerLockConn.Get(&gpbackupBlockedLockCount, checkLockQuery)
                                if gpbackupBlockedLockCount < 1 {
                                        time.Sleep(100 * time.Millisecond)
                                        iterations--
@@ -194,8 +243,12 @@ var _ = Describe("Deadlock handling", func() {
                        // during the trigger metadata dump so that the test 
can queue a bunch of
                        // AccessExclusiveLock requests against the test 
tables. Afterwards, release the
                        // AccessExclusiveLock on public.foo to let gpbackup go 
to the trigger metadata dump.
-                       anotherConn.MustExec(`BEGIN; LOCK TABLE 
pg_catalog.pg_trigger IN ACCESS EXCLUSIVE MODE`)
-                       backupConn.MustExec("COMMIT")
+                       triggerLockConn.MustBegin()
+                       defer func() { _ = triggerLockConn.Rollback() }()
+                       triggerLockConn.MustExec(`LOCK TABLE 
pg_catalog.pg_trigger IN ACCESS EXCLUSIVE MODE`)
+                       initialLockConn.MustCommit()
+                       <-releaseTriggerLock
+                       triggerLockConn.MustCommit()
                }()
 
                // Concurrently wait for gpbackup to block on the trigger 
metadata dump section. Once we
@@ -204,7 +257,9 @@ var _ = Describe("Deadlock handling", func() {
                dataTables := []string{`public."FOObar"`, "public.foo", 
"public.holds", "public.sales", "public.bigtable",
                        "schema2.ao1", "schema2.ao2", "schema2.foo2", 
"schema2.foo3", "schema2.returns"}
                for _, dataTable := range dataTables {
+                       wg.Add(1)
                        go func(dataTable string) {
+                               defer wg.Done()
                                accessExclusiveLockConn := 
testutils.SetupTestDbConn("testdb")
                                defer accessExclusiveLockConn.Close()
 
@@ -226,23 +281,32 @@ var _ = Describe("Deadlock handling", func() {
 
                                // Queue an AccessExclusiveLock request on a 
test table which will later
                                // result in a detected deadlock during the 
gpbackup data dump section.
-                               
accessExclusiveLockConn.MustExec(fmt.Sprintf(`BEGIN; LOCK TABLE %s IN ACCESS 
EXCLUSIVE MODE; COMMIT`, dataTable))
+                               accessExclusiveLockConn.MustBegin()
+                               defer func() { _ = 
accessExclusiveLockConn.Rollback() }()
+                               
accessExclusiveLockConn.MustExec(fmt.Sprintf(`LOCK TABLE %s IN ACCESS EXCLUSIVE 
MODE`, dataTable))
+                               accessExclusiveLockConn.MustCommit()
                        }(dataTable)
                }
 
                // Concurrently wait for all AccessExclusiveLock requests on 
all 10 test tables to block.
                // Once that happens, release the AccessExclusiveLock on 
pg_catalog.pg_trigger to unblock
                // gpbackup and let gpbackup move forward to the data dump 
section.
-               var accessExclBlockedLockCount int
+               accessExclBlockedLockCountChan := make(chan int, 1)
+               wg.Add(1)
                go func() {
+                       defer wg.Done()
+                       accessExclBlockedLockConn := 
testutils.SetupTestDbConn("testdb")
+                       defer accessExclBlockedLockConn.Close()
+
                        // Query to check for ungranted AccessExclusiveLock 
requests on our test tables
                        checkLockQuery := `SELECT count(*) FROM pg_locks WHERE 
granted = 'f' AND mode = 'AccessExclusiveLock'`
 
                        // Wait up to 10 seconds
+                       var accessExclBlockedLockCount int
                        iterations := 100
                        for iterations > 0 {
-                               _ = backupConn.Get(&accessExclBlockedLockCount, 
checkLockQuery)
-                               if accessExclBlockedLockCount < 10 {
+                               _ = 
accessExclBlockedLockConn.Get(&accessExclBlockedLockCount, checkLockQuery)
+                               if accessExclBlockedLockCount < len(dataTables) 
{
                                        time.Sleep(100 * time.Millisecond)
                                        iterations--
                                } else {
@@ -251,15 +315,18 @@ var _ = Describe("Deadlock handling", func() {
                        }
 
                        // Unblock gpbackup by releasing AccessExclusiveLock on 
pg_catalog.pg_trigger
-                       anotherConn.MustExec("COMMIT")
+                       accessExclBlockedLockCountChan <- 
accessExclBlockedLockCount
+                       releaseTriggerLockOnce.Do(func() { 
close(releaseTriggerLock) })
                }()
 
                // gpbackup has finished
-               output, _ := cmd.CombinedOutput()
+               output, err := cmd.CombinedOutput()
                stdout := string(output)
+               accessExclBlockedLockCount := <-accessExclBlockedLockCountChan
+               Expect(err).ToNot(HaveOccurred(), "%s", stdout)
 
                // Check that 10 deadlock traps were placed during the test
-               Expect(accessExclBlockedLockCount).To(Equal(10))
+               Expect(accessExclBlockedLockCount).To(Equal(len(dataTables)))
                // No non-main worker should have been able to run COPY due to 
deadlock detection
                for i := 1; i < 2; i++ {
                        expectedLockString := fmt.Sprintf("[DEBUG]:-Worker %d: 
LOCK TABLE ", i)
diff --git a/end_to_end/signal_handler_test.go 
b/end_to_end/signal_handler_test.go
index e10c399e..95c7a07e 100644
--- a/end_to_end/signal_handler_test.go
+++ b/end_to_end/signal_handler_test.go
@@ -2,15 +2,60 @@ package end_to_end_test
 
 import (
        "math/rand"
+       "os"
        "os/exec"
        "time"
 
+       "github.com/apache/cloudberry-backup/testutils"
        "github.com/apache/cloudberry-go-libs/testhelper"
        . "github.com/onsi/ginkgo/v2"
        . "github.com/onsi/gomega"
        "golang.org/x/sys/unix"
 )
 
+func gpbackupWithBlockedFoo2LockSignal(cmd *exec.Cmd, sig os.Signal, 
checkLockQuery string) ([]byte, int, int) {
+       lockConn := testutils.SetupTestDbConn("testdb")
+       defer lockConn.Close()
+       lockConn.MustBegin()
+       defer func() { _ = lockConn.Rollback() }()
+       lockConn.MustExec("LOCK TABLE schema2.foo2 IN ACCESS EXCLUSIVE MODE")
+
+       beforeLockCountChan := make(chan int, 1)
+       signalDone := make(chan struct{})
+       go func() {
+               defer close(signalDone)
+               lockCheckConn := testutils.SetupTestDbConn("testdb")
+               defer lockCheckConn.Close()
+
+               var beforeLockCount int
+               iterations := 50
+               for iterations > 0 {
+                       _ = lockCheckConn.Get(&beforeLockCount, checkLockQuery)
+                       if beforeLockCount < 1 {
+                               time.Sleep(100 * time.Millisecond)
+                               iterations--
+                       } else {
+                               break
+                       }
+               }
+               beforeLockCountChan <- beforeLockCount
+               if cmd.Process != nil {
+                       _ = cmd.Process.Signal(sig)
+               }
+       }()
+
+       output, _ := cmd.CombinedOutput()
+       <-signalDone
+       beforeLockCount := <-beforeLockCountChan
+
+       afterLockCountConn := testutils.SetupTestDbConn("testdb")
+       defer afterLockCountConn.Close()
+       var afterLockCount int
+       _ = afterLockCountConn.Get(&afterLockCount, checkLockQuery)
+
+       return output, beforeLockCount, afterLockCount
+}
+
 var _ = Describe("Signal handler tests", func() {
        BeforeEach(func() {
                end_to_end_setup()
@@ -90,8 +135,6 @@ var _ = Describe("Signal handler tests", func() {
                        // Query to see if gpbackup lock acquire on 
schema2.foo2 is blocked
                        checkLockQuery := `SELECT count(*) FROM pg_locks l, 
pg_class c, pg_namespace n WHERE l.relation = c.oid AND n.oid = c.relnamespace 
AND n.nspname = 'schema2' AND c.relname = 'foo2' AND l.granted = 'f'`
 
-                       // Acquire AccessExclusiveLock on schema2.foo2 to 
prevent gpbackup from acquiring AccessShareLock
-                       backupConn.MustExec("BEGIN; LOCK TABLE schema2.foo2 IN 
ACCESS EXCLUSIVE MODE")
                        args := []string{
                                "--dbname", "testdb",
                                "--backup-dir", backupDir,
@@ -100,29 +143,12 @@ var _ = Describe("Signal handler tests", func() {
 
                        // Wait up to 5 seconds for gpbackup to block on 
acquiring AccessShareLock.
                        // Once blocked, we send a SIGINT to cancel gpbackup.
-                       var beforeLockCount int
-                       go func() {
-                               iterations := 50
-                               for iterations > 0 {
-                                       _ = backupConn.Get(&beforeLockCount, 
checkLockQuery)
-                                       if beforeLockCount < 1 {
-                                               time.Sleep(100 * 
time.Millisecond)
-                                               iterations--
-                                       } else {
-                                               break
-                                       }
-                               }
-                               _ = cmd.Process.Signal(unix.SIGINT)
-                       }()
-                       output, _ := cmd.CombinedOutput()
+                       output, beforeLockCount, afterLockCount := 
gpbackupWithBlockedFoo2LockSignal(cmd, unix.SIGINT, checkLockQuery)
                        Expect(beforeLockCount).To(Equal(1))
 
                        // After gpbackup has been canceled, we should no 
longer see a blocked SQL
                        // session trying to acquire AccessShareLock on foo2.
-                       var afterLockCount int
-                       _ = backupConn.Get(&afterLockCount, checkLockQuery)
                        Expect(afterLockCount).To(Equal(0))
-                       backupConn.MustExec("ROLLBACK")
 
                        stdout := string(output)
                        Expect(stdout).To(ContainSubstring("Received an 
interrupt signal, aborting backup process"))
@@ -140,8 +166,6 @@ var _ = Describe("Signal handler tests", func() {
                        // Query to see if gpbackup lock acquire on 
schema2.foo2 is blocked
                        checkLockQuery := `SELECT count(*) FROM pg_locks l, 
pg_class c, pg_namespace n WHERE l.relation = c.oid AND n.oid = c.relnamespace 
AND n.nspname = 'schema2' AND c.relname = 'foo2' AND l.granted = 'f'`
 
-                       // Acquire AccessExclusiveLock on schema2.foo2 to 
prevent gpbackup from acquiring AccessShareLock
-                       backupConn.MustExec("BEGIN; LOCK TABLE schema2.foo2 IN 
ACCESS EXCLUSIVE MODE")
                        args := []string{
                                "--dbname", "testdb",
                                "--backup-dir", backupDir,
@@ -151,29 +175,12 @@ var _ = Describe("Signal handler tests", func() {
 
                        // Wait up to 5 seconds for gpbackup to block on 
acquiring AccessShareLock.
                        // Once blocked, we send a SIGINT to cancel gpbackup.
-                       var beforeLockCount int
-                       go func() {
-                               iterations := 50
-                               for iterations > 0 {
-                                       _ = backupConn.Get(&beforeLockCount, 
checkLockQuery)
-                                       if beforeLockCount < 1 {
-                                               time.Sleep(100 * 
time.Millisecond)
-                                               iterations--
-                                       } else {
-                                               break
-                                       }
-                               }
-                               _ = cmd.Process.Signal(unix.SIGINT)
-                       }()
-                       output, _ := cmd.CombinedOutput()
+                       output, beforeLockCount, afterLockCount := 
gpbackupWithBlockedFoo2LockSignal(cmd, unix.SIGINT, checkLockQuery)
                        Expect(beforeLockCount).To(Equal(1))
 
                        // After gpbackup has been canceled, we should no 
longer see a blocked SQL
                        // session trying to acquire AccessShareLock on foo2.
-                       var afterLockCount int
-                       _ = backupConn.Get(&afterLockCount, checkLockQuery)
                        Expect(afterLockCount).To(Equal(0))
-                       backupConn.MustExec("ROLLBACK")
 
                        stdout := string(output)
                        Expect(stdout).To(ContainSubstring("Received an 
interrupt signal, aborting backup process"))
@@ -319,8 +326,6 @@ var _ = Describe("Signal handler tests", func() {
                        // Query to see if gpbackup lock acquire on 
schema2.foo2 is blocked
                        checkLockQuery := `SELECT count(*) FROM pg_locks l, 
pg_class c, pg_namespace n WHERE l.relation = c.oid AND n.oid = c.relnamespace 
AND n.nspname = 'schema2' AND c.relname = 'foo2' AND l.granted = 'f'`
 
-                       // Acquire AccessExclusiveLock on schema2.foo2 to 
prevent gpbackup from acquiring AccessShareLock
-                       backupConn.MustExec("BEGIN; LOCK TABLE schema2.foo2 IN 
ACCESS EXCLUSIVE MODE")
                        args := []string{
                                "--dbname", "testdb",
                                "--backup-dir", backupDir,
@@ -329,29 +334,12 @@ var _ = Describe("Signal handler tests", func() {
 
                        // Wait up to 5 seconds for gpbackup to block on 
acquiring AccessShareLock.
                        // Once blocked, we send a SIGTERM to cancel gpbackup.
-                       var beforeLockCount int
-                       go func() {
-                               iterations := 50
-                               for iterations > 0 {
-                                       _ = backupConn.Get(&beforeLockCount, 
checkLockQuery)
-                                       if beforeLockCount < 1 {
-                                               time.Sleep(100 * 
time.Millisecond)
-                                               iterations--
-                                       } else {
-                                               break
-                                       }
-                               }
-                               _ = cmd.Process.Signal(unix.SIGTERM)
-                       }()
-                       output, _ := cmd.CombinedOutput()
+                       output, beforeLockCount, afterLockCount := 
gpbackupWithBlockedFoo2LockSignal(cmd, unix.SIGTERM, checkLockQuery)
                        Expect(beforeLockCount).To(Equal(1))
 
                        // After gpbackup has been canceled, we should no 
longer see a blocked SQL
                        // session trying to acquire AccessShareLock on foo2.
-                       var afterLockCount int
-                       _ = backupConn.Get(&afterLockCount, checkLockQuery)
                        Expect(afterLockCount).To(Equal(0))
-                       backupConn.MustExec("ROLLBACK")
 
                        stdout := string(output)
                        Expect(stdout).To(ContainSubstring("Received a 
termination signal, aborting backup process"))
@@ -369,8 +357,6 @@ var _ = Describe("Signal handler tests", func() {
                        // Query to see if gpbackup lock acquire on 
schema2.foo2 is blocked
                        checkLockQuery := `SELECT count(*) FROM pg_locks l, 
pg_class c, pg_namespace n WHERE l.relation = c.oid AND n.oid = c.relnamespace 
AND n.nspname = 'schema2' AND c.relname = 'foo2' AND l.granted = 'f'`
 
-                       // Acquire AccessExclusiveLock on schema2.foo2 to 
prevent gpbackup from acquiring AccessShareLock
-                       backupConn.MustExec("BEGIN; LOCK TABLE schema2.foo2 IN 
ACCESS EXCLUSIVE MODE")
                        args := []string{
                                "--dbname", "testdb",
                                "--backup-dir", backupDir,
@@ -380,29 +366,12 @@ var _ = Describe("Signal handler tests", func() {
 
                        // Wait up to 5 seconds for gpbackup to block on 
acquiring AccessShareLock.
                        // Once blocked, we send a SIGTERM to cancel gpbackup.
-                       var beforeLockCount int
-                       go func() {
-                               iterations := 50
-                               for iterations > 0 {
-                                       _ = backupConn.Get(&beforeLockCount, 
checkLockQuery)
-                                       if beforeLockCount < 1 {
-                                               time.Sleep(100 * 
time.Millisecond)
-                                               iterations--
-                                       } else {
-                                               break
-                                       }
-                               }
-                               _ = cmd.Process.Signal(unix.SIGTERM)
-                       }()
-                       output, _ := cmd.CombinedOutput()
+                       output, beforeLockCount, afterLockCount := 
gpbackupWithBlockedFoo2LockSignal(cmd, unix.SIGTERM, checkLockQuery)
                        Expect(beforeLockCount).To(Equal(1))
 
                        // After gpbackup has been canceled, we should no 
longer see a blocked SQL
                        // session trying to acquire AccessShareLock on foo2.
-                       var afterLockCount int
-                       _ = backupConn.Get(&afterLockCount, checkLockQuery)
                        Expect(afterLockCount).To(Equal(0))
-                       backupConn.MustExec("ROLLBACK")
 
                        stdout := string(output)
                        Expect(stdout).To(ContainSubstring("Received a 
termination signal, aborting backup process"))


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to