Skip to content

Commit dfd9b27

Browse files
Merge pull request #1307 from percona/release-2.14.0
Release 2.14.0
2 parents bf61ba0 + d3436ef commit dfd9b27

1,003 files changed

Lines changed: 72672 additions & 37636 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

.github/workflows/ci.yml

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -38,9 +38,15 @@ concurrency:
3838
group: ${{ github.workflow }}-${{ github.ref }}
3939
cancel-in-progress: true
4040

41+
permissions: read-all
42+
4143
jobs:
4244
test:
4345
runs-on: ubuntu-latest
46+
permissions:
47+
contents: read
48+
checks: write
49+
actions: write
4450
timeout-minutes: 180
4551
strategy:
4652
fail-fast: false
@@ -130,6 +136,9 @@ jobs:
130136
if: ${{ always() }}
131137
needs: test
132138
runs-on: ubuntu-latest
139+
permissions:
140+
contents: read
141+
actions: read
133142
steps:
134143
- uses: actions/checkout@v4
135144
- name: Set up Go

.github/workflows/codecov.yml

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,8 @@ on:
1010
- main
1111
- dev
1212

13+
permissions: read-all
14+
1315
jobs:
1416
go-test:
1517
name: runner / go-test

.github/workflows/reviewdog.yml

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,10 @@
11
name: reviewdog
22
on: [pull_request]
3+
4+
permissions:
5+
contents: read
6+
pull-requests: write
7+
38
jobs:
49
shellcheck:
510
name: runner / shellcheck

.github/workflows/trivy.yml

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,8 @@ on:
44
branches: ["main", "dev"]
55
pull_request:
66
branches: ["main", "dev"]
7+
permissions: read-all
8+
79
jobs:
810
scan:
911
name: Trivy

SECURITY.md

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
1+
# Reporting security vulnerabilities
2+
3+
Please report any vulnerabilities privately by emailing our security team at [security@percona.com](mailto:security@percona.com).
4+
Do not create public Jira issues for security vulnerabilities or disclose details in publicly accessible trackers.
5+
6+
Once a vulnerability has been reported, it should be reviewed, recorded, assigned, and fixed or addressed as appropriate.
7+
Based initially on an internal severity assessment (for example, using CVSS or our internal severity ratings; if a CVE is later assigned, its criticality may inform or adjust this assessment and may be reduced internally with certain compensating controls that mitigate the severity level), the following guidelines apply to ensure that issues are resolved or addressed in some concrete manner:
8+
9+
1. Critical – immediately, but no later than 14 days;
10+
2. High – as soon as possible, no later than 30 days;
11+
3. Medium – within 60 days; and
12+
4. Low – make best efforts to patch low-rated vulnerabilities within 90 days.
13+
14+
Should you have a legitimate test case that might include one of the above, please contact [security@percona.com](mailto:security@percona.com) detailing your proposed test, expected outcome, and proposed timelines.
15+
16+
For details, see [Percona Security](https://www.percona.com/security).

cmd/pbm-agent/pitr.go

Lines changed: 64 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -57,6 +57,58 @@ func (a *Agent) removePitr() {
5757
a.setPitr(nil)
5858
}
5959

60+
// waitForPITRSlicerStop polls until PITR config is disabled and the OpLock is stale.
61+
func (a *Agent) waitForPITRSlicerStop(ctx context.Context, rs string, l log.LogEvent) error {
62+
ctx, cancel := context.WithTimeout(ctx, 120*time.Second)
63+
defer cancel()
64+
65+
tk := time.NewTicker(pitrHb)
66+
defer tk.Stop()
67+
68+
l.Debug("waiting for PITR config disabled and OpLock released")
69+
for {
70+
select {
71+
case <-ctx.Done():
72+
return errors.Wrap(ctx.Err(), "done waiting for oplog slicer to stop")
73+
case <-tk.C:
74+
// Config is disabled
75+
cfg, err := config.GetConfig(ctx, a.leadConn)
76+
if err != nil {
77+
return errors.Wrap(err, "get config")
78+
}
79+
if cfg.PITR.Enabled {
80+
continue
81+
}
82+
83+
// No active lock
84+
locks, err := lock.GetOpLocks(ctx, a.leadConn, &lock.LockHeader{
85+
Type: ctrl.CmdPITR,
86+
Replset: rs,
87+
})
88+
if err != nil {
89+
return errors.Wrap(err, "get PITR op locks")
90+
}
91+
92+
ts, err := topo.GetClusterTime(ctx, a.leadConn)
93+
if err != nil {
94+
return errors.Wrap(err, "get cluster time")
95+
}
96+
97+
active := false
98+
for i := range locks {
99+
if locks[i].Heartbeat.T+defs.StaleFrameSec >= ts.T {
100+
active = true
101+
break
102+
}
103+
}
104+
if !active {
105+
l.Info("PITR slicer stopped")
106+
return nil
107+
}
108+
}
109+
}
110+
}
111+
60112
func (a *Agent) getPitr() *currentPitr {
61113
a.slicerMx.Lock()
62114
defer a.slicerMx.Unlock()
@@ -819,6 +871,18 @@ func (a *Agent) pitrActivityMonitor(ctx context.Context) {
819871
continue
820872
}
821873

874+
// If any RS reported an error, let the error monitor handle it.
875+
// Acting here would set StatusReconfig, masking the error and
876+
// preventing proper error handling by the slicers on stop.
877+
rsErrors, err := oplog.GetReplSetsWithStatus(ctx, a.leadConn, oplog.StatusError)
878+
if err != nil && !errors.Is(err, errors.ErrNotFound) {
879+
l.Error("activity check RS errors: %v", err)
880+
continue
881+
}
882+
if len(rsErrors) > 0 {
883+
continue
884+
}
885+
822886
ackedAgents, err := oplog.GetAgentsWithACK(ctx, a.leadConn)
823887
if err != nil {
824888
l.Error("activity get acked agents", err)

cmd/pbm-agent/restore.go

Lines changed: 9 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -78,39 +78,24 @@ func (a *Agent) Restore(ctx context.Context, r *ctrl.RestoreCmd, opid ctrl.OPID,
7878
l.Error("release lock: %v", err)
7979
}
8080
}()
81+
}
8182

83+
// disable pitr from leader node
84+
if nodeInfo.IsClusterLeader() {
8285
err = config.SetConfigVar(ctx, a.leadConn, "pitr.enabled", "false")
8386
if err != nil {
8487
l.Error("disable oplog slicer: %v", err)
8588
} else {
8689
l.Info("oplog slicer disabled")
8790
}
88-
a.removePitr()
8991
}
9092

91-
// stop balancer during the restore
92-
if a.brief.Sharded && nodeInfo.IsClusterLeader() {
93-
bs, err := topo.GetBalancerStatus(ctx, a.leadConn)
94-
if err != nil {
95-
l.Error("get balancer status: %v", err)
96-
return
97-
}
98-
99-
if bs.IsOn() {
100-
err := topo.SetBalancerStatus(ctx, a.leadConn, topo.BalancerModeOff)
101-
if err != nil {
102-
l.Error("set balancer off: %v", err)
103-
}
104-
105-
l.Debug("waiting for balancer off")
106-
bs := topo.WaitForBalancerDisabled(ctx, a.leadConn, time.Second*30, l)
107-
if bs.IsDisabled() {
108-
l.Debug("balancer is disabled")
109-
} else {
110-
l.Warning("balancer is not disabled: balancer mode: %s, in balancer round: %t",
111-
bs.Mode, bs.InBalancerRound)
112-
}
113-
}
93+
// Cancel the local slicer if running.
94+
a.removePitr()
95+
// Wait for the slicer to stop on every node.
96+
if err := a.waitForPITRSlicerStop(ctx, nodeInfo.SetName, l); err != nil {
97+
l.Error("unable to stop PITR slicer: %v", err)
98+
return
11499
}
115100

116101
var bcpType defs.BackupType

cmd/pbm/backup.go

Lines changed: 62 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -170,9 +170,10 @@ func runBackup(
170170
}
171171
fmt.Printf("Starting backup %q%s", b.name, pinfo)
172172
}
173-
startCtx, cancel := context.WithTimeout(ctx, cfg.Backup.Timeouts.StartingStatus())
174-
defer cancel()
175-
err = waitForBcpStatus(startCtx, conn, b.name, showProgress)
173+
err = waitForBcpStatus(ctx, conn, b.name, showProgress)
174+
if showProgress {
175+
fmt.Println()
176+
}
176177
if err != nil {
177178
return nil, errors.Wrap(err, "wait for backup status")
178179
}
@@ -213,7 +214,7 @@ func runBackup(
213214
}
214215

215216
if showProgress {
216-
fmt.Printf("\nWaiting for '%s' backup...", b.name)
217+
fmt.Printf("Waiting for '%s' backup...", b.name)
217218
}
218219
s, err := waitBackup(ctx, conn, b.name, defs.StatusDone, showProgress)
219220
if s != nil && showProgress {
@@ -272,6 +273,10 @@ func waitBackup(
272273
case defs.StatusError:
273274
return &bcp.Status, bcp.Error()
274275
}
276+
277+
if err := checkBackupStale(ctx, conn, bcp); err != nil {
278+
return &bcp.Status, err
279+
}
275280
}
276281

277282
if showProgress {
@@ -280,25 +285,69 @@ func waitBackup(
280285
}
281286
}
282287

283-
func waitForBcpStatus(ctx context.Context, conn connect.Client, bcpName string, showProgress bool) error {
288+
func checkBackupStale(ctx context.Context, conn connect.Client, bcp *backup.BackupMeta) error {
289+
clusterTime, err := topo.GetClusterTime(ctx, conn)
290+
if err != nil {
291+
return errors.Wrap(err, "read cluster time")
292+
}
293+
if bcp.Hb.T+defs.StaleFrameSec < clusterTime.T {
294+
rs := ""
295+
for _, s := range bcp.Replsets {
296+
rs += fmt.Sprintf("\n- %s: %v", s.Name, s.Status)
297+
if s.Error != "" {
298+
rs += ": " + s.Error
299+
}
300+
}
301+
return errors.Errorf(
302+
"backup stuck at %q status, last heartbeat: %d%s",
303+
bcp.Status, bcp.Hb.T, rs)
304+
}
305+
return nil
306+
}
307+
308+
func waitForBcpExists(ctx context.Context, conn connect.Client, bcpName string, showProgress bool) error {
284309
tk := time.NewTicker(time.Second)
285310
defer tk.Stop()
311+
to := time.After(defs.WaitBackupStart)
286312

287-
var bmeta *backup.BackupMeta
288313
for {
289314
select {
290315
case <-tk.C:
291316
if showProgress {
292317
fmt.Print(".")
293318
}
294-
var err error
295-
bmeta, err = backup.NewDBManager(conn).GetBackupByName(ctx, bcpName)
319+
_, err := backup.NewDBManager(conn).GetBackupByName(ctx, bcpName)
296320
if errors.Is(err, errors.ErrNotFound) {
297321
continue
298322
}
299323
if err != nil {
300324
return errors.Wrap(err, "get backup metadata")
301325
}
326+
return nil
327+
case <-to:
328+
return errors.New("no progress from leader, backup metadata not found")
329+
}
330+
}
331+
}
332+
333+
func waitForBcpStatus(ctx context.Context, conn connect.Client, bcpName string, showProgress bool) error {
334+
if err := waitForBcpExists(ctx, conn, bcpName, showProgress); err != nil {
335+
return err
336+
}
337+
338+
tk := time.NewTicker(time.Second)
339+
defer tk.Stop()
340+
341+
for {
342+
select {
343+
case <-tk.C:
344+
if showProgress {
345+
fmt.Print(".")
346+
}
347+
bmeta, err := backup.NewDBManager(conn).GetBackupByName(ctx, bcpName)
348+
if err != nil {
349+
return errors.Wrap(err, "get backup metadata")
350+
}
302351
switch bmeta.Status {
303352
case defs.StatusRunning, defs.StatusDumpDone, defs.StatusDone, defs.StatusCancelled:
304353
return nil
@@ -315,22 +364,12 @@ func waitForBcpStatus(ctx context.Context, conn connect.Client, bcpName string,
315364
}
316365
return errors.Errorf("status error on %s", rs)
317366
}
318-
case <-ctx.Done():
319-
if bmeta == nil {
320-
return errors.New("no progress from leader, backup metadata not found")
321-
}
322-
rs := ""
323-
for _, s := range bmeta.Replsets {
324-
rs += fmt.Sprintf("- Backup on replicaset \"%s\" in state: %v\n", s.Name, s.Status)
325-
if s.Error != "" {
326-
rs += ": " + s.Error
327-
}
328-
}
329-
if rs == "" {
330-
rs = "<no replset has started backup>\n"
331-
}
332367

333-
return errors.New("no confirmation that backup has successfully started. Replsets status:\n" + rs)
368+
if err := checkBackupStale(ctx, conn, bmeta); err != nil {
369+
return err
370+
}
371+
case <-ctx.Done():
372+
return ctx.Err()
334373
}
335374
}
336375
}

cmd/pbm/common.go

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,12 @@
11
package main
22

33
import (
4+
"bufio"
45
"context"
56
"encoding/json"
67
"fmt"
8+
"os"
9+
"strings"
710
"time"
811

912
"github.com/percona/percona-backup-mongodb/pbm/connect"
@@ -15,6 +18,33 @@ import (
1518

1619
var errWaitTimeout = errors.New("Operation is in progress. Check pbm status and logs")
1720

21+
var errUserCanceled = errors.New("canceled")
22+
23+
func askConfirmation(question string) error {
24+
fi, err := os.Stdin.Stat()
25+
if err != nil {
26+
return errors.Wrap(err, "stat stdin")
27+
}
28+
if (fi.Mode() & os.ModeCharDevice) == 0 {
29+
return errors.New("no tty")
30+
}
31+
32+
fmt.Printf("%s [y/N] ", question)
33+
34+
scanner := bufio.NewScanner(os.Stdin)
35+
scanner.Scan()
36+
if err := scanner.Err(); err != nil {
37+
return errors.Wrap(err, "read stdin")
38+
}
39+
40+
switch strings.TrimSpace(scanner.Text()) {
41+
case "yes", "Yes", "YES", "Y", "y":
42+
return nil
43+
}
44+
45+
return errUserCanceled
46+
}
47+
1848
func sendCmd(ctx context.Context, conn connect.Client, cmd ctrl.Cmd) error {
1949
cmd.TS = time.Now().UTC().Unix()
2050
_, err := conn.CmdStreamCollection().InsertOne(ctx, cmd)

0 commit comments

Comments
 (0)