diff --git a/backup/backup.go b/backup/backup.go index 28123e63..ec225a4d 100644 --- a/backup/backup.go +++ b/backup/backup.go @@ -45,6 +45,7 @@ func DoSetup() { SetLoggerVerbosity() gplog.Verbose("Backup Command: %s", os.Args) gplog.Info("gpbackup version = %s", GetVersion()) + SetSnapshotAttemptsFromEnvironment() utils.CheckGpexpandRunning(utils.BackupPreventedByGpexpandMessage) timestamp := history.CurrentTimestamp() @@ -129,6 +130,7 @@ func DoBackup() { gplog.Info("Gathering table state information") metadataTables, dataTables := RetrieveAndProcessTables() + backupReport.SkippedDataTables = sortedSkippedDataTables() dataTables, numExtOrForeignTables := GetBackupDataSet(dataTables) if len(dataTables) == 0 && !backupReport.MetadataOnly { gplog.Warn("No tables in backup set contain data. Performing metadata-only backup instead.") @@ -194,8 +196,9 @@ func DoBackup() { } printDataBackupWarnings(numExtOrForeignTables) + printSkippedDataTableWarnings() if MustGetFlagBool(options.WITH_STATS) { - backupStatistics(metadataTables) + backupStatistics(tablesWithBackedUpData(metadataTables)) } globalTOC.WriteToFileAndMakeReadOnly(globalFPInfo.GetTOCFilePath()) diff --git a/backup/data.go b/backup/data.go index 1a6d72dc..36195537 100644 --- a/backup/data.go +++ b/backup/data.go @@ -7,6 +7,7 @@ package backup import ( "errors" "fmt" + "sort" "strings" "sync" "sync/atomic" @@ -380,6 +381,38 @@ func printDataBackupWarnings(numExtTables int64) { } } +func printSkippedDataTableWarnings() { + if len(skippedDataTables) > 0 { + gplog.Warn("Data of %d table(s) not backed up because they changed on disk after the backup snapshot was taken.", len(skippedDataTables)) + gplog.Warn("See the backup report or %s for the list of tables.", gplog.GetLogFilePath()) + } +} + +// tablesWithBackedUpData leaves out the relations whose data changed after the +// snapshot or was left out because of that: statistics describe data, and a +// restore creates these tables empty. +func tablesWithBackedUpData(tables []Table) []Table { + if len(changedRelations) == 0 { + return tables + } + kept := make([]Table, 0, len(tables)) + for _, table := range tables { + if !changedRelations[table.FQN()] { + kept = append(kept, table) + } + } + return kept +} + +func sortedSkippedDataTables() []string { + tables := make([]string, 0, len(skippedDataTables)) + for fqn := range skippedDataTables { + tables = append(tables, fqn) + } + sort.Strings(tables) + return tables +} + // Remove external/foreign tables from the data backup set func GetBackupDataSet(tables []Table) ([]Table, int64) { var backupDataSet []Table diff --git a/backup/display_report_test.go b/backup/display_report_test.go index a0763420..3909c1ec 100644 --- a/backup/display_report_test.go +++ b/backup/display_report_test.go @@ -65,6 +65,16 @@ var _ = Describe("display-report internal tests", func() { Expect(backupError).To(Equal("could not dispatch to segment seg0\nconnection refused: server closed")) }) + It("keeps the data-not-backed-up line out of the backup error text", func() { + text := "backup status: Failure\n" + + "backup error: could not dispatch to segment seg0\n\n" + + "data not backed up: public.ao_t, public.heap_t\n\n" + + "count of database objects in backup:\ntables 1\n" + fields, _, backupError := parseReportText(text) + Expect(backupError).To(Equal("could not dispatch to segment seg0")) + Expect(fields).To(HaveKeyWithValue("data_not_backed_up", "public.ao_t, public.heap_t")) + }) + It("folds a colonless continuation line into the most recently seen key", func() { text := "incremental backup set:\n20260101000000\n20260102000000\n\ncount of database objects in backup:\ntables 1\n" fields, objectCounts, _ := parseReportText(text) diff --git a/backup/global_variables.go b/backup/global_variables.go index 06ad8733..755eea3a 100644 --- a/backup/global_variables.go +++ b/backup/global_variables.go @@ -33,6 +33,13 @@ const ( /* * Non-flag variables */ +// SnapshotAttemptsEnvVar names the environment variable that tunes how many +// snapshots a backup may take when tables change under it; see +// SetSnapshotAttemptsFromEnvironment. +const SnapshotAttemptsEnvVar = "WHPGBACKUP_SNAPSHOT_ATTEMPTS" + +const defaultSnapshotAttempts = 3 + var ( backupReport *report.Report connectionPool *dbconn.DBConn @@ -48,6 +55,18 @@ var ( filterRelationClause string quotedRoleNames map[string]string backupSnapshot string + // How many snapshots lockBackupSet may try before it gives up on tables + // that keep changing under it and skips their data. + maxSnapshotAttempts = defaultSnapshotAttempts + // The relations still changing on disk after the last snapshot attempt + // and the locked tables above them, by FQN. Their data cannot be read + // under the snapshot, and a changed relation's pg_aoseg name is stale. + changedRelations map[string]bool + // FQNs of the tables left out of the data backup set because of that. + skippedDataTables map[string]bool + // The filter options the backup started with, kept so the include and + // exclude lists can be resolved again under a new snapshot. + filterOptions *options.Options /* * Used for synchronizing DoCleanup. In DoInit() we increment the group * and then wait for at least one DoCleanup to finish, either in DoTeardown @@ -163,6 +182,18 @@ func SetFilterRelationClause(filterClause string) { filterRelationClause = filterClause } +func SetMaxSnapshotAttempts(attempts int) { + maxSnapshotAttempts = attempts +} + +func GetMaxSnapshotAttempts() int { + return maxSnapshotAttempts +} + +func GetSkippedDataTables() map[string]bool { + return skippedDataTables +} + func SetQuotedRoleNames(quotedRoles map[string]string) { quotedRoleNames = quotedRoles } diff --git a/backup/queries_incremental.go b/backup/queries_incremental.go index 646d151e..7ff5ad44 100644 --- a/backup/queries_incremental.go +++ b/backup/queries_incremental.go @@ -35,8 +35,15 @@ type aoSegTable struct { func getAllModCounts(connectionPool *dbconn.DBConn) map[string]int64 { var segTableFQNs = getAOSegTableFQNs(connectionPool) modCounts := make(map[string]int64) - for aoTableFQN, segTable := range segTableFQNs { - modCounts[aoTableFQN] = getModCount(connectionPool, segTable) + for aoTableFQN, segTableFQN := range segTableFQNs { + // A relation that changed on disk since the snapshot, or a table whose + // data was left out because a partition did, gets no incremental + // entry: its data is not backed up, and a changed relation's segment + // relation name comes from the snapshot and may no longer exist. + if changedRelations[aoTableFQN] { + continue + } + modCounts[aoTableFQN] = getModCount(connectionPool, segTableFQN) } return modCounts } diff --git a/backup/queries_relation_test.go b/backup/queries_relation_test.go index a487b8f9..f0815e91 100644 --- a/backup/queries_relation_test.go +++ b/backup/queries_relation_test.go @@ -61,4 +61,102 @@ var _ = Describe("backup internal tests", func() { structmatcher.ExpectStructsToMatch(&expectedResult[0], &result[0]) }) }) + Describe("GetRelationsChangedSinceSnapshot", func() { + tables := []backup.Relation{ + {SchemaOid: 2200, Oid: 101, Schema: "public", Name: "unchanged"}, + {SchemaOid: 2200, Oid: 102, Schema: "public", Name: "rewritten"}, + {SchemaOid: 2200, Oid: 103, Schema: "public", Name: "dropped"}, + {SchemaOid: 2200, Oid: 104, Schema: "public", Name: "parted"}, + } + header := []string{"oid", "parentoid", "schema", "name", "snapshotrelfilenode", "currentrelfilenode"} + // The partition root has no storage of its own; its data lives in leaf 201. + unchangedRows := func() *sqlmock.Rows { + return sqlmock.NewRows(header). + AddRow(101, nil, "public", "unchanged", 1001, 1001). + AddRow(102, nil, "public", "rewritten", 1002, 1002). + AddRow(103, nil, "public", "dropped", 1003, 1003). + AddRow(104, nil, "public", "parted", 0, nil). + AddRow(201, 104, "public", "parted_1_prt_1", 2001, 2001) + } + + It("returns nothing without querying when there are no tables", func() { + changed := backup.GetRelationsChangedSinceSnapshot(connectionPool, []backup.Relation{}) + Expect(changed).To(BeEmpty()) + Expect(mock.ExpectationsWereMet()).To(Succeed()) + }) + It("returns nothing when every relation still has the relfilenode the snapshot saw", func() { + mock.ExpectQuery(`WITH RECURSIVE (.*)`).WillReturnRows(unchangedRows()) + + changed := backup.GetRelationsChangedSinceSnapshot(connectionPool, tables) + Expect(changed).To(BeEmpty()) + }) + It("reports a table whose live relfilenode differs from the snapshot as rewritten", func() { + rows := sqlmock.NewRows(header). + AddRow(101, nil, "public", "unchanged", 1001, 1001). + AddRow(102, nil, "public", "rewritten", 1002, 2002). + AddRow(103, nil, "public", "dropped", 1003, 1003) + mock.ExpectQuery(`WITH RECURSIVE (.*)`).WillReturnRows(rows) + + changed := backup.GetRelationsChangedSinceSnapshot(connectionPool, tables) + Expect(changed).To(HaveLen(1)) + Expect(changed[0].Relation).To(Equal(tables[1])) + Expect(changed[0].Dropped).To(BeFalse()) + Expect(changed[0].Ancestors).To(BeEmpty()) + }) + It("reports a table with no live relfilenode as dropped", func() { + rows := sqlmock.NewRows(header). + AddRow(101, nil, "public", "unchanged", 1001, 1001). + AddRow(102, nil, "public", "rewritten", 1002, 1002). + AddRow(103, nil, "public", "dropped", 1003, nil) + mock.ExpectQuery(`WITH RECURSIVE (.*)`).WillReturnRows(rows) + + changed := backup.GetRelationsChangedSinceSnapshot(connectionPool, tables) + Expect(changed).To(HaveLen(1)) + Expect(changed[0].Relation).To(Equal(tables[2])) + Expect(changed[0].Dropped).To(BeTrue()) + }) + It("reports a changed partition below a locked table with the table as its ancestor", func() { + rows := sqlmock.NewRows(header). + AddRow(104, nil, "public", "parted", 0, nil). + AddRow(201, 104, "public", "parted_1_prt_1", 2001, 2002) + mock.ExpectQuery(`WITH RECURSIVE (.*)`).WillReturnRows(rows) + + changed := backup.GetRelationsChangedSinceSnapshot(connectionPool, tables) + Expect(changed).To(HaveLen(1)) + Expect(changed[0].FQN()).To(Equal("public.parted_1_prt_1")) + Expect(changed[0].Dropped).To(BeFalse()) + Expect(changed[0].Ancestors).To(Equal([]backup.Relation{tables[3]})) + Expect(changed[0].Describe()).To(Equal("public.parted_1_prt_1 (rewritten or truncated, partition of public.parted)")) + }) + It("does not compare relations without storage", func() { + rows := sqlmock.NewRows(header). + AddRow(104, nil, "public", "parted", 0, nil) + mock.ExpectQuery(`WITH RECURSIVE (.*)`).WillReturnRows(rows) + + changed := backup.GetRelationsChangedSinceSnapshot(connectionPool, tables) + Expect(changed).To(BeEmpty()) + }) + It("describes a changed table with its reason", func() { + rewritten := backup.ChangedRelation{Relation: tables[1], Dropped: false} + dropped := backup.ChangedRelation{Relation: tables[2], Dropped: true} + Expect(rewritten.Reason()).To(Equal("rewritten or truncated")) + Expect(rewritten.Describe()).To(Equal("public.rewritten (rewritten or truncated)")) + Expect(dropped.Reason()).To(Equal("dropped and recreated")) + Expect(dropped.Describe()).To(Equal("public.dropped (dropped and recreated)")) + }) + It("lists locked tables in input order, each followed by its changed partitions", func() { + rows := sqlmock.NewRows(header). + AddRow(201, 104, "public", "parted_1_prt_1", 2001, 2002). + AddRow(104, nil, "public", "parted", 0, nil). + AddRow(103, nil, "public", "dropped", 1003, nil). + AddRow(102, nil, "public", "rewritten", 1002, 2002) + mock.ExpectQuery(`WITH RECURSIVE (.*)`).WillReturnRows(rows) + + changed := backup.GetRelationsChangedSinceSnapshot(connectionPool, tables) + Expect(changed).To(HaveLen(3)) + Expect(changed[0].Name).To(Equal("rewritten")) + Expect(changed[1].Name).To(Equal("dropped")) + Expect(changed[2].Name).To(Equal("parted_1_prt_1")) + }) + }) }) diff --git a/backup/queries_relations.go b/backup/queries_relations.go index bed1ae2e..efa9e91b 100644 --- a/backup/queries_relations.go +++ b/backup/queries_relations.go @@ -8,6 +8,7 @@ package backup import ( "database/sql" "fmt" + "sort" "strconv" "strings" @@ -518,6 +519,154 @@ func LockTables(connectionPool *dbconn.DBConn, tables []Relation) { progressBar.Finish() } +// ChangedRelation is a relation whose storage no longer matches what the +// backup snapshot describes: it was rewritten, truncated, or dropped (and +// possibly recreated under the same name) after the snapshot was taken. It is +// either one of the locked tables or a partition below one of them; Ancestors +// lists the inheritance chain up to the locked table, nearest first. +type ChangedRelation struct { + Relation + Dropped bool + Ancestors []Relation +} + +// Reason says what happened to the relation, for warnings. +func (c ChangedRelation) Reason() string { + if c.Dropped { + return "dropped and recreated" + } + return "rewritten or truncated" +} + +// Describe names the relation, its Reason, and the table it is a partition of. +func (c ChangedRelation) Describe() string { + if len(c.Ancestors) > 0 { + root := c.Ancestors[len(c.Ancestors)-1] + return fmt.Sprintf("%s (%s, partition of %s)", c.FQN(), c.Reason(), root.FQN()) + } + return fmt.Sprintf("%s (%s)", c.FQN(), c.Reason()) +} + +/* + * GetRelationsChangedSinceSnapshot reports which of the given locked tables, + * or of the partitions below them, changed on disk between the backup + * snapshot and the locks now held. + * + * The catalog rows come from the backup snapshot, while pg_relation_filenode() + * looks the relation up in the live catalog, the same way LOCK TABLE and the + * later COPY do. A table rewritten by ALTER TABLE ... SET WITH (REORGANIZE=true) + * or SET DISTRIBUTED BY, a heap table rewritten by VACUUM FULL, or a truncated + * table has a new relfilenode; a table dropped (and possibly recreated under + * the same name) has none. Either way the snapshot can no longer see the + * table's data: the old files are gone, and for append-optimized tables the + * pg_aoseg relation named under the snapshot no longer exists. + * + * Partitions are followed through pg_inherits because a backup that does not + * use --leaf-partition-data locks and copies the parent only, while the data + * lives in the children; LOCK TABLE on the parent locks them too, so the + * check is race-free for them as well. Relations without storage (relfilenode + * 0, e.g. partition roots) are not compared. Locked tables come back in input + * order, each followed by its changed partitions. + */ +func GetRelationsChangedSinceSnapshot(connectionPool *dbconn.DBConn, tables []Relation) []ChangedRelation { + if len(tables) == 0 { + return nil + } + type relationStorage struct { + Oid uint32 + ParentOid sql.NullInt64 + Schema string + Name string + SnapshotRelfilenode uint32 + CurrentRelfilenode sql.NullInt64 + } + oids := make([]string, 0, len(tables)) + for _, table := range tables { + oids = append(oids, strconv.FormatUint(uint64(table.Oid), 10)) + } + query := fmt.Sprintf(` + WITH RECURSIVE inheritance(oid, parentoid) AS ( + SELECT c.oid, NULL::oid + FROM pg_class c + WHERE c.oid IN (%s) + UNION ALL + SELECT i.inhrelid, i.inhparent + FROM pg_inherits i + JOIN inheritance h ON i.inhparent = h.oid + ) + SELECT DISTINCT h.oid, + h.parentoid, + quote_ident(n.nspname) AS schema, + quote_ident(c.relname) AS name, + c.relfilenode AS snapshotrelfilenode, + pg_catalog.pg_relation_filenode(c.oid) AS currentrelfilenode + FROM inheritance h + JOIN pg_class c ON c.oid = h.oid + JOIN pg_namespace n ON n.oid = c.relnamespace`, strings.Join(oids, ", ")) + rows := make([]relationStorage, 0) + err := connectionPool.Select(&rows, query) + gplog.FatalOnError(err) + + locked := make(map[uint32]int, len(tables)) + for i, table := range tables { + locked[table.Oid] = i + } + relations := make(map[uint32]Relation) + parents := make(map[uint32]uint32) + dropped := make(map[uint32]bool) + for _, row := range rows { + if _, isLocked := locked[row.Oid]; isLocked { + relations[row.Oid] = tables[locked[row.Oid]] + } else { + relations[row.Oid] = Relation{Oid: row.Oid, Schema: row.Schema, Name: row.Name} + } + // A partition that is itself locked appears both on its own and + // below its parent; keep the parent link. + if row.ParentOid.Valid { + parents[row.Oid] = uint32(row.ParentOid.Int64) + } + if row.SnapshotRelfilenode == 0 { + continue + } + if !row.CurrentRelfilenode.Valid { + dropped[row.Oid] = true + } else if uint32(row.CurrentRelfilenode.Int64) != row.SnapshotRelfilenode { + dropped[row.Oid] = false + } + } + + ancestorsOf := func(oid uint32) []Relation { + ancestors := make([]Relation, 0) + for parent, ok := parents[oid]; ok; parent, ok = parents[parent] { + ancestors = append(ancestors, relations[parent]) + } + return ancestors + } + rootIndex := func(oid uint32) int { + root := oid + for parent, ok := parents[root]; ok; parent, ok = parents[root] { + root = parent + } + return locked[root] + } + + changed := make([]ChangedRelation, 0) + for oid, isDropped := range dropped { + changed = append(changed, ChangedRelation{Relation: relations[oid], Dropped: isDropped, Ancestors: ancestorsOf(oid)}) + } + sort.Slice(changed, func(i, j int) bool { + ri, rj := rootIndex(changed[i].Oid), rootIndex(changed[j].Oid) + if ri != rj { + return ri < rj + } + if len(changed[i].Ancestors) != len(changed[j].Ancestors) { + return len(changed[i].Ancestors) < len(changed[j].Ancestors) + } + return changed[i].Oid < changed[j].Oid + }) + return changed +} + // GenerateTableBatches batches tables to reduce network congestion and // resource contention. Returns an array of batches where a batch of tables is // a single string with comma separated tables diff --git a/backup/queries_statistics.go b/backup/queries_statistics.go index bc19625a..8919b5b1 100644 --- a/backup/queries_statistics.go +++ b/backup/queries_statistics.go @@ -53,6 +53,9 @@ type AttributeStatistic struct { } func GetAttributeStatistics(connectionPool *dbconn.DBConn, tables []Table) map[uint32][]AttributeStatistic { + if len(tables) == 0 { + return map[uint32][]AttributeStatistic{} + } inheritClause := "" statSlotClause := "" if connectionPool.Version.AtLeast("6") { @@ -137,6 +140,9 @@ type TupleStatistic struct { } func GetTupleStatistics(connectionPool *dbconn.DBConn, tables []Table) map[uint32]TupleStatistic { + if len(tables) == 0 { + return map[uint32]TupleStatistic{} + } tablenames := make([]string, 0) for _, table := range tables { tablenames = append(tablenames, table.FQN()) diff --git a/backup/queries_statistics_test.go b/backup/queries_statistics_test.go new file mode 100644 index 00000000..3f0c089b --- /dev/null +++ b/backup/queries_statistics_test.go @@ -0,0 +1,26 @@ +package backup_test + +import ( + "github.com/greenplum-db/gpbackup/backup" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +var _ = Describe("backup/queries_statistics", func() { + Describe("statistics query builders with an empty table set", func() { + // A backup that skips every table's data (all of them kept changing) + // passes no tables here. The builders must return nothing without + // running a query, since an empty IN () list is a syntax error. + It("GetAttributeStatistics returns no statistics and runs no query", func() { + result := backup.GetAttributeStatistics(connectionPool, []backup.Table{}) + Expect(result).To(BeEmpty()) + Expect(mock.ExpectationsWereMet()).To(Succeed()) + }) + It("GetTupleStatistics returns no statistics and runs no query", func() { + result := backup.GetTupleStatistics(connectionPool, []backup.Table{}) + Expect(result).To(BeEmpty()) + Expect(mock.ExpectationsWereMet()).To(Succeed()) + }) + }) +}) diff --git a/backup/validate.go b/backup/validate.go index b368576f..39de4cc8 100644 --- a/backup/validate.go +++ b/backup/validate.go @@ -2,6 +2,7 @@ package backup import ( "fmt" + "strconv" "github.com/greenplum-db/gpbackup/filepath" "github.com/greenplum-db/gpbackup/history" @@ -11,6 +12,7 @@ import ( "github.com/spf13/pflag" "github.com/warehouse-pg/common-go-libs/dbconn" "github.com/warehouse-pg/common-go-libs/gplog" + "github.com/warehouse-pg/common-go-libs/operating" ) /* @@ -22,6 +24,7 @@ import ( */ func ValidateAndProcessFilterLists(opts *options.Options) { gplog.Verbose("Validating Tables and Schemas exist in Database") + filterOptions = opts // pre-create these so we can save the processed versions of our filters IncludedRelationFqns = make([]options.Relation, 0) @@ -216,3 +219,25 @@ func validateFromTimestamp(fromTimestamp string) { "previous backup.", fromTimestampFPInfo.Timestamp), "") } } + +/* + * SetSnapshotAttemptsFromEnvironment reads how many snapshots lockBackupSet + * may take before it gives up on tables that keep changing. The value comes + * from the WHPGBACKUP_SNAPSHOT_ATTEMPTS environment variable, so a site can + * tune it without a flag: unset or empty keeps the default, 1 means no retry, + * and anything that is not a whole number of at least 1 stops the backup + * before it starts. + */ +func SetSnapshotAttemptsFromEnvironment() { + value := operating.System.Getenv(SnapshotAttemptsEnvVar) + if value == "" { + maxSnapshotAttempts = defaultSnapshotAttempts + return + } + attempts, err := strconv.Atoi(value) + if err != nil || attempts < 1 { + gplog.Fatal(errors.Errorf(`%s must be a whole number of at least 1, got "%s"`, SnapshotAttemptsEnvVar, value), "") + } + maxSnapshotAttempts = attempts + gplog.Verbose("Snapshot attempts set to %d by %s", attempts, SnapshotAttemptsEnvVar) +} diff --git a/backup/validate_test.go b/backup/validate_test.go index acde8325..fddc8d81 100644 --- a/backup/validate_test.go +++ b/backup/validate_test.go @@ -1,15 +1,18 @@ package backup_test import ( + "fmt" "strings" "github.com/DATA-DOG/go-sqlmock" "github.com/greenplum-db/gpbackup/backup" "github.com/greenplum-db/gpbackup/options" "github.com/spf13/cobra" + "github.com/warehouse-pg/common-go-libs/operating" "github.com/warehouse-pg/common-go-libs/testhelper" . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" ) var _ = Describe("backup/validate tests", func() { @@ -287,3 +290,57 @@ var _ = Describe("backup/validate tests", func() { ) }) }) + +var _ = Describe("SetSnapshotAttemptsFromEnvironment", func() { + var value string + BeforeEach(func() { + operating.System.Getenv = func(key string) string { + if key == backup.SnapshotAttemptsEnvVar { + return value + } + return "" + } + }) + AfterEach(func() { + operating.System = operating.InitializeSystemFunctions() + backup.SetMaxSnapshotAttempts(3) + }) + It("keeps the default when the variable is unset", func() { + value = "" + backup.SetMaxSnapshotAttempts(7) + backup.SetSnapshotAttemptsFromEnvironment() + Expect(backup.GetMaxSnapshotAttempts()).To(Equal(3)) + }) + It("uses the number of attempts given", func() { + value = "5" + backup.SetSnapshotAttemptsFromEnvironment() + Expect(backup.GetMaxSnapshotAttempts()).To(Equal(5)) + }) + It("accepts 1, which disables the retry", func() { + value = "1" + backup.SetSnapshotAttemptsFromEnvironment() + Expect(backup.GetMaxSnapshotAttempts()).To(Equal(1)) + }) + DescribeTable("rejects anything that is not a whole number of at least 1", + func(given string) { + value = given + defer testhelper.ShouldPanicWithMessage( + fmt.Sprintf(`%s must be a whole number of at least 1, got "%s"`, backup.SnapshotAttemptsEnvVar, given)) + backup.SetSnapshotAttemptsFromEnvironment() + }, + Entry("zero", "0"), + Entry("negative", "-2"), + Entry("not a number", "abc"), + Entry("fraction", "2.5"), + Entry("surrounding space", " 4"), + ) + It("leaves the attempt count alone when it rejects a value", func() { + value = "0" + backup.SetMaxSnapshotAttempts(3) + defer func() { + _ = recover() + Expect(backup.GetMaxSnapshotAttempts()).To(Equal(3)) + }() + backup.SetSnapshotAttemptsFromEnvironment() + }) +}) diff --git a/backup/wrappers.go b/backup/wrappers.go index dbd37668..b9cc05ba 100644 --- a/backup/wrappers.go +++ b/backup/wrappers.go @@ -14,6 +14,7 @@ import ( "github.com/greenplum-db/gpbackup/utils" "github.com/nightlyone/lockfile" "github.com/pkg/errors" + "github.com/spf13/pflag" "github.com/warehouse-pg/common-go-libs/cluster" "github.com/warehouse-pg/common-go-libs/dbconn" "github.com/warehouse-pg/common-go-libs/gplog" @@ -207,14 +208,7 @@ func createBackupDirectoriesOnAllHosts() { */ func RetrieveAndProcessTables() ([]Table, []Table) { - includedRelations := GetIncludedUserTableRelations(connectionPool, IncludedRelationFqns) - tableRelations := ConvertRelationsOptionsToBackup(includedRelations) - - // Query extension config dump tables early so we can lock them and - // include them in the single ConstructDefinitionsForTables call. - configDumpRelations, configDumpFilterConds := GetExtensionConfigDumpRelations(connectionPool) - tableRelations = append(tableRelations, configDumpRelations...) - LockTables(connectionPool, tableRelations) + tableRelations, configDumpRelations, configDumpFilterConds := lockBackupSet() if connectionPool.Version.AtLeast("6") { tableRelations = append(tableRelations, GetForeignTableRelations(connectionPool)...) @@ -245,11 +239,153 @@ func RetrieveAndProcessTables() ([]Table, []Table) { metadataTables, dataTables := SplitTablesByPartitionType(regularTables, IncludedRelationFqns) dataTables = append(dataTables, configDumpTables...) + dataTables = removeChangedDataTables(dataTables) objectCounts["Tables"] = len(metadataTables) return metadataTables, dataTables } +/* + * lockBackupSet resolves the tables in the backup set and takes ACCESS SHARE + * locks on them. The catalog is read under the backup snapshot, but LOCK TABLE + * resolves names against the live catalog, so a table rewritten (ALTER TABLE + * ... SET WITH (REORGANIZE=true), SET DISTRIBUTED BY, VACUUM FULL of a heap + * table), truncated, or dropped and recreated between the two can no longer + * be read consistently under the snapshot: its old storage is gone. Nothing + * can change once the locks are held, so the set, partitions included, is + * verified after locking. If anything changed, the snapshot and the + * resolution are redone, up to maxSnapshotAttempts times. Relations still + * changing after the last attempt are recorded so that their data, and the + * data of the locked tables above them, is left out of the backup with a + * warning; whatever metadata the backup carries for them is unaffected. + */ +func lockBackupSet() ([]Relation, []Relation, map[uint32]string) { + changedRelations = make(map[string]bool) + skippedDataTables = make(map[string]bool) + for attempt := 1; ; attempt++ { + includedRelations := GetIncludedUserTableRelations(connectionPool, IncludedRelationFqns) + tableRelations := ConvertRelationsOptionsToBackup(includedRelations) + + // Query extension config dump tables early so we can lock them and + // include them in the single ConstructDefinitionsForTables call. + configDumpRelations, configDumpFilterConds := GetExtensionConfigDumpRelations(connectionPool) + tableRelations = append(tableRelations, configDumpRelations...) + LockTables(connectionPool, tableRelations) + + // A metadata-only backup reads nothing but the catalog, which the + // snapshot keeps consistent whatever happens to the tables' storage. + // The storage check compares the snapshot's relfilenode with the live + // one, so it is meaningful only when a synchronized snapshot is in use; + // older servers take none, and pg_relation_filenode does not exist + // before GPDB 6, so the check is skipped there. + if MustGetFlagBool(options.METADATA_ONLY) || connectionPool.Version.Before(SNAPSHOT_GPDB_MIN_VERSION) { + return tableRelations, configDumpRelations, configDumpFilterConds + } + changed := GetRelationsChangedSinceSnapshot(connectionPool, tableRelations) + if len(changed) == 0 { + return tableRelations, configDumpRelations, configDumpFilterConds + } + + if attempt >= maxSnapshotAttempts { + for _, relation := range changed { + gplog.Warn("Table %s changed on disk after the backup snapshot was taken and was still changing "+ + "after %d attempt(s); its data will not be backed up.", relation.Describe(), attempt) + changedRelations[relation.FQN()] = true + for _, ancestor := range relation.Ancestors { + changedRelations[ancestor.FQN()] = true + } + } + return tableRelations, configDumpRelations, configDumpFilterConds + } + + descriptions := make([]string, len(changed)) + for i, relation := range changed { + descriptions[i] = relation.Describe() + } + gplog.Warn("Table(s) %s changed on disk after the backup snapshot was taken; taking a new snapshot (attempt %d of %d)", + strings.Join(descriptions, ", "), attempt+1, maxSnapshotAttempts) + restartBackupTransactions() + } +} + +/* + * restartBackupTransactions gives every open backup connection a new + * transaction on a fresh synchronized snapshot and resolves the include and + * exclude lists again under it. Nothing has been written when this runs. + * + * The session GUCs are set inside the transaction at connection setup, so the + * rollback discards them; they are set again the same way, in the same order. + */ +func restartBackupTransactions() { + backupSnapshot = "" + for connNum := 0; connNum < connectionPool.NumConns; connNum++ { + if connectionPool.Tx[connNum] == nil { + continue + } + // The old transaction is discarded either way; a rollback error only + // means its connection is already gone, and Begin reconnects. + err := connectionPool.Rollback(connNum) + if err != nil { + gplog.Warn("Connection %d: %s", connNum, err) + } + connectionPool.MustBegin(connNum) + if connectionPool.Version.AtLeast(SNAPSHOT_GPDB_MIN_VERSION) { + if connNum == 0 { + snapshot, err := GetSynchronizedSnapshot(connectionPool) + gplog.FatalOnError(err) + backupSnapshot = snapshot + } else if backupSnapshot != "" { + err := SetSynchronizedSnapshot(connectionPool, connNum, backupSnapshot) + gplog.FatalOnError(err) + } + } + SetSessionGUCs(connNum) + } + + refreshFilterLists() +} + +/* + * refreshFilterLists resolves the include and exclude lists again under the + * current snapshot by repeating what the backup did when it started, from the + * names the user gave: validate them, then expand the include list to the + * partitions and dependent tables found now. A table dropped and recreated + * under the same name is found by its new OID, and a partition created in the + * window is picked up. As at the start, an included table that no longer + * exists is fatal. + */ +func refreshFilterLists() { + filterOptions.ResetIncludedRelations() + if flag := cmdFlags.Lookup(options.INCLUDE_RELATION); flag != nil { + err := flag.Value.(pflag.SliceValue).Replace(filterOptions.GetOriginalIncludedTables()) + gplog.FatalOnError(err) + } + ValidateAndProcessFilterLists(filterOptions) + includeOids := GetOidsFromRelationList(IncludedRelationFqns) + err := ExpandIncludesForPartitions(connectionPool, filterOptions, includeOids, cmdFlags) + gplog.FatalOnError(err) + SetFilterRelationClause("") +} + +// removeChangedDataTables takes the tables that changed on disk since the +// snapshot, or that have a partition that did, out of the data backup set. A +// table is removed here only; nothing else about it is touched. +func removeChangedDataTables(dataTables []Table) []Table { + if len(changedRelations) == 0 { + return dataTables + } + kept := make([]Table, 0, len(dataTables)) + for _, table := range dataTables { + if changedRelations[table.FQN()] { + gplog.Warn("Data for table %s not backed up.", table.FQN()) + skippedDataTables[table.FQN()] = true + continue + } + kept = append(kept, table) + } + return kept +} + func retrieveFunctions(sortables *[]Sortable, metadataMap MetadataMap) ([]Function, map[uint32]FunctionInfo) { gplog.Verbose("Retrieving function information") functionMetadata := GetMetadataForObjectType(connectionPool, TYPE_FUNCTION) diff --git a/end_to_end/concurrent_ddl_test.go b/end_to_end/concurrent_ddl_test.go new file mode 100644 index 00000000..2518b2de --- /dev/null +++ b/end_to_end/concurrent_ddl_test.go @@ -0,0 +1,472 @@ +package end_to_end_test + +import ( + "fmt" + "os" + "os/exec" + "time" + + "github.com/greenplum-db/gpbackup/testutils" + "github.com/greenplum-db/gpbackup/toc" + "github.com/warehouse-pg/common-go-libs/dbconn" + "github.com/warehouse-pg/common-go-libs/testhelper" + "gopkg.in/yaml.v2" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +/* + * gpbackup exports its snapshot before it locks the tables. A table whose + * storage is replaced in that window (ALTER TABLE ... SET WITH (REORGANIZE=true), + * SET DISTRIBUTED BY, TRUNCATE, VACUUM FULL of a heap table, or DROP + CREATE + * under the same name) can no longer be read consistently under that snapshot: its old files + * are gone, so COPY sees no rows, and for AO tables the pg_aoseg helper named + * under the snapshot no longer exists. + * + * These tests drive that window deterministically: a second session holds the + * rewrite open, gpbackup blocks in LockTables behind it, the session commits, + * and gpbackup must notice the change after the lock is granted. + */ + +const ( + changedTableWarning = "changed on disk after the backup snapshot was taken" + newSnapshotWarning = "taking a new snapshot" + dataSkippedWarning = "not backed up" + // The environment variable that tunes how many snapshots gpbackup takes. + snapshotAttemptsEnvVar = "WHPGBACKUP_SNAPSHOT_ATTEMPTS" +) + +// holdRewrite opens a transaction on a fresh connection, runs the given +// statements, and leaves the transaction open so its locks stay held. +func holdRewrite(statements string) *dbconn.DBConn { + conn := testutils.SetupTestDbConn("testdb") + conn.MustBegin() + conn.MustExec(statements) + return conn +} + +// releaseRewrite rolls back a rewrite that is still open (a failed assertion +// may have skipped the commit) so a blocked gpbackup can finish, then closes +// the connection. +func releaseRewrite(conn *dbconn.DBConn) { + if conn.Tx[0] != nil { + _ = conn.Rollback() + } + conn.Close() +} + +// startBackup runs gpbackup in the background with the given extra flags and +// returns a channel that is closed when it exits, plus pointers to its output +// and error. +func startBackup(extraArgs ...string) (<-chan struct{}, *string, *error) { + return startBackupEnv(nil, extraArgs...) +} + +// startBackupEnv is startBackup with extra environment variables (KEY=value). +func startBackupEnv(env []string, extraArgs ...string) (<-chan struct{}, *string, *error) { + done := make(chan struct{}) + var output string + var err error + go func() { + defer GinkgoRecover() + output, err = runBackupEnv(env, extraArgs...) + close(done) + }() + return done, &output, &err +} + +// waitForBackupExit blocks until the background gpbackup exits, so a failed +// test never leaves a gpbackup process behind holding locks on testdb. +func waitForBackupExit(done <-chan struct{}) { + select { + case <-done: + case <-time.After(5 * time.Minute): + Fail("gpbackup did not exit within 5 minutes") + } +} + +// waitForBackupLockWait waits until gpbackup's LOCK TABLE is blocked behind +// the transaction holding a rewrite. gpbackup tags its sessions with +// application_name gpbackup_, and a blocked LOCK TABLE shows up as +// an ungranted AccessShareLock request of one of those sessions. +func waitForBackupLockWait(conn *dbconn.DBConn) { + query := `SELECT count(*) FROM pg_locks l JOIN pg_stat_activity a ON l.pid = a.pid + WHERE a.application_name LIKE 'gpbackup_%' AND l.granted = 'f' AND l.mode = 'AccessShareLock'` + var waiters int + for i := 0; i < 300; i++ { + _ = conn.Get(&waiters, query) + if waiters > 0 { + return + } + time.Sleep(100 * time.Millisecond) + } + Fail("gpbackup did not block in LockTables within 30s") +} + +// waitForLockWaiter waits until some session is queued for a lock of the +// given mode on the table, e.g. a rewrite queued behind the current holder. +func waitForLockWaiter(conn *dbconn.DBConn, schema string, table string, mode string) { + query := fmt.Sprintf(`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 = '%s' AND c.relname = '%s' AND l.granted = 'f' AND l.mode = '%s'`, + schema, table, mode) + var waiters int + for i := 0; i < 300; i++ { + _ = conn.Get(&waiters, query) + if waiters > 0 { + return + } + time.Sleep(100 * time.Millisecond) + } + Fail(fmt.Sprintf("no session queued for %s on %s.%s within 30s", mode, schema, table)) +} + +// runBackup runs gpbackup with the given extra flags and returns its combined +// output and exit error. +func runBackup(extraArgs ...string) (string, error) { + return runBackupEnv(nil, extraArgs...) +} + +// runBackupEnv is runBackup with extra environment variables (KEY=value). +func runBackupEnv(env []string, extraArgs ...string) (string, error) { + args := append([]string{"--verbose", "--dbname", "testdb", "--backup-dir", backupDir}, extraArgs...) + cmd := exec.Command(gpbackupPath, args...) + cmd.Env = append(os.Environ(), env...) + output, err := cmd.CombinedOutput() + return string(output), err +} + +// runLeafPartitionBackup runs gpbackup with --leaf-partition-data, so the AO +// modcount queries run. +func runLeafPartitionBackup() (string, error) { + return runBackup("--leaf-partition-data") +} + +func readTOC(timestamp string) *toc.TOC { + tocStruct := &toc.TOC{} + err := yaml.Unmarshal(getMetdataFileContents(backupDir, timestamp, "toc.yaml"), tocStruct) + Expect(err).ToNot(HaveOccurred()) + return tocStruct +} + +func findDataEntry(tocStruct *toc.TOC, schema string, name string) *toc.CoordinatorDataEntry { + for i := range tocStruct.DataEntries { + entry := &tocStruct.DataEntries[i] + if entry.Schema == schema && entry.Name == name { + return entry + } + } + return nil +} + +func hasTableDDL(tocStruct *toc.TOC, schema string, name string) bool { + for _, entry := range tocStruct.PredataEntries { + if entry.ObjectType == toc.OBJ_TABLE && entry.Schema == schema && entry.Name == name { + return true + } + } + return false +} + +var _ = Describe("Tables changed between snapshot and lock", func() { + BeforeEach(func() { + if useOldBackupVersion { + Skip("This test is not needed for old backup versions") + } + end_to_end_setup() + }) + AfterEach(func() { + end_to_end_teardown() + }) + + It("backs up normally when no concurrent DDL happens", func() { + output, err := runLeafPartitionBackup() + Expect(err).ToNot(HaveOccurred(), output) + Expect(output).To(ContainSubstring("Backup completed successfully")) + Expect(output).ToNot(ContainSubstring(changedTableWarning)) + Expect(output).ToNot(ContainSubstring("[CRITICAL]")) + + tocStruct := readTOC(getBackupTimestamp(output)) + Expect(findDataEntry(tocStruct, "schema2", "ao1").RowsCopied).To(Equal(int64(1000))) + Expect(tocStruct.IncrementalMetadata.AO).To(HaveKey("schema2.ao1")) + }) + + It("retries with a new snapshot when AO and heap tables are rewritten in the window", func() { + // AO table: the rewrite replaces its pg_aoseg helper, which is what the + // customer hit. Heap table: the rewrite replaces its relfilenode, which + // today produces an empty table in the backup with no error at all. + rewriter := holdRewrite(`ALTER TABLE schema2.ao1 SET WITH (reorganize=true); + ALTER TABLE public.foo SET WITH (reorganize=true)`) + done, output, err := startBackup("--leaf-partition-data") + defer waitForBackupExit(done) + defer releaseRewrite(rewriter) + + waitForBackupLockWait(backupConn) + rewriter.MustCommit() + waitForBackupExit(done) + + Expect(*err).ToNot(HaveOccurred(), *output) + Expect(*output).To(ContainSubstring("Backup completed successfully")) + Expect(*output).ToNot(ContainSubstring("[CRITICAL]")) + Expect(*output).To(ContainSubstring(changedTableWarning)) + Expect(*output).To(ContainSubstring("schema2.ao1")) + Expect(*output).To(ContainSubstring("public.foo")) + Expect(*output).To(ContainSubstring(newSnapshotWarning)) + Expect(*output).ToNot(ContainSubstring(dataSkippedWarning)) + + // The second attempt saw the rewritten storage, so both tables carry + // their full data and the AO table keeps its incremental entry. + timestamp := getBackupTimestamp(*output) + tocStruct := readTOC(timestamp) + Expect(findDataEntry(tocStruct, "schema2", "ao1").RowsCopied).To(Equal(int64(1000))) + Expect(findDataEntry(tocStruct, "public", "foo").RowsCopied).To(Equal(int64(40000))) + Expect(tocStruct.IncrementalMetadata.AO).To(HaveKey("schema2.ao1")) + + // The metadata written after the retry must restore like any other. + gprestore(gprestorePath, restoreHelperPath, timestamp, + "--redirect-db", "restoredb", "--backup-dir", backupDir) + assertDataRestored(restoreConn, map[string]int{ + "schema2.ao1": 1000, + "public.foo": 40000, + }) + }) + + It("retries when a table is truncated and reloaded in the window", func() { + // Without detection the COPY under the stale snapshot would return no + // rows, silently backing up an empty table. + rewriter := holdRewrite(`TRUNCATE schema2.ao2; INSERT INTO schema2.ao2 SELECT generate_series(1, 5)`) + defer testhelper.AssertQueryRuns(backupConn, + "TRUNCATE schema2.ao2; INSERT INTO schema2.ao2 SELECT generate_series(1, 1000)") + done, output, err := startBackup("--leaf-partition-data") + defer waitForBackupExit(done) + defer releaseRewrite(rewriter) + + waitForBackupLockWait(backupConn) + rewriter.MustCommit() + waitForBackupExit(done) + + Expect(*err).ToNot(HaveOccurred(), *output) + Expect(*output).To(ContainSubstring("Backup completed successfully")) + Expect(*output).To(ContainSubstring(changedTableWarning)) + Expect(*output).To(ContainSubstring("schema2.ao2")) + + timestamp := getBackupTimestamp(*output) + tocStruct := readTOC(timestamp) + Expect(findDataEntry(tocStruct, "schema2", "ao2").RowsCopied).To(Equal(int64(5))) + + gprestore(gprestorePath, restoreHelperPath, timestamp, + "--redirect-db", "restoredb", "--backup-dir", backupDir) + assertDataRestored(restoreConn, map[string]int{"schema2.ao2": 5}) + }) + + It("retries when a table is dropped and recreated under the same name in the window", func() { + // LOCK TABLE by name lands on the new table while the snapshot still + // enumerates the old OID. The retry re-resolves the include list. + rewriter := holdRewrite(`DROP TABLE schema2.ao1; + CREATE TABLE schema2.ao1 (i integer) WITH (appendonly=true); + INSERT INTO schema2.ao1 SELECT generate_series(1, 7)`) + defer testhelper.AssertQueryRuns(backupConn, + "TRUNCATE schema2.ao1; INSERT INTO schema2.ao1 SELECT generate_series(1, 1000)") + done, output, err := startBackup("--leaf-partition-data") + defer waitForBackupExit(done) + defer releaseRewrite(rewriter) + + waitForBackupLockWait(backupConn) + rewriter.MustCommit() + waitForBackupExit(done) + + Expect(*err).ToNot(HaveOccurred(), *output) + Expect(*output).To(ContainSubstring("Backup completed successfully")) + Expect(*output).To(ContainSubstring(changedTableWarning)) + Expect(*output).To(ContainSubstring("schema2.ao1")) + + timestamp := getBackupTimestamp(*output) + tocStruct := readTOC(timestamp) + Expect(findDataEntry(tocStruct, "schema2", "ao1").RowsCopied).To(Equal(int64(7))) + Expect(tocStruct.IncrementalMetadata.AO).To(HaveKey("schema2.ao1")) + + gprestore(gprestorePath, restoreHelperPath, timestamp, + "--redirect-db", "restoredb", "--backup-dir", backupDir) + assertDataRestored(restoreConn, map[string]int{"schema2.ao1": 7}) + }) + + It("expands an included partitioned table again when it is dropped and recreated with more partitions in the window", func() { + // The include list is expanded to the partitions when the backup + // starts. The retry must expand it again, or the partition created in + // the window has no DDL in the backup and the restore fails on it. + testhelper.AssertQueryRuns(backupConn, `CREATE TABLE schema2.relayout (id int, d int) DISTRIBUTED BY (id) + PARTITION BY RANGE (d) (START (1) END (3) EVERY (1)); + INSERT INTO schema2.relayout SELECT g, (g % 2) + 1 FROM generate_series(1, 20) g`) + defer testhelper.AssertQueryRuns(backupConn, "DROP TABLE schema2.relayout") + rewriter := holdRewrite(`DROP TABLE schema2.relayout; + CREATE TABLE schema2.relayout (id int, d int) DISTRIBUTED BY (id) + PARTITION BY RANGE (d) (START (1) END (4) EVERY (1)); + INSERT INTO schema2.relayout SELECT g, (g % 3) + 1 FROM generate_series(1, 30) g`) + done, output, err := startBackup("--include-table", "schema2.relayout") + defer waitForBackupExit(done) + defer releaseRewrite(rewriter) + + waitForBackupLockWait(backupConn) + rewriter.MustCommit() + waitForBackupExit(done) + + Expect(*err).ToNot(HaveOccurred(), *output) + Expect(*output).To(ContainSubstring("Backup completed successfully")) + Expect(*output).ToNot(ContainSubstring("[CRITICAL]")) + Expect(*output).To(ContainSubstring(changedTableWarning)) + Expect(*output).ToNot(ContainSubstring(dataSkippedWarning)) + + timestamp := getBackupTimestamp(*output) + tocStruct := readTOC(timestamp) + Expect(findDataEntry(tocStruct, "schema2", "relayout").RowsCopied).To(Equal(int64(30))) + + // A filtered backup carries no CREATE SCHEMA. + testhelper.AssertQueryRuns(restoreConn, "CREATE SCHEMA schema2") + gprestore(gprestorePath, restoreHelperPath, timestamp, + "--redirect-db", "restoredb", "--backup-dir", backupDir) + assertDataRestored(restoreConn, map[string]int{ + "schema2.relayout": 30, + "schema2.relayout_1_prt_3": 10, + }) + }) + + It("retries when a partition is truncated while the backup copies the parent", func() { + // Without --leaf-partition-data the parent is locked and copied and the + // partitions are only reached through it. Truncate and reload one + // partition so its storage changes while the row count stays the same. + truncate := "TRUNCATE schema2.returns_1_prt_jan17" + if backupConn.Version.Before("7") { + truncate = "ALTER TABLE schema2.returns TRUNCATE PARTITION jan17" + } + rewriter := holdRewrite(fmt.Sprintf(`CREATE TEMP TABLE saved_jan17 AS SELECT * FROM schema2.returns_1_prt_jan17; + %s; INSERT INTO schema2.returns SELECT * FROM saved_jan17`, truncate)) + done, output, err := startBackup() + defer waitForBackupExit(done) + defer releaseRewrite(rewriter) + + waitForBackupLockWait(backupConn) + rewriter.MustCommit() + waitForBackupExit(done) + + Expect(*err).ToNot(HaveOccurred(), *output) + Expect(*output).To(ContainSubstring("Backup completed successfully")) + Expect(*output).ToNot(ContainSubstring("[CRITICAL]")) + Expect(*output).To(ContainSubstring(changedTableWarning)) + Expect(*output).To(ContainSubstring("schema2.returns_1_prt_jan17")) + Expect(*output).ToNot(ContainSubstring(dataSkippedWarning)) + + timestamp := getBackupTimestamp(*output) + tocStruct := readTOC(timestamp) + Expect(findDataEntry(tocStruct, "schema2", "returns").RowsCopied).To(Equal(int64(6))) + + gprestore(gprestorePath, restoreHelperPath, timestamp, + "--redirect-db", "restoredb", "--backup-dir", backupDir) + assertDataRestored(restoreConn, map[string]int{"schema2.returns": 6}) + }) + + It("takes only as many snapshots as WHPGBACKUP_SNAPSHOT_ATTEMPTS allows", func() { + // One attempt means no retry: the first change is final and the + // table's data is skipped. + rewriter := holdRewrite("ALTER TABLE schema2.ao1 SET WITH (reorganize=true)") + done, output, err := startBackupEnv([]string{snapshotAttemptsEnvVar + "=1"}, "--leaf-partition-data") + defer waitForBackupExit(done) + defer releaseRewrite(rewriter) + + waitForBackupLockWait(backupConn) + rewriter.MustCommit() + waitForBackupExit(done) + + Expect(*err).ToNot(HaveOccurred(), *output) + Expect(*output).To(ContainSubstring("Backup completed successfully")) + Expect(*output).ToNot(ContainSubstring(newSnapshotWarning)) + Expect(*output).To(ContainSubstring("still changing after 1 attempt(s)")) + Expect(*output).To(ContainSubstring(dataSkippedWarning)) + Expect(*output).To(ContainSubstring("schema2.ao1")) + + timestamp := getBackupTimestamp(*output) + tocStruct := readTOC(timestamp) + Expect(findDataEntry(tocStruct, "schema2", "ao1")).To(BeNil()) + Expect(findDataEntry(tocStruct, "schema2", "ao2").RowsCopied).To(Equal(int64(1000))) + }) + + It("rejects an invalid WHPGBACKUP_SNAPSHOT_ATTEMPTS before the backup starts", func() { + for _, value := range []string{"0", "-1", "abc", "2.5"} { + output, err := runBackupEnv([]string{snapshotAttemptsEnvVar + "=" + value}) + Expect(err).To(HaveOccurred(), output) + Expect(output).To(ContainSubstring(snapshotAttemptsEnvVar + " must be a whole number of at least 1")) + Expect(output).ToNot(ContainSubstring("Backup Timestamp")) + } + }) + + It("skips the data and statistics of a table that keeps changing after every attempt and completes", func() { + // Each attempt must find the table rewritten again. Queue the next + // rewrite behind the current one before releasing it: gpbackup's + // AccessShareLock request is ahead of it in the lock queue, so the + // attempt is granted, fails its check, rolls back, and only then does + // the queued rewrite start and block the next attempt. + const attempts = 3 + rewriter := holdRewrite("ALTER TABLE schema2.ao1 SET WITH (reorganize=true)") + done, output, err := startBackup("--leaf-partition-data", "--with-stats") + defer waitForBackupExit(done) + defer releaseRewrite(rewriter) + + for attempt := 1; attempt <= attempts; attempt++ { + waitForBackupLockWait(backupConn) + if attempt == attempts { + rewriter.MustCommit() + break + } + next := testutils.SetupTestDbConn("testdb") + defer releaseRewrite(next) + queued := make(chan struct{}) + go func(conn *dbconn.DBConn) { + defer GinkgoRecover() + conn.MustBegin() + // LOCK first so the rewrite is queued as a single lock request; + // the ALTER then runs once the lock is granted. + conn.MustExec(`LOCK TABLE schema2.ao1 IN ACCESS EXCLUSIVE MODE; + ALTER TABLE schema2.ao1 SET WITH (reorganize=true)`) + close(queued) + }(next) + waitForLockWaiter(backupConn, "schema2", "ao1", "AccessExclusiveLock") + rewriter.MustCommit() + <-queued + rewriter = next + } + waitForBackupExit(done) + + Expect(*err).ToNot(HaveOccurred(), *output) + Expect(*output).To(ContainSubstring("Backup completed successfully")) + Expect(*output).ToNot(ContainSubstring("[CRITICAL]")) + Expect(*output).To(ContainSubstring(dataSkippedWarning)) + Expect(*output).To(ContainSubstring("schema2.ao1")) + + timestamp := getBackupTimestamp(*output) + tocStruct := readTOC(timestamp) + // DDL kept, data and incremental entry dropped, other tables untouched. + Expect(hasTableDDL(tocStruct, "schema2", "ao1")).To(BeTrue()) + Expect(findDataEntry(tocStruct, "schema2", "ao1")).To(BeNil()) + Expect(tocStruct.IncrementalMetadata.AO).ToNot(HaveKey("schema2.ao1")) + Expect(findDataEntry(tocStruct, "schema2", "ao2").RowsCopied).To(Equal(int64(1000))) + Expect(tocStruct.IncrementalMetadata.AO).To(HaveKey("schema2.ao2")) + + reportContents := string(getMetdataFileContents(backupDir, timestamp, "report")) + Expect(reportContents).To(ContainSubstring("schema2.ao1")) + + // Statistics describe data the backup does not have; the skipped table + // gets none, the others keep theirs. + statistics := string(getMetdataFileContents(backupDir, timestamp, "statistics.sql")) + Expect(statistics).ToNot(ContainSubstring("schema2.ao1")) + Expect(statistics).To(ContainSubstring("schema2.ao2")) + + // The backup restores: the skipped table exists and is empty. + gprestore(gprestorePath, restoreHelperPath, timestamp, + "--redirect-db", "restoredb", "--backup-dir", backupDir) + assertDataRestored(restoreConn, map[string]int{ + "schema2.ao1": 0, + "schema2.ao2": 1000, + "public.foo": 40000, + }) + }) +}) diff --git a/integration/changed_relations_test.go b/integration/changed_relations_test.go new file mode 100644 index 00000000..41d1e75b --- /dev/null +++ b/integration/changed_relations_test.go @@ -0,0 +1,369 @@ +package integration + +import ( + "fmt" + "strconv" + "strings" + + "github.com/greenplum-db/gpbackup/backup" + "github.com/greenplum-db/gpbackup/options" + "github.com/greenplum-db/gpbackup/testutils" + "github.com/spf13/cobra" + "github.com/warehouse-pg/common-go-libs/dbconn" + "github.com/warehouse-pg/common-go-libs/gplog" + "github.com/warehouse-pg/common-go-libs/testhelper" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +/* + * A table whose storage is replaced between the backup snapshot and LockTables + * cannot be read consistently under that snapshot. These tests run the DDL on + * a second connection after the backup transaction has fixed its snapshot and + * before the tables are locked, which is exactly the window gpbackup exposes. + */ +var _ = Describe("Tables changed between snapshot and lock", func() { + const ( + aoTable = "public.ao_changed" + heapTable = "public.heap_changed" + stableTable = "public.ao_stable" + partTable = "public.part_changed" + partLeaf = "public.part_changed_1_prt_2" + ) + var ddlConn *dbconn.DBConn + + // beginBackupTransaction opens the backup transaction on the pool and runs + // a first catalog read so the repeatable-read snapshot is fixed now, before + // the concurrent DDL of each test. + // gpbackup sets its session GUCs inside this transaction, so they are + // part of what a retry has to preserve. + beginBackupTransaction := func() { + connectionPool.MustBegin(0) + if connectionPool.Version.AtLeast(backup.SNAPSHOT_GPDB_MIN_VERSION) { + snapshot, err := backup.GetSynchronizedSnapshot(connectionPool) + Expect(err).ToNot(HaveOccurred()) + backup.SetBackupSnapshot(snapshot) + } + backup.SetSessionGUCs(0) + } + + // truncateLeafSQL empties the second partition and reloads the same rows, + // so the leaf gets a new relfilenode while the row count stays the same. + truncateLeafSQL := func() string { + truncate := fmt.Sprintf("TRUNCATE %s", partLeaf) + if connectionPool.Version.Before("7") { + truncate = fmt.Sprintf("ALTER TABLE %s TRUNCATE PARTITION FOR (RANK(2))", partTable) + } + return fmt.Sprintf("%s; INSERT INTO %s SELECT g, 2 FROM generate_series(1, 10) g", truncate, partTable) + } + + searchPath := func() string { + return dbconn.MustSelectString(connectionPool, "SELECT current_setting('search_path') AS string") + } + + lockedRelations := func() []backup.Relation { + included := backup.GetIncludedUserTableRelations(connectionPool, backup.IncludedRelationFqns) + relations := backup.ConvertRelationsOptionsToBackup(included) + backup.LockTables(connectionPool, relations) + return relations + } + + fqns := func(changed []backup.ChangedRelation) []string { + result := make([]string, 0, len(changed)) + for _, rel := range changed { + result = append(result, rel.FQN()) + } + return result + } + + dataTableFQNs := func(tables []backup.Table) []string { + result := make([]string, 0, len(tables)) + for _, table := range tables { + result = append(result, table.FQN()) + } + return result + } + + BeforeEach(func() { + gplog.SetVerbosity(gplog.LOGERROR) // turn off the progress bar in LockTables + var rootCmd = &cobra.Command{} + backup.DoInit(rootCmd) // initializes objectCounts, but also replaces the test logger + _, stderr, logFile = testhelper.SetupTestLogger() + backup.UseCmdFlags(backupCmdFlags) + backup.SetMaxSnapshotAttempts(3) + // The suite keeps search_path at pg_catalog; start each test from the + // server default so only the GUCs set inside the transaction count. + testhelper.AssertQueryRuns(connectionPool, "RESET search_path") + + testhelper.AssertQueryRuns(connectionPool, fmt.Sprintf(` + CREATE TABLE %s (i int) WITH (appendonly=true) DISTRIBUTED BY (i); + INSERT INTO %s SELECT generate_series(1, 1000); + CREATE TABLE %s (i int) DISTRIBUTED BY (i); + INSERT INTO %s SELECT generate_series(1, 100); + CREATE TABLE %s (i int) WITH (appendonly=true) DISTRIBUTED BY (i); + INSERT INTO %s SELECT generate_series(1, 10); + CREATE TABLE %s (id int, d int) WITH (appendonly=true) DISTRIBUTED BY (id) + PARTITION BY RANGE (d) (START (1) END (4) EVERY (1)); + INSERT INTO %s SELECT g, (g %% 3) + 1 FROM generate_series(1, 30) g;`, + aoTable, aoTable, heapTable, heapTable, stableTable, stableTable, partTable, partTable)) + + for _, table := range []string{aoTable, heapTable, stableTable, partTable} { + _ = backupCmdFlags.Set(options.INCLUDE_RELATION, table) + } + opts, err := options.NewOptions(backupCmdFlags) + Expect(err).ToNot(HaveOccurred()) + backup.ValidateAndProcessFilterLists(opts) + includeOids := backup.GetOidsFromRelationList(backup.IncludedRelationFqns) + Expect(backup.ExpandIncludesForPartitions(connectionPool, opts, includeOids, backupCmdFlags)).To(Succeed()) + + ddlConn = testutils.SetupTestDbConn("testdb") + }) + AfterEach(func() { + if connectionPool.Tx[0] != nil { + _ = connectionPool.Rollback(0) + } + backup.SetBackupSnapshot("") + backup.SetMaxSnapshotAttempts(3) + testhelper.AssertQueryRuns(connectionPool, "SET search_path TO pg_catalog") + ddlConn.Close() + testhelper.AssertQueryRuns(connectionPool, fmt.Sprintf( + "DROP TABLE IF EXISTS %s; DROP TABLE IF EXISTS %s; DROP TABLE IF EXISTS %s; DROP TABLE IF EXISTS %s", + aoTable, heapTable, stableTable, partTable)) + }) + + Describe("GetRelationsChangedSinceSnapshot", func() { + It("reports nothing when no table changed", func() { + beginBackupTransaction() + relations := lockedRelations() + + changed := backup.GetRelationsChangedSinceSnapshot(connectionPool, relations) + Expect(changed).To(BeEmpty()) + }) + It("reports an AO table rewritten by ALTER TABLE SET WITH (reorganize=true)", func() { + beginBackupTransaction() + testhelper.AssertQueryRuns(ddlConn, fmt.Sprintf("ALTER TABLE %s SET WITH (reorganize=true)", aoTable)) + relations := lockedRelations() + + changed := backup.GetRelationsChangedSinceSnapshot(connectionPool, relations) + Expect(fqns(changed)).To(ConsistOf(aoTable)) + Expect(changed[0].Dropped).To(BeFalse()) + }) + It("reports a truncated heap table", func() { + beginBackupTransaction() + testhelper.AssertQueryRuns(ddlConn, fmt.Sprintf("TRUNCATE %s", heapTable)) + relations := lockedRelations() + + changed := backup.GetRelationsChangedSinceSnapshot(connectionPool, relations) + Expect(fqns(changed)).To(ConsistOf(heapTable)) + Expect(changed[0].Dropped).To(BeFalse()) + }) + It("reports a heap table rewritten by VACUUM FULL", func() { + beginBackupTransaction() + testhelper.AssertQueryRuns(ddlConn, fmt.Sprintf("VACUUM FULL %s", heapTable)) + relations := lockedRelations() + + changed := backup.GetRelationsChangedSinceSnapshot(connectionPool, relations) + Expect(fqns(changed)).To(ConsistOf(heapTable)) + Expect(changed[0].Dropped).To(BeFalse()) + }) + It("reports a truncated partition below a locked table with the table as its ancestor", func() { + // Without --leaf-partition-data the backup locks and copies the + // parent; the check has to reach the partition through pg_inherits. + beginBackupTransaction() + testhelper.AssertQueryRuns(ddlConn, truncateLeafSQL()) + relations := lockedRelations() + + changed := backup.GetRelationsChangedSinceSnapshot(connectionPool, relations) + Expect(fqns(changed)).To(ConsistOf(partLeaf)) + Expect(changed[0].Dropped).To(BeFalse()) + Expect(changed[0].Ancestors).ToNot(BeEmpty()) + Expect(changed[0].Ancestors[len(changed[0].Ancestors)-1].FQN()).To(Equal(partTable)) + }) + It("reports a table dropped and recreated under the same name as dropped", func() { + beginBackupTransaction() + testhelper.AssertQueryRuns(ddlConn, fmt.Sprintf( + "DROP TABLE %s; CREATE TABLE %s (i int) WITH (appendonly=true) DISTRIBUTED BY (i)", aoTable, aoTable)) + relations := lockedRelations() + + changed := backup.GetRelationsChangedSinceSnapshot(connectionPool, relations) + Expect(fqns(changed)).To(ConsistOf(aoTable)) + Expect(changed[0].Dropped).To(BeTrue()) + }) + It("reports every changed table when several change at once", func() { + beginBackupTransaction() + testhelper.AssertQueryRuns(ddlConn, fmt.Sprintf( + "ALTER TABLE %s SET WITH (reorganize=true); ALTER TABLE %s SET WITH (reorganize=true)", aoTable, heapTable)) + relations := lockedRelations() + + changed := backup.GetRelationsChangedSinceSnapshot(connectionPool, relations) + Expect(fqns(changed)).To(ConsistOf(aoTable, heapTable)) + }) + }) + + Describe("RetrieveAndProcessTables", func() { + It("backs up all tables with no warning when nothing changed", func() { + logStart := len(logFile.Contents()) + beginBackupTransaction() + + metadataTables, dataTables := backup.RetrieveAndProcessTables() + + Expect(dataTableFQNs(dataTables)).To(ConsistOf(aoTable, heapTable, stableTable, partTable)) + Expect(dataTableFQNs(metadataTables)).To(ContainElements(aoTable, heapTable, stableTable, partTable)) + Expect(backup.GetSkippedDataTables()).To(BeEmpty()) + Expect(string(logFile.Contents()[logStart:])).ToNot(ContainSubstring("changed on disk")) + + aoMetadata := backup.GetAOIncrementalMetadata(connectionPool) + Expect(aoMetadata).To(HaveKey(aoTable)) + Expect(aoMetadata).To(HaveKey(stableTable)) + }) + It("retries with a new snapshot and backs up a table rewritten in the window", func() { + logStart := len(logFile.Contents()) + beginBackupTransaction() + testhelper.AssertQueryRuns(ddlConn, fmt.Sprintf( + "ALTER TABLE %s SET WITH (reorganize=true); ALTER TABLE %s SET WITH (reorganize=true)", aoTable, heapTable)) + + _, dataTables := backup.RetrieveAndProcessTables() + + Expect(dataTableFQNs(dataTables)).To(ConsistOf(aoTable, heapTable, stableTable, partTable)) + Expect(backup.GetSkippedDataTables()).To(BeEmpty()) + logged := string(logFile.Contents()[logStart:]) + Expect(logged).To(ContainSubstring("[WARNING]:-")) + Expect(logged).To(ContainSubstring("changed on disk after the backup snapshot was taken")) + Expect(logged).To(ContainSubstring(aoTable)) + Expect(logged).To(ContainSubstring(heapTable)) + Expect(logged).To(ContainSubstring("taking a new snapshot (attempt 2 of 3)")) + Expect(logged).ToNot(ContainSubstring("not backed up")) + + // The new transaction carries the same session GUCs as the first + // one, otherwise object names would be dumped unqualified. + Expect(searchPath()).To(Equal("pg_catalog")) + + // The new snapshot sees the rewritten storage: the AO helper resolves + // and the row counts are the current ones. + aoMetadata := backup.GetAOIncrementalMetadata(connectionPool) + Expect(aoMetadata).To(HaveKey(aoTable)) + Expect(dbconn.MustSelectString(connectionPool, + fmt.Sprintf("SELECT count(*)::text AS string FROM %s", heapTable))).To(Equal("100")) + }) + It("re-resolves the include list when a table was dropped and recreated in the window", func() { + oldOid := backup.IncludedRelationFqns[0].Oid + for _, rel := range backup.IncludedRelationFqns { + if rel.Schema == "public" && rel.Name == "ao_changed" { + oldOid = rel.Oid + } + } + beginBackupTransaction() + testhelper.AssertQueryRuns(ddlConn, fmt.Sprintf(`DROP TABLE %s; + CREATE TABLE %s (i int) WITH (appendonly=true) DISTRIBUTED BY (i); + INSERT INTO %s SELECT generate_series(1, 7)`, aoTable, aoTable, aoTable)) + + _, dataTables := backup.RetrieveAndProcessTables() + + Expect(dataTableFQNs(dataTables)).To(ConsistOf(aoTable, heapTable, stableTable, partTable)) + for _, table := range dataTables { + if table.FQN() == aoTable { + Expect(table.Oid).ToNot(Equal(oldOid)) + } + } + Expect(searchPath()).To(Equal("pg_catalog")) + Expect(backup.GetAOIncrementalMetadata(connectionPool)).To(HaveKey(aoTable)) + Expect(dbconn.MustSelectString(connectionPool, + fmt.Sprintf("SELECT count(*)::text AS string FROM %s", aoTable))).To(Equal("7")) + }) + It("expands the include list again when a partitioned table is dropped and recreated with more partitions in the window", func() { + // The include list was expanded to the partitions when the backup + // started; the retry must expand it again from the user's names, + // or a partition created in the window is missed. + beginBackupTransaction() + testhelper.AssertQueryRuns(ddlConn, fmt.Sprintf(`DROP TABLE %s; + CREATE TABLE %s (id int, d int) WITH (appendonly=true) DISTRIBUTED BY (id) + PARTITION BY RANGE (d) (START (1) END (5) EVERY (1)); + INSERT INTO %s SELECT g, (g %% 4) + 1 FROM generate_series(1, 40) g`, partTable, partTable, partTable)) + + _, dataTables := backup.RetrieveAndProcessTables() + + Expect(dataTableFQNs(dataTables)).To(ConsistOf(aoTable, heapTable, stableTable, partTable)) + Expect(backup.GetSkippedDataTables()).To(BeEmpty()) + included := make([]string, 0, len(backup.IncludedRelationFqns)) + for _, relation := range backup.IncludedRelationFqns { + included = append(included, relation.Schema+"."+relation.Name) + } + Expect(included).To(ContainElement(partTable)) + if connectionPool.Version.AtLeast("7") { + // 7.x lists the partitions themselves; the new fourth one is there. + Expect(included).To(ContainElement(partTable + "_1_prt_4")) + } + // Every OID in the list is a live relation under the new snapshot. + oids := backup.GetOidsFromRelationList(backup.IncludedRelationFqns) + Expect(dbconn.MustSelectString(connectionPool, fmt.Sprintf( + "SELECT count(*)::text AS string FROM pg_class WHERE oid IN (%s)", strings.Join(oids, ",")))). + To(Equal(strconv.Itoa(len(oids)))) + Expect(searchPath()).To(Equal("pg_catalog")) + Expect(dbconn.MustSelectString(connectionPool, + fmt.Sprintf("SELECT count(*)::text AS string FROM %s", partTable))).To(Equal("40")) + }) + It("skips the data of a table still changed after the last attempt", func() { + backup.SetMaxSnapshotAttempts(1) + logStart := len(logFile.Contents()) + beginBackupTransaction() + testhelper.AssertQueryRuns(ddlConn, fmt.Sprintf("ALTER TABLE %s SET WITH (reorganize=true)", aoTable)) + + metadataTables, dataTables := backup.RetrieveAndProcessTables() + + Expect(dataTableFQNs(dataTables)).To(ConsistOf(heapTable, stableTable, partTable)) + Expect(dataTableFQNs(metadataTables)).To(ContainElements(aoTable, heapTable, stableTable, partTable)) + Expect(backup.GetSkippedDataTables()).To(HaveKey(aoTable)) + logged := string(logFile.Contents()[logStart:]) + Expect(logged).To(ContainSubstring("[WARNING]:-")) + Expect(logged).To(ContainSubstring(fmt.Sprintf("Data for table %s not backed up", aoTable))) + + // The stale pg_aoseg name of the skipped table must not be queried. + aoMetadata := backup.GetAOIncrementalMetadata(connectionPool) + Expect(aoMetadata).ToNot(HaveKey(aoTable)) + Expect(aoMetadata).To(HaveKey(stableTable)) + }) + It("retries when a partition of a table copied through its parent is truncated in the window", func() { + beginBackupTransaction() + testhelper.AssertQueryRuns(ddlConn, truncateLeafSQL()) + + _, dataTables := backup.RetrieveAndProcessTables() + + Expect(dataTableFQNs(dataTables)).To(ConsistOf(aoTable, heapTable, stableTable, partTable)) + Expect(backup.GetSkippedDataTables()).To(BeEmpty()) + // Without --leaf-partition-data the incremental entry is the parent's + // on 6.x, where the parent has storage, and the leaf's on 7.x, where + // it has none. + aoMetadata := backup.GetAOIncrementalMetadata(connectionPool) + if connectionPool.Version.Before("7") { + Expect(aoMetadata).To(HaveKey(partTable)) + } else { + Expect(aoMetadata).To(HaveKey(partLeaf)) + } + Expect(dbconn.MustSelectString(connectionPool, + fmt.Sprintf("SELECT count(*)::text AS string FROM %s", partTable))).To(Equal("30")) + }) + It("skips the data of the parent when a partition keeps changing", func() { + // The parent is the table whose COPY reads the partition, so it is + // the one that has to leave the data set. + backup.SetMaxSnapshotAttempts(1) + logStart := len(logFile.Contents()) + beginBackupTransaction() + testhelper.AssertQueryRuns(ddlConn, truncateLeafSQL()) + + metadataTables, dataTables := backup.RetrieveAndProcessTables() + + Expect(dataTableFQNs(dataTables)).To(ConsistOf(aoTable, heapTable, stableTable)) + Expect(dataTableFQNs(metadataTables)).To(ContainElement(partTable)) + Expect(backup.GetSkippedDataTables()).To(HaveKey(partTable)) + logged := string(logFile.Contents()[logStart:]) + Expect(logged).To(ContainSubstring(fmt.Sprintf("partition of %s", partTable))) + Expect(logged).To(ContainSubstring(fmt.Sprintf("Data for table %s not backed up", partTable))) + + aoMetadata := backup.GetAOIncrementalMetadata(connectionPool) + Expect(aoMetadata).ToNot(HaveKey(partLeaf)) + Expect(aoMetadata).ToNot(HaveKey(partTable)) + Expect(aoMetadata).To(HaveKey(stableTable)) + }) + }) +}) diff --git a/integration/integration_suite_test.go b/integration/integration_suite_test.go index f2837c64..880f031b 100644 --- a/integration/integration_suite_test.go +++ b/integration/integration_suite_test.go @@ -139,12 +139,18 @@ var _ = AfterSuite(func() { Fail("Could not remove /tmp/testdir* directories on 1 or more hosts") } } + connection1 := testutils.SetupTestDbConn("template1") if connectionPool != nil { connectionPool.Close() + // The pool's client connections are closed above, but the server + // backends exit asynchronously and an idle one can linger. DROP + // DATABASE waits only about five seconds for other backends before it + // fails, so terminate any backend still on testdb first, from this + // template1 connection, to keep dropdb from racing that timeout. + _, _ = connection1.Exec("SELECT pg_terminate_backend(pid) FROM pg_stat_activity WHERE datname = 'testdb' AND pid <> pg_backend_pid()") err := exec.Command("dropdb", "testdb").Run() Expect(err).To(BeNil()) } - connection1 := testutils.SetupTestDbConn("template1") testhelper.AssertQueryRuns(connection1, "DROP ROLE testrole") testhelper.AssertQueryRuns(connection1, "DROP ROLE anothertestrole") connection1.Close() diff --git a/options/options.go b/options/options.go index b0f24d77..9fd2bf76 100644 --- a/options/options.go +++ b/options/options.go @@ -152,6 +152,12 @@ func (o *Options) AddIncludedRelation(relation string) { o.IncludedRelations = append(o.IncludedRelations, relation) } +// ResetIncludedRelations drops the relations added by AddIncludedRelation, +// leaving the ones the user gave. +func (o *Options) ResetIncludedRelations() { + o.IncludedRelations = append([]string(nil), o.originalIncludedRelations...) +} + type Relation struct { SchemaOid uint32 Oid uint32 diff --git a/report/report.go b/report/report.go index e08a2f4e..a5efe29d 100644 --- a/report/report.go +++ b/report/report.go @@ -29,6 +29,9 @@ import ( type Report struct { BackupParamsString string DatabaseSize string + // Tables whose data was left out of the backup because their storage kept + // changing between the backup snapshot and the table locks. + SkippedDataTables []string history.BackupConfig } @@ -151,6 +154,14 @@ func (report *Report) WriteBackupReportFile(reportFilename string, timestamp str LineInfo{}, LineInfo{Key: "backup status:", Value: history.BackupStatusSucceed}) } + // A paragraph of its own: the report parser reads the value of "backup + // error:" up to the next blank line, so this line must not follow it + // directly. + if len(report.SkippedDataTables) > 0 { + reportInfo = append(reportInfo, + LineInfo{}, + LineInfo{Key: "data not backed up:", Value: strings.Join(report.SkippedDataTables, ", ")}) + } reportInfo = append(reportInfo, LineInfo{}) if report.DatabaseSize != "" { reportInfo = append(reportInfo, diff --git a/report/report_test.go b/report/report_test.go index 24a21084..1329647b 100644 --- a/report/report_test.go +++ b/report/report_test.go @@ -152,6 +152,17 @@ count of database objects in backup: sequences 1 tables 42 types 1000`)) + }) + It("lists the tables whose data was not backed up in a paragraph of their own", func() { + backupReport.SkippedDataTables = []string{"public.ao_t", "public.heap_t"} + backupReport.WriteBackupReportFile("filename", timestamp, endtime, objectCounts, "Cannot access /tmp/backups: Permission denied") + Expect(buffer).To(Say(`backup status: Failure +backup error: Cannot access /tmp/backups: Permission denied + +data not backed up: public\.ao_t, public\.heap_t + +database size: 42 MB +segment count: 3`)) }) It("writes a report without database size information", func() { backupReport.DatabaseSize = ""