Skip to content

Commit 1706883

Browse files
Allow remote-config to be generated from a standby
`kubectl dba remote-config` creates the replication user on the source before rendering the AppBinding and secrets, and that step needs a writable primary: CREATE ROLE / ALTER ROLE / GRANT are catalog writes a hot standby rejects. There are sources with no primary to exec into. A remote replica acting as the source of a chained replica is the clearest case: every member is a standby, and the replication role is already there anyway, carried over by physical WAL replication. Creating it is impossible and also unnecessary. Nothing in the generated config is primary-specific - the host comes from -d, the port from --port, the secret is synthesised and the client certificate's CN is the username - so only the side effect is coupled to the primary, not the artifact. Add --skip-user-creation, which leaves the source catalog alone and takes the given password at face value. It also drops the `status.phase == Ready` gate, since rendering manifests only reads the CR and a degraded source is exactly when a DR config is wanted. The role is still verified where possible: reading pg_roles (mysql.user for MySQL) works on a standby, so a missing role or one lacking the REPLICATION attribute fails with the SQL to run by hand, while an unreachable database only warns. Also fixes pod lookup, which returned nil - reported as success - when no pod matched, so the command wrote an auth secret for a user it never created and exited 0. It now fails with an actionable message, and skips pods that do not run the database container, such as the Postgres arbiter. Signed-off-by: souravbiswassanto <saurov@appscode.com>
1 parent a738b49 commit 1706883

7 files changed

Lines changed: 382 additions & 47 deletions

File tree

‎pkg/common/mysql.go‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -47,7 +47,9 @@ type MySQLOpts struct {
4747
ErrWriter *bytes.Buffer
4848
}
4949

50-
func NewMySQLOpts(f cmdutil.Factory, dbName, namespace string) (*MySQLOpts, error) {
50+
func NewMySQLOpts(f cmdutil.Factory, dbName, namespace string, options ...OptionFunc) (*MySQLOpts, error) {
51+
cfg := buildOptionConfig(options)
52+
5153
config, err := f.ToRESTConfig()
5254
if err != nil {
5355
return nil, err
@@ -75,7 +77,7 @@ func NewMySQLOpts(f cmdutil.Factory, dbName, namespace string) (*MySQLOpts, erro
7577
return nil, err
7678
}
7779

78-
if db.Status.Phase != dbapi.DatabasePhaseReady {
80+
if !cfg.skipReadinessCheck && db.Status.Phase != dbapi.DatabasePhaseReady {
7981
return nil, fmt.Errorf("MySQL %s/%s is not ready", namespace, dbName)
8082
}
8183

‎pkg/common/options.go‎

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,43 @@
1+
/*
2+
Copyright AppsCode Inc. and Contributors
3+
4+
Licensed under the AppsCode Community License 1.0.0 (the "License");
5+
you may not use this file except in compliance with the License.
6+
You may obtain a copy of the License at
7+
8+
https://github.com/appscode/licenses/raw/1.0.0/AppsCode-Community-1.0.0.md
9+
10+
Unless required by applicable law or agreed to in writing, software
11+
distributed under the License is distributed on an "AS IS" BASIS,
12+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
See the License for the specific language governing permissions and
14+
limitations under the License.
15+
*/
16+
17+
package common
18+
19+
// OptionFunc tunes how a per-engine Opts value is built.
20+
type OptionFunc func(*optionConfig)
21+
22+
type optionConfig struct {
23+
skipReadinessCheck bool
24+
}
25+
26+
// SkipReadinessCheck drops the `status.phase == Ready` gate on the database CR.
27+
// Commands that only read the CR to render manifests do not need a healthy
28+
// database, and refusing to run against a degraded source is actively unhelpful
29+
// for disaster-recovery workflows, where the source being unhealthy is the
30+
// whole reason the command is being run.
31+
func SkipReadinessCheck() OptionFunc {
32+
return func(cfg *optionConfig) {
33+
cfg.skipReadinessCheck = true
34+
}
35+
}
36+
37+
func buildOptionConfig(opts []OptionFunc) optionConfig {
38+
var cfg optionConfig
39+
for _, opt := range opts {
40+
opt(&cfg)
41+
}
42+
return cfg
43+
}

‎pkg/common/postgres.go‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -49,7 +49,9 @@ type PostgresOpts struct {
4949
ErrWriter *bytes.Buffer
5050
}
5151

52-
func NewPostgresOpts(f cmdutil.Factory, dbName, namespace string) (*PostgresOpts, error) {
52+
func NewPostgresOpts(f cmdutil.Factory, dbName, namespace string, options ...OptionFunc) (*PostgresOpts, error) {
53+
cfg := buildOptionConfig(options)
54+
5355
config, err := f.ToRESTConfig()
5456
if err != nil {
5557
return nil, err
@@ -77,7 +79,7 @@ func NewPostgresOpts(f cmdutil.Factory, dbName, namespace string) (*PostgresOpts
7779
return nil, err
7880
}
7981

80-
if db.Status.Phase != dbapi.DatabasePhaseReady {
82+
if !cfg.skipReadinessCheck && db.Status.Phase != dbapi.DatabasePhaseReady {
8183
return nil, fmt.Errorf("postgres %s/%s is not ready", namespace, dbName)
8284
}
8385

‎pkg/remote_replica/mysql.go‎

Lines changed: 72 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ import (
2121
"fmt"
2222
"log"
2323
"os"
24+
"strings"
2425
"time"
2526

2627
dbapi "kubedb.dev/apimachinery/apis/kubedb/v1"
@@ -31,9 +32,9 @@ import (
3132
core "k8s.io/api/core/v1"
3233
kerr "k8s.io/apimachinery/pkg/api/errors"
3334
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
34-
"k8s.io/apimachinery/pkg/labels"
3535
"k8s.io/apimachinery/pkg/util/sets"
3636
"k8s.io/apimachinery/pkg/util/wait"
37+
"k8s.io/klog/v2"
3738
cmdutil "k8s.io/kubectl/pkg/cmd/util"
3839
cm_util "kmodules.xyz/cert-manager-util/certmanager/v1"
3940
kutil "kmodules.xyz/client-go"
@@ -59,7 +60,7 @@ const (
5960

6061
func MysqlAPP(f cmdutil.Factory) *cobra.Command {
6162
var userName, password, dns, ns string
62-
var yes bool
63+
var yes, skipUserCreation bool
6364
cmd := cobra.Command{
6465
Use: "mysql",
6566
Short: desLong,
@@ -71,12 +72,16 @@ func MysqlAPP(f cmdutil.Factory) *cobra.Command {
7172
if len(args) == 0 {
7273
log.Fatal("no database name given")
7374
}
74-
if err := userPrompt(yes); err != nil {
75-
log.Fatal(err)
75+
// Nothing on the source is altered when user creation is skipped, so
76+
// the "password will be altered" confirmation has nothing to confirm.
77+
if !skipUserCreation {
78+
if err := userPrompt(yes); err != nil {
79+
log.Fatal(err)
80+
}
7681
}
7782
var buffer []byte
7883

79-
buffer, err := generateMySQLConfig(f, userName, password, dns, ns, args[0])
84+
buffer, err := generateMySQLConfig(f, userName, password, dns, ns, args[0], skipUserCreation)
8085
if err != nil {
8186
log.Fatal(err)
8287
}
@@ -111,13 +116,22 @@ func MysqlAPP(f cmdutil.Factory) *cobra.Command {
111116
log.Fatal(err)
112117
}
113118
cmd.PersistentFlags().BoolVarP(&yes, "yes", "y", false, "permission for alter password for the remote replica")
119+
cmd.PersistentFlags().BoolVar(&skipUserCreation, "skip-user-creation", false,
120+
"do not create or alter the replication user on the source; assume it already exists with the given password. "+
121+
"Needed when the source has no writable primary to exec into")
114122
return &cmd
115123
}
116124

117-
func generateMySQLConfig(f cmdutil.Factory, userName string, password string, dns string, ns string, dbname string) ([]byte, error) {
125+
func generateMySQLConfig(f cmdutil.Factory, userName string, password string, dns string, ns string, dbname string, skipUserCreation bool) ([]byte, error) {
118126
var buffer []byte
119127

120-
opts, err := common.NewMySQLOpts(f, dbname, ns)
128+
var dbOpts []common.OptionFunc
129+
if skipUserCreation {
130+
// Rendering manifests only reads the CR; a non-Ready source should not
131+
// block a disaster recovery config from being generated.
132+
dbOpts = append(dbOpts, common.SkipReadinessCheck())
133+
}
134+
opts, err := common.NewMySQLOpts(f, dbname, ns, dbOpts...)
121135
if err != nil {
122136
return nil, fmt.Errorf("failed to get db %s, err:%v", dbname, err)
123137
}
@@ -127,7 +141,7 @@ func generateMySQLConfig(f cmdutil.Factory, userName string, password string, dn
127141
return nil, fmt.Errorf("failed to get appbinding %v", err)
128142
}
129143

130-
authBuff, authSecretName, err := generateMySQLAuthSecret(userName, password, ns, opts)
144+
authBuff, authSecretName, err := generateMySQLAuthSecret(userName, password, ns, skipUserCreation, opts)
131145
if err != nil {
132146
return nil, fmt.Errorf("failed to generate auth secret ,%v", err)
133147
}
@@ -198,15 +212,20 @@ func generateMySQLTlsSecret(userName string, apb *appApi.AppBinding, ns string,
198212
return buffer, tlsSecret.Name, nil
199213
}
200214

201-
func generateMySQLAuthSecret(userName string, password string, ns string, opts *common.MySQLOpts) ([]byte, string, error) {
202-
if userName != opts.Username {
215+
func generateMySQLAuthSecret(userName string, password string, ns string, skipUserCreation bool, opts *common.MySQLOpts) ([]byte, string, error) {
216+
switch {
217+
case userName == opts.Username:
218+
password = opts.Pass
219+
case skipUserCreation:
220+
if err := verifyMySQLUser(opts, userName); err != nil {
221+
return nil, "", err
222+
}
223+
default:
203224
// generate user if not present
204225
err := generateMySQLUser(opts, userName, password)
205226
if err != nil {
206227
return nil, "", fmt.Errorf("failed to generate user err:%v", err)
207228
}
208-
} else {
209-
password = opts.Pass
210229
}
211230
// generate auth secret
212231
AuthSecret := core.Secret{
@@ -235,17 +254,48 @@ func generateMySQLAuthSecret(userName string, password string, ns string, opts *
235254
return buffer, AuthSecret.Name, nil
236255
}
237256

238-
func generateMySQLUser(opts *common.MySQLOpts, name string, password string) error {
239-
label := opts.DB.OffshootLabels()
240-
if *opts.DB.Spec.Replicas > 1 {
241-
label["kubedb.com/role"] = "primary"
257+
// verifyMySQLUser checks the replication user without writing anything. Reading
258+
// mysql.user works on a read replica, so any running member will do. An
259+
// unreachable database downgrades to a warning: generating the config is still
260+
// the useful outcome.
261+
func verifyMySQLUser(opts *common.MySQLOpts, name string) error {
262+
pod, err := pickDBPod(opts.Client, opts.DB.Namespace, opts.DB.OffshootLabels(), MySQLContainerName, false)
263+
if err != nil {
264+
klog.Warningf("skipping verification of user %q: %v", name, err)
265+
return nil
242266
}
243267

244-
pods, err := opts.Client.CoreV1().Pods(opts.DB.Namespace).List(context.TODO(), metav1.ListOptions{
245-
LabelSelector: labels.Set.String(label),
246-
})
247-
if err != nil || len(pods.Items) == 0 {
248-
return err
268+
query := fmt.Sprintf("export MYSQL_PWD='%s' && mysql -uroot -N -B -e \"SELECT Repl_slave_priv FROM mysql.user WHERE user='%s'\"", opts.Pass, name)
269+
out, err := exec_util.ExecIntoPod(opts.Config, pod,
270+
exec_util.Command("bash", "-c", query),
271+
exec_util.Container(MySQLContainerName),
272+
)
273+
if err != nil {
274+
klog.Warningf("skipping verification of user %q: failed to query %s: %v", name, pod.Name, err)
275+
return nil
276+
}
277+
278+
switch strings.TrimSpace(out) {
279+
case "":
280+
return fmt.Errorf("user %q does not exist on %s/%s; create it on the writable primary first:\n "+
281+
"CREATE USER %s IDENTIFIED BY '<password>'; GRANT REPLICATION SLAVE, CLONE_ADMIN, BACKUP_ADMIN ON *.* TO '%s'@'%%';",
282+
name, opts.DB.Namespace, opts.DB.Name, name, name)
283+
case "N":
284+
return fmt.Errorf("user %q exists on %s/%s but lacks REPLICATION SLAVE; run on the writable primary:\n "+
285+
"GRANT REPLICATION SLAVE, CLONE_ADMIN, BACKUP_ADMIN ON *.* TO '%s'@'%%';",
286+
name, opts.DB.Namespace, opts.DB.Name, name)
287+
}
288+
289+
fmt.Printf("user %q verified on %s (replication granted); leaving the source catalog untouched\n", name, pod.Name)
290+
return nil
291+
}
292+
293+
func generateMySQLUser(opts *common.MySQLOpts, name string, password string) error {
294+
// DDL only: a clustered source must be addressed through its primary.
295+
pod, err := pickDBPod(opts.Client, opts.DB.Namespace, opts.DB.OffshootLabels(), MySQLContainerName, *opts.DB.Spec.Replicas > 1)
296+
if err != nil {
297+
return fmt.Errorf("%v; pass --skip-user-creation to generate the config against a standby "+
298+
"whose replication user already exists", err)
249299
}
250300
query := fmt.Sprintf("export MYSQL_PWD='%s' && mysql -uroot -e \"create user if not exists %s; alter user %s identified by '%s';"+
251301
"GRANT REPLICATION SLAVE, CLONE_ADMIN, BACKUP_ADMIN ON *.* TO '%s'@'%%' WITH GRANT OPTION; \"", opts.Pass,
@@ -257,7 +307,7 @@ func generateMySQLUser(opts *common.MySQLOpts, name string, password string) err
257307
container,
258308
}
259309

260-
_, err = exec_util.ExecIntoPod(opts.Config, &pods.Items[0], options...)
310+
_, err = exec_util.ExecIntoPod(opts.Config, pod, options...)
261311
if err != nil {
262312
return err
263313
}

‎pkg/remote_replica/pods.go‎

Lines changed: 98 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,98 @@
1+
/*
2+
Copyright AppsCode Inc. and Contributors
3+
4+
Licensed under the AppsCode Community License 1.0.0 (the "License");
5+
you may not use this file except in compliance with the License.
6+
You may obtain a copy of the License at
7+
8+
https://github.com/appscode/licenses/raw/1.0.0/AppsCode-Community-1.0.0.md
9+
10+
Unless required by applicable law or agreed to in writing, software
11+
distributed under the License is distributed on an "AS IS" BASIS,
12+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
See the License for the specific language governing permissions and
14+
limitations under the License.
15+
*/
16+
17+
package remote_replica
18+
19+
import (
20+
"context"
21+
"fmt"
22+
23+
"kubedb.dev/apimachinery/apis/kubedb"
24+
25+
core "k8s.io/api/core/v1"
26+
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
27+
"k8s.io/apimachinery/pkg/labels"
28+
"k8s.io/client-go/kubernetes"
29+
)
30+
31+
// Containers holding the database server itself, the ones a client query has to
32+
// be exec-ed into.
33+
const (
34+
PostgresContainerName = "postgres"
35+
MySQLContainerName = "mysql"
36+
)
37+
38+
// pickDBPod returns a running database pod to exec into.
39+
//
40+
// primaryOnly is set when the caller intends to run DDL. CREATE ROLE / ALTER
41+
// ROLE / GRANT write to the shared catalog, and a hot standby rejects them with
42+
// "cannot execute ... in a read-only transaction", so only a pod labelled
43+
// kubedb.com/role=primary qualifies. Read-only catalog queries run fine on any
44+
// member, so verification passes primaryOnly=false and merely prefers the
45+
// primary when one happens to be labelled.
46+
func pickDBPod(client kubernetes.Interface, ns string, selector map[string]string, container string, primaryOnly bool) (*core.Pod, error) {
47+
pods, err := client.CoreV1().Pods(ns).List(context.TODO(), metav1.ListOptions{
48+
LabelSelector: labels.Set(selector).String(),
49+
})
50+
if err != nil {
51+
return nil, err
52+
}
53+
54+
pod, err := selectPod(pods.Items, container, primaryOnly)
55+
if err != nil {
56+
return nil, fmt.Errorf("in namespace %s: %v", ns, err)
57+
}
58+
return pod, nil
59+
}
60+
61+
// selectPod holds the choice itself, kept free of any client so it can be
62+
// exercised directly. Only pods that actually carry `container` are eligible: a
63+
// KubeDB cluster runs helper pods under the same offshoot labels - the Postgres
64+
// arbiter, for one, runs pg-coordinator alone - and exec-ing a database query
65+
// into those fails in a way that is hard to read.
66+
func selectPod(pods []core.Pod, container string, primaryOnly bool) (*core.Pod, error) {
67+
var usable []core.Pod
68+
for i := range pods {
69+
if pods[i].Status.Phase == core.PodRunning && hasRunningContainer(&pods[i], container) {
70+
usable = append(usable, pods[i])
71+
}
72+
}
73+
if len(usable) == 0 {
74+
// Reported explicitly: an empty list used to be swallowed and treated as
75+
// success, which produced a config for a user that was never created.
76+
return nil, fmt.Errorf("none of the %d database pod(s) has a running %q container", len(pods), container)
77+
}
78+
79+
for i := range usable {
80+
if usable[i].Labels[kubedb.LabelRole] == kubedb.DatabasePodPrimary {
81+
return &usable[i], nil
82+
}
83+
}
84+
if primaryOnly {
85+
return nil, fmt.Errorf("no pod is labelled %s=%s, so there is no writable primary to run DDL against",
86+
kubedb.LabelRole, kubedb.DatabasePodPrimary)
87+
}
88+
return &usable[0], nil
89+
}
90+
91+
func hasRunningContainer(pod *core.Pod, container string) bool {
92+
for _, cs := range pod.Status.ContainerStatuses {
93+
if cs.Name == container {
94+
return cs.State.Running != nil
95+
}
96+
}
97+
return false
98+
}

0 commit comments

Comments
 (0)