forked from rook/rook
This fixes a comment typo caught bt the codespell CI check. Signed-off-by: Michael Adam <obnox@samba.org>
1440 lines
54 KiB
Go
1440 lines
54 KiB
Go
/*
|
|
Copyright 2016 The Rook Authors. All rights reserved.
|
|
|
|
Licensed under the Apache License, Version 2.0 (the "License");
|
|
you may not use this file except in compliance with the License.
|
|
You may obtain a copy of the License at
|
|
|
|
http://www.apache.org/licenses/LICENSE-2.0
|
|
|
|
Unless required by applicable law or agreed to in writing, software
|
|
distributed under the License is distributed on an "AS IS" BASIS,
|
|
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
See the License for the specific language governing permissions and
|
|
limitations under the License.
|
|
*/
|
|
|
|
package object
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"os"
|
|
"reflect"
|
|
"sort"
|
|
"strconv"
|
|
"strings"
|
|
"syscall"
|
|
|
|
"github.com/pkg/errors"
|
|
cephv1 "github.com/rook/rook/pkg/apis/ceph.rook.io/v1"
|
|
"github.com/rook/rook/pkg/clusterd"
|
|
cephclient "github.com/rook/rook/pkg/daemon/ceph/client"
|
|
cephver "github.com/rook/rook/pkg/operator/ceph/version"
|
|
"github.com/rook/rook/pkg/operator/k8sutil"
|
|
"github.com/rook/rook/pkg/util"
|
|
"github.com/rook/rook/pkg/util/exec"
|
|
"github.com/rook/rook/pkg/util/log"
|
|
"golang.org/x/sync/errgroup"
|
|
v1 "k8s.io/api/core/v1"
|
|
kerrors "k8s.io/apimachinery/pkg/api/errors"
|
|
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
|
"k8s.io/apimachinery/pkg/types"
|
|
"k8s.io/apimachinery/pkg/util/sets"
|
|
validation "k8s.io/apimachinery/pkg/util/validation"
|
|
"sigs.k8s.io/controller-runtime/pkg/client"
|
|
)
|
|
|
|
const (
|
|
rootPool = ".rgw.root"
|
|
|
|
// AppName is the name Rook uses for the object store's application
|
|
AppName = "rook-ceph-rgw"
|
|
bucketProvisionerName = "ceph.rook.io/bucket"
|
|
AccessKeyName = "access-key"
|
|
SecretKeyName = "secret-key"
|
|
svcDNSSuffix = "svc"
|
|
rgwRadosPoolPgNum = "8"
|
|
rgwApplication = "rgw"
|
|
)
|
|
|
|
var (
|
|
metadataPools = []string{
|
|
// .rgw.root (rootPool) is appended to this slice where needed
|
|
"rgw.control",
|
|
"rgw.meta",
|
|
"rgw.log",
|
|
"rgw.buckets.index",
|
|
"rgw.buckets.non-ec",
|
|
"rgw.otp",
|
|
}
|
|
dataPoolName = "rgw.buckets.data"
|
|
|
|
// An user with system privileges for dashboard service
|
|
DashboardUser = "dashboard-admin"
|
|
|
|
zonePoolNSSuffix = map[string]string{
|
|
"domain_root": ".meta.root",
|
|
"control_pool": ".control",
|
|
"gc_pool": ".log.gc",
|
|
"lc_pool": ".log.lc",
|
|
"log_pool": ".log",
|
|
"intent_log_pool": ".log.intent",
|
|
"usage_log_pool": ".log.usage",
|
|
"roles_pool": ".meta.roles",
|
|
"reshard_pool": ".log.reshard",
|
|
"user_keys_pool": ".meta.users.keys",
|
|
"user_email_pool": ".meta.users.email",
|
|
"user_swift_pool": ".meta.users.swift",
|
|
"user_uid_pool": ".meta.users.uid",
|
|
"otp_pool": ".otp",
|
|
"notif_pool": ".log.notif",
|
|
"topics_pool": ".meta.topics", // introduced in Ceph v19
|
|
"account_pool": ".meta.account", // introduced in Ceph v19
|
|
"group_pool": ".meta.group", // introduced in Ceph v19
|
|
"restore_pool": ".log.restore", // introduced in Ceph v20
|
|
"bucket_logging_pool": ".log.bucket-logging", // introduced in Ceph v20
|
|
"dedup_pool": ".dedup", // introduced in Ceph v20
|
|
}
|
|
)
|
|
|
|
type idType struct {
|
|
ID string `json:"id"`
|
|
}
|
|
|
|
type zoneGroupType struct {
|
|
MasterZoneID string `json:"master_zone"`
|
|
IsMaster bool `json:"is_master"`
|
|
Zones []zoneType `json:"zones"`
|
|
Endpoints []string `json:"endpoints"`
|
|
}
|
|
|
|
type zoneType struct {
|
|
Name string `json:"name"`
|
|
Endpoints []string `json:"endpoints"`
|
|
}
|
|
|
|
type realmType struct {
|
|
Realms []string `json:"realms"`
|
|
}
|
|
|
|
// allow commitConfigChanges to be overridden for unit testing
|
|
var commitConfigChanges = CommitConfigChanges
|
|
|
|
func deleteRealmAndPools(objContext *Context, spec cephv1.ObjectStoreSpec) error {
|
|
if spec.IsMultisite() {
|
|
// since pools for object store are created by the zone, the object store only needs to be removed from the zone
|
|
err := removeObjectStoreFromMultisite(objContext, spec)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
return deleteSingleSiteRealmAndPools(objContext, spec)
|
|
}
|
|
|
|
func removeObjectStoreFromMultisite(objContext *Context, spec cephv1.ObjectStoreSpec) error {
|
|
// get list of endpoints not including the endpoint of the object-store for the zone
|
|
zoneEndpointsList, isEndpointAlreadyExists, err := getZoneEndpoints(objContext, objContext.Endpoint)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// The endpoint is present in zone, hence remove it
|
|
if isEndpointAlreadyExists {
|
|
realmArg := fmt.Sprintf("--rgw-realm=%s", objContext.Realm)
|
|
zoneGroupArg := fmt.Sprintf("--rgw-zonegroup=%s", objContext.ZoneGroup)
|
|
zoneEndpoints := strings.Join(zoneEndpointsList, ",")
|
|
endpointArg := fmt.Sprintf("--endpoints=%s", zoneEndpoints)
|
|
|
|
zoneIsMaster, err := CheckZoneIsMaster(objContext)
|
|
if err != nil {
|
|
return errors.Wrapf(err, "failed to determine if zone %q is master", objContext.Zone)
|
|
}
|
|
|
|
zoneGroupIsMaster := false
|
|
if zoneIsMaster {
|
|
_, err = RunAdminCommandNoMultisite(objContext, false, "zonegroup", "modify", realmArg, zoneGroupArg, endpointArg)
|
|
if err != nil {
|
|
if kerrors.IsNotFound(err) {
|
|
return err
|
|
}
|
|
return errors.Wrapf(err, "failed to remove object store %q endpoint from rgw zone group %q", objContext.Name, objContext.ZoneGroup)
|
|
}
|
|
log.NamedDebug(objContext.NsName(), logger, "endpoint %q was removed from zone group %q. the remaining endpoints in the zone group are %q", objContext.Endpoint, objContext.ZoneGroup, zoneEndpoints)
|
|
|
|
// check if zone group is master only if zone is master for creating the system user
|
|
zoneGroupIsMaster, err = checkZoneGroupIsMaster(objContext)
|
|
if err != nil {
|
|
return errors.Wrapf(err, "failed to find out whether zone group %q is the master zone group", objContext.ZoneGroup)
|
|
}
|
|
}
|
|
|
|
_, err = runAdminCommand(objContext, false, "zone", "modify", endpointArg)
|
|
if err != nil {
|
|
return errors.Wrapf(err, "failed to remove object store %q endpoint from rgw zone %q", objContext.Name, spec.Zone.Name)
|
|
}
|
|
log.NamedDebug(objContext.NsName(), logger, "endpoint %q was removed from zone %q. the remaining endpoints in the zone are %q", objContext.Endpoint, objContext.Zone, zoneEndpoints)
|
|
|
|
if zoneIsMaster && zoneGroupIsMaster && zoneEndpoints == "" {
|
|
log.NamedWarning(objContext.NsName(), logger, "WARNING: No other zone in realm %q can commit to the period or pull the realm until you create another object-store in zone %q", objContext.Realm, objContext.Zone)
|
|
}
|
|
|
|
// this will notify other zones of changes if there are multi-zones
|
|
if err := commitConfigChanges(objContext); err != nil {
|
|
return errors.Wrapf(err, "failed to commit config changes after removing CephObjectStore %q from multi-site", objContext.NsName())
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func deleteSingleSiteRealmAndPools(objContext *Context, spec cephv1.ObjectStoreSpec) error {
|
|
stores, err := getObjectStores(objContext)
|
|
if err != nil {
|
|
return errors.Wrap(err, "failed to detect object stores during deletion")
|
|
}
|
|
if len(stores) == 0 {
|
|
log.NamedInfo(objContext.NsName(), logger, "did not find object store %q, nothing to delete", objContext.Name)
|
|
return nil
|
|
}
|
|
log.NamedInfo(objContext.NsName(), logger, "Found stores %v when deleting store %s", stores, objContext.Name)
|
|
|
|
err = deleteRealm(objContext)
|
|
if err != nil {
|
|
return errors.Wrap(err, "failed to delete realm")
|
|
}
|
|
|
|
lastStore := false
|
|
if len(stores) == 1 && stores[0] == objContext.Name {
|
|
lastStore = true
|
|
}
|
|
|
|
if !spec.PreservePoolsOnDelete {
|
|
if EmptyPool(spec.DataPool) && EmptyPool(spec.MetadataPool) {
|
|
log.NamedInfo(objContext.NsName(), logger, "skipping removal of pools since not specified in the object store")
|
|
return nil
|
|
}
|
|
err = DeletePools(objContext, lastStore, objContext.Name)
|
|
if err != nil {
|
|
return errors.Wrap(err, "failed to delete object store pools")
|
|
}
|
|
} else {
|
|
log.NamedInfo(objContext.NsName(), logger, "PreservePoolsOnDelete is set in object store %s. Pools not deleted", objContext.Name)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// This is used for quickly getting the name of the realm, zone group, and zone for an object-store to pass into a Context
|
|
func getMultisiteForObjectStore(ctx context.Context, clusterdContext *clusterd.Context, spec *cephv1.ObjectStoreSpec, namespace, name string) (string, string, string, error) {
|
|
if spec.IsExternal() {
|
|
// In https://github.com/rook/rook/issues/6342, it was determined that
|
|
// a multisite context isn't needed for external mode CephObjectStores.
|
|
// The context is only needed for managing an object store, which isn't
|
|
// happening in external mode.
|
|
return "", "default", "default", nil
|
|
}
|
|
if spec.IsMultisite() {
|
|
zone, err := clusterdContext.RookClientset.CephV1().CephObjectZones(namespace).Get(ctx, spec.Zone.Name, metav1.GetOptions{})
|
|
if err != nil {
|
|
return "", "", "", errors.Wrapf(err, "failed to find zone for object-store %q", name)
|
|
}
|
|
|
|
zonegroup, err := clusterdContext.RookClientset.CephV1().CephObjectZoneGroups(namespace).Get(ctx, zone.Spec.ZoneGroup, metav1.GetOptions{})
|
|
if err != nil {
|
|
return "", "", "", errors.Wrapf(err, "failed to find zone group for object-store %q", name)
|
|
}
|
|
|
|
realm, err := clusterdContext.RookClientset.CephV1().CephObjectRealms(namespace).Get(ctx, zonegroup.Spec.Realm, metav1.GetOptions{})
|
|
if err != nil {
|
|
return "", "", "", errors.Wrapf(err, "failed to find realm for object-store %q", name)
|
|
}
|
|
|
|
return realm.Name, zonegroup.Name, zone.Name, nil
|
|
}
|
|
|
|
return name, name, name, nil
|
|
}
|
|
|
|
func CheckZoneIsMaster(objContext *Context) (bool, error) {
|
|
log.NamedDebug(objContext.NsName(), logger, "checking if zone %v is the master zone", objContext.Zone)
|
|
realmArg := fmt.Sprintf("--rgw-realm=%s", objContext.Realm)
|
|
zoneGroupArg := fmt.Sprintf("--rgw-zonegroup=%s", objContext.ZoneGroup)
|
|
zoneArg := fmt.Sprintf("--rgw-zone=%s", objContext.Zone)
|
|
|
|
zoneGroupJson, err := RunAdminCommandNoMultisite(objContext, true, "zonegroup", "get", realmArg, zoneGroupArg)
|
|
if err != nil {
|
|
// This handles the case where the pod we use to exec command (act as a proxy) is not found/ready yet
|
|
// The caller can nicely handle the error and not overflow the op logs with misleading error messages
|
|
if kerrors.IsNotFound(err) {
|
|
return false, err
|
|
}
|
|
return false, errors.Wrap(err, "failed to get rgw zone group")
|
|
}
|
|
zoneGroupOutput, err := DecodeZoneGroupConfig(zoneGroupJson)
|
|
if err != nil {
|
|
return false, errors.Wrap(err, "failed to parse zonegroup get json")
|
|
}
|
|
log.NamedDebug(objContext.NsName(), logger, "got master zone ID for zone group %v", objContext.ZoneGroup)
|
|
|
|
zoneOutput, err := RunAdminCommandNoMultisite(objContext, true, "zone", "get", realmArg, zoneGroupArg, zoneArg)
|
|
if err != nil {
|
|
// This handles the case where the pod we use to exec command (act as a proxy) is not found/ready yet
|
|
// The caller can nicely handle the error and not overflow the op logs with misleading error messages
|
|
if kerrors.IsNotFound(err) {
|
|
return false, err
|
|
}
|
|
return false, errors.Wrap(err, "failed to get rgw zone")
|
|
}
|
|
zoneID, err := decodeID(zoneOutput)
|
|
if err != nil {
|
|
return false, errors.Wrap(err, "failed to parse zone id")
|
|
}
|
|
log.NamedDebug(objContext.NsName(), logger, "got zone ID for zone %v", objContext.Zone)
|
|
|
|
if zoneID == zoneGroupOutput.MasterZoneID {
|
|
log.NamedDebug(objContext.NsName(), logger, "zone is master")
|
|
return true, nil
|
|
}
|
|
|
|
log.NamedDebug(objContext.NsName(), logger, "zone is not master")
|
|
return false, nil
|
|
}
|
|
|
|
func checkZoneGroupIsMaster(objContext *Context) (bool, error) {
|
|
log.NamedDebug(objContext.NsName(), logger, "checking if zone group %v is the master zone group", objContext.ZoneGroup)
|
|
realmArg := fmt.Sprintf("--rgw-realm=%s", objContext.Realm)
|
|
zoneGroupArg := fmt.Sprintf("--rgw-zonegroup=%s", objContext.ZoneGroup)
|
|
|
|
zoneGroupOutput, err := RunAdminCommandNoMultisite(objContext, true, "zonegroup", "get", realmArg, zoneGroupArg)
|
|
if err != nil {
|
|
// This handles the case where the pod we use to exec command (act as a proxy) is not found/ready yet
|
|
// The caller can nicely handle the error and not overflow the op logs with misleading error messages
|
|
if kerrors.IsNotFound(err) {
|
|
return false, err
|
|
}
|
|
return false, errors.Wrap(err, "failed to get rgw zone group")
|
|
}
|
|
|
|
zoneGroupJson, err := DecodeZoneGroupConfig(zoneGroupOutput)
|
|
if err != nil {
|
|
return false, errors.Wrap(err, "failed to parse master zone id")
|
|
}
|
|
|
|
return zoneGroupJson.IsMaster, nil
|
|
}
|
|
|
|
func DecodeSecret(secret *v1.Secret, keyName string) (string, error) {
|
|
realmKey, ok := secret.Data[keyName]
|
|
|
|
if !ok {
|
|
return "", fmt.Errorf("failed to find key %q in secret %q data. user likely created or modified the secret manually and should add the missing key back into the secret", keyName, secret.Name)
|
|
}
|
|
|
|
return string(realmKey), nil
|
|
}
|
|
|
|
func GetRealmKeySecret(ctx context.Context, clusterdContext *clusterd.Context, realmName types.NamespacedName) (*v1.Secret, error) {
|
|
realmSecretName := realmName.Name + "-keys"
|
|
realmSecret, err := clusterdContext.Clientset.CoreV1().Secrets(realmName.Namespace).Get(ctx, realmSecretName, metav1.GetOptions{})
|
|
if err != nil {
|
|
return nil, errors.Wrapf(err, "failed to get CephObjectRealm %q keys secret", realmName.String())
|
|
}
|
|
log.NamedDebug(realmName, logger, "found keys secret for CephObjectRealm %q", realmName.String())
|
|
return realmSecret, nil
|
|
}
|
|
|
|
func GetRealmKeyArgsFromSecret(realmSecret *v1.Secret, realmName types.NamespacedName) (string, string, error) {
|
|
accessKey, err := DecodeSecret(realmSecret, AccessKeyName)
|
|
if err != nil {
|
|
return "", "", errors.Wrapf(err, "failed to decode CephObjectRealm %q access key from secret %q", realmName.String(), realmSecret.Name)
|
|
}
|
|
secretKey, err := DecodeSecret(realmSecret, SecretKeyName)
|
|
if err != nil {
|
|
return "", "", errors.Wrapf(err, "failed to decode CephObjectRealm %q secret key from secret %q", realmName.String(), realmSecret.Name)
|
|
}
|
|
log.NamedDebug(realmName, logger, "decoded keys for realm %q", realmName.String())
|
|
|
|
accessKeyArg := fmt.Sprintf("--access-key=%s", accessKey)
|
|
secretKeyArg := fmt.Sprintf("--secret-key=%s", secretKey)
|
|
|
|
return accessKeyArg, secretKeyArg, nil
|
|
}
|
|
|
|
func GetRealmKeyArgs(ctx context.Context, clusterdContext *clusterd.Context, realmName, namespace string) (string, string, error) {
|
|
nsName := types.NamespacedName{Namespace: namespace, Name: realmName}
|
|
realmNsName := types.NamespacedName{Namespace: namespace, Name: realmName}
|
|
log.NamedDebug(nsName, logger, "getting keys for realm %q", realmNsName.String())
|
|
|
|
secret, err := GetRealmKeySecret(ctx, clusterdContext, realmNsName)
|
|
if err != nil {
|
|
return "", "", err
|
|
}
|
|
|
|
return GetRealmKeyArgsFromSecret(secret, realmNsName)
|
|
}
|
|
|
|
func getZoneEndpoints(objContext *Context, serviceEndpoint string) ([]string, bool, error) {
|
|
log.NamedDebug(objContext.NsName(), logger, "getting current endpoints for zone %v", objContext.Zone)
|
|
realmArg := fmt.Sprintf("--rgw-realm=%s", objContext.Realm)
|
|
zoneGroupArg := fmt.Sprintf("--rgw-zonegroup=%s", objContext.ZoneGroup)
|
|
isEndpointAlreadyExists := false
|
|
|
|
zoneGroupOutput, err := RunAdminCommandNoMultisite(objContext, true, "zonegroup", "get", realmArg, zoneGroupArg)
|
|
if err != nil {
|
|
// This handles the case where the pod we use to exec command (act as a proxy) is not found/ready yet
|
|
// The caller can nicely handle the error and not overflow the op logs with misleading error messages
|
|
return []string{}, isEndpointAlreadyExists, errorOrIsNotFound(err, "failed to get rgw zone group %q", objContext.Name)
|
|
}
|
|
zoneGroupJson, err := DecodeZoneGroupConfig(zoneGroupOutput)
|
|
if err != nil {
|
|
return []string{}, isEndpointAlreadyExists, errors.Wrap(err, "failed to parse zones list")
|
|
}
|
|
|
|
zoneEndpointsList := []string{}
|
|
for _, zone := range zoneGroupJson.Zones {
|
|
if zone.Name == objContext.Zone {
|
|
for _, endpoint := range zone.Endpoints {
|
|
// in case object-store operator code is rereconciled, zone modify could get run again with serviceEndpoint added again
|
|
if endpoint != serviceEndpoint {
|
|
zoneEndpointsList = append(zoneEndpointsList, endpoint)
|
|
} else {
|
|
isEndpointAlreadyExists = true
|
|
}
|
|
}
|
|
break
|
|
}
|
|
}
|
|
|
|
return zoneEndpointsList, isEndpointAlreadyExists, nil
|
|
}
|
|
|
|
func createMultisiteConfigurations(objContext *Context, store *cephv1.CephObjectStore, configType, configTypeArg string, args ...string) error {
|
|
args = append([]string{configType}, args...)
|
|
args = append(args, configTypeArg)
|
|
// get the multisite config before creating
|
|
configTypeArgs := []string{configType, "get", configTypeArg}
|
|
if configType == "zonegroup" {
|
|
configTypeArgs = append(configTypeArgs, fmt.Sprintf("--rgw-realm=%s", objContext.Realm))
|
|
}
|
|
output, getConfigErr := RunAdminCommandNoMultisite(objContext, true, configTypeArgs...)
|
|
if getConfigErr == nil {
|
|
return nil
|
|
}
|
|
|
|
if kerrors.IsNotFound(getConfigErr) {
|
|
// the pod used to exec command (act as a proxy) is not found/ready yet
|
|
// caller can nicely handle error and not overflow logs with misleading error messages
|
|
return getConfigErr
|
|
}
|
|
|
|
code, err := exec.ExtractExitCode(getConfigErr)
|
|
if err != nil {
|
|
return errorOrIsNotFound(getConfigErr, "'radosgw-admin %q get' failed with code %q, for reason %q, error: (%v)", configType, strconv.Itoa(code), output, string(kerrors.ReasonForError(err)))
|
|
}
|
|
// ENOENT means "No such file or directory"
|
|
if code != int(syscall.ENOENT) {
|
|
code := strconv.Itoa(code)
|
|
return errors.Wrapf(getConfigErr, "'radosgw-admin %q get' failed with code %q, for reason %q", configType, code, output)
|
|
}
|
|
|
|
// create the object if it doesn't exist yet
|
|
if store.Spec.DefaultRealm && objContext.clusterInfo.CephVersion.IsAtLeast(cephver.Squid) {
|
|
log.NamedInfo(objContext.NsName(), logger, "marking object store %q as default realm", store.Namespace+"/"+store.Name)
|
|
args = append(args, "--default")
|
|
}
|
|
output, err = RunAdminCommandNoMultisite(objContext, false, args...)
|
|
if err != nil {
|
|
return errorOrIsNotFound(err, "failed to create ceph %q %q, for reason %q", configType, configTypeArg, output)
|
|
}
|
|
log.NamedDebug(objContext.NsName(), logger, "created %q %q", configType, configTypeArg)
|
|
|
|
return nil
|
|
}
|
|
|
|
func createNonMultisiteStore(objContext *Context, endpointArg string, store *cephv1.CephObjectStore) error {
|
|
log.NamedDebug(objContext.NsName(), logger, "creating realm, zone group, zone for object-store %v", objContext.Name)
|
|
|
|
realmArg := fmt.Sprintf("--rgw-realm=%s", objContext.Realm)
|
|
zoneGroupArg := fmt.Sprintf("--rgw-zonegroup=%s", objContext.ZoneGroup)
|
|
zoneArg := fmt.Sprintf("--rgw-zone=%s", objContext.Zone)
|
|
|
|
err := createMultisiteConfigurations(objContext, store, "realm", realmArg, "create")
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
err = createMultisiteConfigurations(objContext, store, "zonegroup", zoneGroupArg, "create", "--master", realmArg, endpointArg)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
err = createMultisiteConfigurations(objContext, store, "zone", zoneArg, "create", "--master", endpointArg, realmArg, zoneGroupArg)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
log.NamedInfo(objContext.NsName(), logger, "Object store %q: realm=%s, zonegroup=%s, zone=%s", objContext.Name, objContext.Realm, objContext.ZoneGroup, objContext.Zone)
|
|
|
|
// Configure the zone for RADOS namespaces
|
|
err = ConfigureSharedPoolsForZone(objContext, store.Spec.SharedPools)
|
|
if err != nil {
|
|
return errors.Wrapf(err, "failed to configure rados namespaces for zone")
|
|
}
|
|
|
|
if err := commitConfigChanges(objContext); err != nil {
|
|
return errors.Wrapf(err, "failed to commit config changes after creating multisite config for CephObjectStore %q", objContext.NsName())
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func JoinMultisite(objContext *Context, endpointArg, zoneEndpoints, namespace string) error {
|
|
log.NamedDebug(objContext.NsName(), logger, "joining zone %v", objContext.Zone)
|
|
realmArg := fmt.Sprintf("--rgw-realm=%s", objContext.Realm)
|
|
zoneGroupArg := fmt.Sprintf("--rgw-zonegroup=%s", objContext.ZoneGroup)
|
|
zoneArg := fmt.Sprintf("--rgw-zone=%s", objContext.Zone)
|
|
|
|
zoneIsMaster, err := CheckZoneIsMaster(objContext)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
zoneGroupIsMaster := false
|
|
|
|
if zoneIsMaster {
|
|
// endpoints that are part of a master zone are supposed to be the endpoints for a zone group
|
|
_, err := RunAdminCommandNoMultisite(objContext, false, "zonegroup", "modify", realmArg, zoneGroupArg, endpointArg)
|
|
if err != nil {
|
|
return errorOrIsNotFound(err, "failed to add object store %q in rgw zone group %q", objContext.Name, objContext.ZoneGroup)
|
|
}
|
|
log.NamedDebug(objContext.NsName(), logger, "endpoints for zonegroup %q are now %q", objContext.ZoneGroup, zoneEndpoints)
|
|
|
|
// check if zone group is master only if zone is master for creating the system user
|
|
zoneGroupIsMaster, err = checkZoneGroupIsMaster(objContext)
|
|
if err != nil {
|
|
return errors.Wrapf(err, "failed to find out whether zone group %q in is the master zone group", objContext.ZoneGroup)
|
|
}
|
|
}
|
|
_, err = RunAdminCommandNoMultisite(objContext, false, "zone", "modify", realmArg, zoneGroupArg, zoneArg, endpointArg)
|
|
if err != nil {
|
|
return errorOrIsNotFound(err, "failed to add object store %q in rgw zone %q", objContext.Name, objContext.Zone)
|
|
}
|
|
log.NamedDebug(objContext.NsName(), logger, "endpoints for zone %q are now %q", objContext.Zone, zoneEndpoints)
|
|
|
|
if err := commitConfigChanges(objContext); err != nil {
|
|
return errors.Wrapf(err, "failed to commit config changes for CephObjectStore %q when joining multisite ", objContext.NsName())
|
|
}
|
|
|
|
log.NamedInfo(objContext.NsName(), logger, "added object store %q to realm %q, zonegroup %q, zone %q", objContext.Name, objContext.Realm, objContext.ZoneGroup, objContext.Zone)
|
|
|
|
// create system user for realm for master zone in master zonegroup for multisite scenario
|
|
if zoneIsMaster && zoneGroupIsMaster {
|
|
err = createSystemUser(objContext, namespace)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func createSystemUser(objContext *Context, namespace string) error {
|
|
uid := objContext.Realm + "-system-user"
|
|
uidArg := fmt.Sprintf("--uid=%s", uid)
|
|
realmArg := fmt.Sprintf("--rgw-realm=%s", objContext.Realm)
|
|
zoneGroupArg := fmt.Sprintf("--rgw-zonegroup=%s", objContext.ZoneGroup)
|
|
zoneArg := fmt.Sprintf("--rgw-zone=%s", objContext.Zone)
|
|
|
|
output, err := RunAdminCommandNoMultisite(objContext, false, "user", "info", uidArg, realmArg, zoneGroupArg, zoneArg)
|
|
if err == nil {
|
|
log.NamedDebug(objContext.NsName(), logger, "realm system user %q has already been created", uid)
|
|
return nil
|
|
}
|
|
|
|
if code, ok := exec.ExitStatus(err); ok && code == int(syscall.EINVAL) {
|
|
log.NamedDebug(objContext.NsName(), logger, "realm system user %q not found, running `radosgw-admin user create`", uid)
|
|
accessKeyArg, secretKeyArg, err := GetRealmKeyArgs(objContext.clusterInfo.Context, objContext.Context, objContext.Realm, namespace)
|
|
if err != nil {
|
|
return errors.Wrap(err, "failed to get keys for realm")
|
|
}
|
|
log.NamedDebug(objContext.NsName(), logger, "found keys to create realm system user %v", uid)
|
|
systemArg := "--system"
|
|
displayNameArg := fmt.Sprintf("--display-name=%s.user", objContext.Realm)
|
|
output, err = RunAdminCommandNoMultisite(objContext, false, "user", "create", realmArg, zoneGroupArg, zoneArg, uidArg, displayNameArg, accessKeyArg, secretKeyArg, systemArg)
|
|
if err != nil {
|
|
return errorOrIsNotFound(err, "failed to create realm system user %q for reason: %q", uid, output)
|
|
}
|
|
log.NamedDebug(objContext.NsName(), logger, "created realm system user %v", uid)
|
|
} else {
|
|
return errorOrIsNotFound(err, "radosgw-admin user info for system user failed with code %d and output %q", strconv.Itoa(code), output)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func configureObjectStore(objContext *Context, store *cephv1.CephObjectStore, zone *cephv1.CephObjectZone) error {
|
|
log.NamedDebug(objContext.NsName(), logger, "setting multisite configuration for object-store %v", store.Name)
|
|
|
|
if store.Spec.IsMultisite() {
|
|
if zone != nil && len(zone.Spec.CustomEndpoints) == 0 {
|
|
// get list of endpoints not including the endpoint of the object-store for the zone
|
|
zoneEndpointsList, isEndpointAlreadyExists, err := getZoneEndpoints(objContext, objContext.Endpoint)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// There is no need to update the Zone endpoints when:
|
|
// - the zone does not have the endpoint and the synchronization is disabled on the objectstore
|
|
// - the zone already have the endpoint and the synchronization is enabled
|
|
if isEndpointAlreadyExists == store.Spec.Gateway.DisableMultisiteSyncTraffic {
|
|
if !isEndpointAlreadyExists {
|
|
zoneEndpointsList = append(zoneEndpointsList, objContext.Endpoint)
|
|
}
|
|
|
|
zoneEndpoints := strings.Join(zoneEndpointsList, ",")
|
|
log.NamedDebug(objContext.NsName(), logger, "Endpoints for zone %q are: %q", objContext.Zone, zoneEndpoints)
|
|
endpointArg := fmt.Sprintf("--endpoints=%s", zoneEndpoints)
|
|
|
|
err = JoinMultisite(objContext, endpointArg, zoneEndpoints, store.Namespace)
|
|
if err != nil {
|
|
return errors.Wrapf(err, "failed join ceph multisite in zone %q", objContext.Zone)
|
|
}
|
|
}
|
|
}
|
|
} else {
|
|
var endpointArg string
|
|
if !store.Spec.Gateway.DisableMultisiteSyncTraffic {
|
|
endpointArg = fmt.Sprintf("--endpoints=%s", objContext.Endpoint)
|
|
}
|
|
err := createNonMultisiteStore(objContext, endpointArg, store)
|
|
if err != nil {
|
|
return errorOrIsNotFound(err, "failed create ceph multisite for object-store %q", objContext.Name)
|
|
}
|
|
}
|
|
|
|
log.NamedInfo(objContext.NsName(), logger, "configuration for object-store %v is complete", store.Name)
|
|
return nil
|
|
}
|
|
|
|
func deleteRealm(context *Context) error {
|
|
realmArg := fmt.Sprintf("--rgw-realm=%s", context.Name)
|
|
zoneGroupArg := fmt.Sprintf("--rgw-zonegroup=%s", context.Name)
|
|
|
|
// Delete in reverse order: zone → zonegroup → realm.
|
|
// Note: radosgw-admin uses "delete" for zone and zonegroup, but "rm" for realm.
|
|
_, err := runAdminCommand(context, false, "zone", "delete")
|
|
if err != nil {
|
|
log.NamedWarning(context.NsName(), logger, "failed to delete rgw zone %q. %v", context.Name, err)
|
|
}
|
|
|
|
_, err = RunAdminCommandNoMultisite(context, false, "zonegroup", "delete", realmArg, zoneGroupArg)
|
|
if err != nil {
|
|
log.NamedWarning(context.NsName(), logger, "failed to delete rgw zonegroup %q. %v", context.Name, err)
|
|
}
|
|
|
|
_, err = RunAdminCommandNoMultisite(context, false, "realm", "rm", realmArg)
|
|
if err != nil {
|
|
log.NamedWarning(context.NsName(), logger, "failed to delete rgw realm %q. %v", context.Name, err)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func decodeID(data string) (string, error) {
|
|
var id idType
|
|
err := json.Unmarshal([]byte(data), &id)
|
|
if err != nil {
|
|
return "", errors.Wrap(err, "failed to unmarshal json")
|
|
}
|
|
|
|
return id.ID, err
|
|
}
|
|
|
|
func DecodeZoneGroupConfig(data string) (zoneGroupType, error) {
|
|
var config zoneGroupType
|
|
err := json.Unmarshal([]byte(data), &config)
|
|
if err != nil {
|
|
return config, errors.Wrap(err, "failed to unmarshal json")
|
|
}
|
|
|
|
return config, err
|
|
}
|
|
|
|
func getObjectStores(context *Context) ([]string, error) {
|
|
output, err := RunAdminCommandNoMultisite(context, true, "realm", "list")
|
|
if err != nil {
|
|
// This handles the case where the pod we use to exec command (act as a proxy) is not found/ready yet
|
|
// The caller can nicely handle the error and not overflow the op logs with misleading error messages
|
|
if kerrors.IsNotFound(err) {
|
|
return []string{}, err
|
|
}
|
|
// exit status 2 indicates the object store does not exist, so return nothing
|
|
if strings.Index(err.Error(), "exit status 2") == 0 {
|
|
return []string{}, nil
|
|
}
|
|
return nil, err
|
|
}
|
|
|
|
var r realmType
|
|
err = json.Unmarshal([]byte(output), &r)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "Failed to unmarshal realms")
|
|
}
|
|
|
|
return r.Realms, nil
|
|
}
|
|
|
|
func DeletePools(ctx *Context, lastStore bool, poolPrefix string) error {
|
|
pools := append(metadataPools, dataPoolName)
|
|
if lastStore {
|
|
pools = append(pools, rootPool)
|
|
}
|
|
|
|
if configurePoolsConcurrently() {
|
|
waitGroup, _ := errgroup.WithContext(ctx.clusterInfo.Context)
|
|
for _, pool := range pools {
|
|
name := poolName(poolPrefix, pool)
|
|
waitGroup.Go(func() error {
|
|
if err := cephclient.DeletePool(ctx.Context, ctx.clusterInfo, name); err != nil {
|
|
return errors.Wrapf(err, "failed to delete pool %q. ", name)
|
|
}
|
|
return nil
|
|
},
|
|
)
|
|
}
|
|
|
|
// Wait for all the pools to be deleted
|
|
if err := waitGroup.Wait(); err != nil {
|
|
log.NamedWarning(ctx.NsName(), logger, "%s", err)
|
|
}
|
|
|
|
} else {
|
|
for _, pool := range pools {
|
|
name := poolName(poolPrefix, pool)
|
|
if err := cephclient.DeletePool(ctx.Context, ctx.clusterInfo, name); err != nil {
|
|
log.NamedWarning(ctx.NsName(), logger, "failed to delete pool %q. %v", name, err)
|
|
}
|
|
}
|
|
}
|
|
|
|
// Delete erasure code profile if any
|
|
erasureCodes, err := cephclient.ListErasureCodeProfiles(ctx.Context, ctx.clusterInfo)
|
|
if err != nil {
|
|
return errors.Wrapf(err, "failed to list erasure code profiles for cluster %s", ctx.clusterInfo.Namespace)
|
|
}
|
|
// cleans up the EC profile for the data pool only. Metadata pools don't support EC (only replication is supported).
|
|
// The profile name must match the full pool name used during creation (e.g., "<store>.rgw.buckets.data_ecprofile").
|
|
dataPool := poolName(poolPrefix, dataPoolName)
|
|
ecProfileName := cephclient.GetErasureCodeProfileForPool(dataPool)
|
|
for i := range erasureCodes {
|
|
if erasureCodes[i] == ecProfileName {
|
|
if err := cephclient.DeleteErasureCodeProfile(ctx.Context, ctx.clusterInfo, ecProfileName); err != nil {
|
|
return errors.Wrapf(err, "failed to delete erasure code profile %s for object store %s", ecProfileName, ctx.Name)
|
|
}
|
|
break
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func allObjectPools(storeName string) []string {
|
|
baseObjPools := append(metadataPools, dataPoolName, rootPool)
|
|
|
|
poolsForThisStore := make([]string, 0, len(baseObjPools))
|
|
for _, p := range baseObjPools {
|
|
poolsForThisStore = append(poolsForThisStore, poolName(storeName, p))
|
|
}
|
|
return poolsForThisStore
|
|
}
|
|
|
|
// Detect if there are pools that do not exist for this object store
|
|
func missingPools(context *Context) ([]string, error) {
|
|
// list pools instead of querying each pool individually. querying each individually makes it
|
|
// hard to determine if an error is because the pool does not exist or because of a connection
|
|
// issue with ceph mons (or some other underlying issue). if listing pools fails, we can be sure
|
|
// it is a connection issue and return an error.
|
|
existingPoolSummaries, err := cephclient.ListPoolSummaries(context.Context, context.clusterInfo)
|
|
if err != nil {
|
|
return []string{}, errors.Wrapf(err, "failed to determine if pools are missing. failed to list pools")
|
|
}
|
|
existingPools := sets.New[string]()
|
|
for _, summary := range existingPoolSummaries {
|
|
existingPools.Insert(summary.Name)
|
|
}
|
|
|
|
missingPools := []string{}
|
|
for _, objPool := range allObjectPools(context.Zone) {
|
|
if !existingPools.Has(objPool) {
|
|
missingPools = append(missingPools, objPool)
|
|
}
|
|
}
|
|
|
|
return missingPools, nil
|
|
}
|
|
|
|
// poolsExistForNonRawOps checks whether all object store pools exist.
|
|
// Non-raw radosgw-admin commands (ones not in RGW raw_storage_ops_list, see:
|
|
// https://github.com/ceph/ceph/blob/649c3131862d310ce0b273dcbe08ebca0c2ec7ec/src/rgw/radosgw-admin/radosgw-admin.cc#L4540-L4567
|
|
// )
|
|
// trigger full RGW initialization which auto-creates default pools — for example user info called by disableRGWDashboard and
|
|
// user create called by NewMultisiteAdminOpsContext during store deletion.
|
|
// Use this function on delete reconciliation path to skip non-raw commands if pools were not created.
|
|
func poolsExistForNonRawOps(context *Context) bool {
|
|
missing, err := missingPools(context)
|
|
if err != nil {
|
|
log.NamedWarning(context.NsName(), logger, "failed to check for missing pools, assuming they don't exist to avoid ghost pool creation: %v", err)
|
|
return false
|
|
}
|
|
if len(missing) > 0 {
|
|
log.NamedInfo(context.NsName(), logger, "some object store pools are missing, skipping non-raw radosgw-admin operations to avoid ghost pool creation: %v", missing)
|
|
return false
|
|
}
|
|
return true
|
|
}
|
|
|
|
func CreateObjectStorePools(context *Context, cluster *cephv1.ClusterSpec, metadataPool, dataPool cephv1.PoolSpec) error {
|
|
if EmptyPool(dataPool) || EmptyPool(metadataPool) {
|
|
log.NamedInfo(context.NsName(), logger, "no pools specified for the CR, checking for their existence...")
|
|
missingPools, err := missingPools(context)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if len(missingPools) > 0 {
|
|
return fmt.Errorf("CR store pools are missing: %v", missingPools)
|
|
}
|
|
|
|
// pools exist, nothing to do
|
|
return nil
|
|
}
|
|
|
|
if err := createSimilarPools(context, append(metadataPools, rootPool), cluster, metadataPool, rgwRadosPoolPgNum); err != nil {
|
|
return errors.Wrap(err, "failed to create metadata pools")
|
|
}
|
|
|
|
if err := createSimilarPools(context, []string{dataPoolName}, cluster, dataPool, cephclient.DefaultPGCount); err != nil {
|
|
return errors.Wrap(err, "failed to create data pool")
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func ConfigureSharedPoolsForZone(objContext *Context, sharedPools cephv1.ObjectSharedPoolsSpec) error {
|
|
if sharedPools.DataPoolName == "" && sharedPools.MetadataPoolName == "" && len(sharedPools.PoolPlacements) == 0 {
|
|
log.NamedDebug(objContext.NsName(), logger, "no shared pools to configure for store")
|
|
return nil
|
|
}
|
|
|
|
log.NamedInfo(objContext.NsName(), logger, "configuring shared pools for object store")
|
|
if err := sharedPoolsExist(objContext, sharedPools); err != nil {
|
|
return errors.Wrapf(err, "object store cannot be configured until shared pools exist")
|
|
}
|
|
|
|
zoneConfig, err := getZoneJSON(objContext)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
zoneUpdated, err := adjustZoneDefaultPools(objContext, zoneConfig, sharedPools)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
zoneUpdated, err = adjustZonePlacementPools(zoneUpdated, sharedPools)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
hasZoneChanged := !reflect.DeepEqual(zoneConfig, zoneUpdated)
|
|
|
|
zoneGroupConfig, err := getZoneGroupJSON(objContext)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defaultPlacement := getDefaultPlacementName(sharedPools)
|
|
zoneGroupUpdated, err := adjustZoneGroupPlacementTargets(zoneGroupConfig, zoneUpdated, defaultPlacement)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
hasZoneGroupChanged := !reflect.DeepEqual(zoneGroupConfig, zoneGroupUpdated)
|
|
|
|
// persist configuration updates:
|
|
if hasZoneChanged {
|
|
log.NamedInfo(objContext.NsName(), logger, "zone config changed: performing zone config updates for %s", objContext.Zone)
|
|
_, err := updateZoneJSON(objContext, zoneUpdated)
|
|
if err != nil {
|
|
return fmt.Errorf("unable to persist zone config update for %s: %w", objContext.Zone, err)
|
|
}
|
|
}
|
|
if hasZoneGroupChanged {
|
|
log.NamedInfo(objContext.NsName(), logger, "zonegroup config changed: performing zonegroup config updates for %s", objContext.ZoneGroup)
|
|
_, err = updateZoneGroupJSON(objContext, zoneGroupUpdated)
|
|
if err != nil {
|
|
return fmt.Errorf("unable to persist zonegroup config update for %s: %w", objContext.ZoneGroup, err)
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func sharedPoolsExist(objContext *Context, sharedPools cephv1.ObjectSharedPoolsSpec) error {
|
|
existingPools, err := cephclient.ListPoolSummaries(objContext.Context, objContext.clusterInfo)
|
|
if err != nil {
|
|
return errors.Wrapf(err, "failed to list pools")
|
|
}
|
|
existing := make(map[string]struct{}, len(existingPools))
|
|
for _, pool := range existingPools {
|
|
existing[pool.Name] = struct{}{}
|
|
}
|
|
// sharedPools.MetadataPoolName, DataPoolName, and sharedPools.PoolPlacements.DataNonECPoolName are optional.
|
|
// ignore optional pools with empty name:
|
|
existing[""] = struct{}{}
|
|
|
|
if _, ok := existing[sharedPools.MetadataPoolName]; !ok {
|
|
return fmt.Errorf("sharedPool do not exist: %s", sharedPools.MetadataPoolName)
|
|
}
|
|
if _, ok := existing[sharedPools.DataPoolName]; !ok {
|
|
return fmt.Errorf("sharedPool do not exist: %s", sharedPools.DataPoolName)
|
|
}
|
|
|
|
for _, pp := range sharedPools.PoolPlacements {
|
|
if _, ok := existing[pp.MetadataPoolName]; !ok {
|
|
return fmt.Errorf("sharedPool does not exist: pool %s for placement %s", pp.MetadataPoolName, pp.Name)
|
|
}
|
|
if _, ok := existing[pp.DataPoolName]; !ok {
|
|
return fmt.Errorf("sharedPool do not exist: pool %s for placement %s", pp.DataPoolName, pp.Name)
|
|
}
|
|
if _, ok := existing[pp.DataNonECPoolName]; !ok {
|
|
return fmt.Errorf("sharedPool do not exist: pool %s for placement %s", pp.DataNonECPoolName, pp.Name)
|
|
}
|
|
for _, sc := range pp.StorageClasses {
|
|
if _, ok := existing[sc.DataPoolName]; !ok {
|
|
return fmt.Errorf("sharedPool do not exist: pool %s for StorageClass %s", sc.DataPoolName, sc.Name)
|
|
}
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func adjustZoneDefaultPools(objContext *Context, zone map[string]interface{}, spec cephv1.ObjectSharedPoolsSpec) (map[string]interface{}, error) {
|
|
name, err := getObjProperty[string](zone, "name")
|
|
if err != nil {
|
|
return nil, fmt.Errorf("unable to get zone name: %w", err)
|
|
}
|
|
|
|
zone, err = deepCopyJson(zone)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("unable to deep copy zone %s: %w", name, err)
|
|
}
|
|
|
|
defaultMetaPool := getDefaultMetadataPool(spec)
|
|
if defaultMetaPool == "" {
|
|
// default pool is not presented in shared pool spec
|
|
return zone, nil
|
|
}
|
|
// add zone namespace to metadata pool to safely share across rgw instances or zones.
|
|
// in non-multisite case zone name equals to rgw instance name
|
|
defaultMetaPool = defaultMetaPool + ":" + name
|
|
for pool, nsSuffix := range zonePoolNSSuffix {
|
|
// replace rgw internal index pools with namespaced metadata pool
|
|
namespacedPool := defaultMetaPool + nsSuffix
|
|
|
|
// check if old pool has data BEFORE overwriting the zone property
|
|
prev, _ := zone[pool].(string)
|
|
if prev != "" && prev != namespacedPool {
|
|
empty, err := checkPoolIsEmpty(objContext, prev)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("zone pool field %q: unable to check old pool %q: %w", pool, prev, err)
|
|
}
|
|
if !empty {
|
|
log.NamedWarning(objContext.NsName(), logger,
|
|
"zone pool field %q is being remapped from %q which still contains data", pool, prev)
|
|
continue
|
|
}
|
|
}
|
|
|
|
prev, err := updateObjProperty(zone, namespacedPool, pool)
|
|
if err != nil {
|
|
log.NamedInfo(objContext.NsName(), logger, "unable to apply rados namespace to shared pool: %v", err)
|
|
continue
|
|
}
|
|
if namespacedPool != prev {
|
|
log.NamedDebug(objContext.NsName(), logger, "update shared pool %s for zone %s: %s -> %s", pool, name, prev, namespacedPool)
|
|
}
|
|
}
|
|
|
|
// check for unknown pool properties in zone json
|
|
for field, val := range zone {
|
|
if _, ok := val.(string); !ok {
|
|
// not a string property
|
|
continue
|
|
}
|
|
if !strings.HasSuffix(field, "_pool") {
|
|
// not a pool property
|
|
continue
|
|
}
|
|
if _, ok := zonePoolNSSuffix[field]; !ok {
|
|
log.NamedWarning(objContext.NsName(), logger, "zone config %q contains unknown pool %q", name, field)
|
|
}
|
|
}
|
|
|
|
return zone, nil
|
|
}
|
|
|
|
func ZoneJsonPoolKeys() sets.Set[string] {
|
|
s := sets.New[string]()
|
|
for k := range zonePoolNSSuffix {
|
|
s.Insert(k)
|
|
}
|
|
return s
|
|
}
|
|
|
|
func checkPoolIsEmpty(objContext *Context, name string) (bool, error) {
|
|
poolName, namespace, isNamespaced := strings.Cut(name, ":")
|
|
|
|
poolExists, err := cephclient.IsPoolPresent(objContext.Context, objContext.clusterInfo, poolName)
|
|
if err != nil {
|
|
return false, fmt.Errorf("failed to check if pool %q exists: %w", poolName, err)
|
|
}
|
|
if !poolExists {
|
|
return true, nil
|
|
}
|
|
|
|
if isNamespaced {
|
|
hasObjects, err := cephclient.RadosNamespaceHasObjects(objContext.Context, objContext.clusterInfo, poolName, namespace)
|
|
if err != nil {
|
|
return false, fmt.Errorf("failed to check pool %q namespace %q: %w", poolName, namespace, err)
|
|
}
|
|
return !hasObjects, nil
|
|
}
|
|
|
|
stats, err := cephclient.GetPoolStats(objContext.Context, objContext.clusterInfo)
|
|
if err != nil {
|
|
return false, fmt.Errorf("failed to get pool stats: %w", err)
|
|
}
|
|
for _, p := range stats.Pools {
|
|
if p.Name == poolName && p.Stats.Objects > 0 {
|
|
return false, nil
|
|
}
|
|
}
|
|
return true, nil
|
|
}
|
|
|
|
// configurePoolsConcurrently checks if operator pod resources are set or not
|
|
func configurePoolsConcurrently() bool {
|
|
// if operator resources are specified return false as it will lead to operator pod killed due to resource limit
|
|
// nolint S1008 (go-staticcheck), we can safely suppress this
|
|
if os.Getenv("OPERATOR_RESOURCES_SPECIFIED") == "true" {
|
|
return false
|
|
}
|
|
return true
|
|
}
|
|
|
|
func createSimilarPools(ctx *Context, pools []string, cluster *cephv1.ClusterSpec, poolSpec cephv1.PoolSpec, pgCount string) error {
|
|
// We have concurrency
|
|
if configurePoolsConcurrently() {
|
|
waitGroup, _ := errgroup.WithContext(ctx.clusterInfo.Context)
|
|
for _, pool := range pools {
|
|
// Avoid the loop reusing the same value with a closure
|
|
pool := pool
|
|
|
|
waitGroup.Go(func() error { return createRGWPool(ctx, cluster, poolSpec, pgCount, pool) })
|
|
}
|
|
return waitGroup.Wait()
|
|
}
|
|
|
|
// No concurrency!
|
|
for _, pool := range pools {
|
|
err := createRGWPool(ctx, cluster, poolSpec, pgCount, pool)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func createRGWPool(ctx *Context, cluster *cephv1.ClusterSpec, poolSpec cephv1.PoolSpec, pgCount, requestedName string) error {
|
|
// create the pool if it doesn't exist yet
|
|
poolSpec.Application = rgwApplication
|
|
pool := cephv1.NamedPoolSpec{
|
|
Name: poolName(ctx.Name, requestedName),
|
|
PoolSpec: poolSpec,
|
|
}
|
|
if err := cephclient.CreatePoolWithPGs(ctx.Context, ctx.clusterInfo, cluster, &pool, pgCount); err != nil {
|
|
return errors.Wrapf(err, "failed to create pool %q", pool.Name)
|
|
}
|
|
// Set the pg_num_min if not the default so the autoscaler won't immediately increase the pg count
|
|
if pgCount != cephclient.DefaultPGCount {
|
|
if err := cephclient.SetPoolProperty(ctx.Context, ctx.clusterInfo, pool.Name, "pg_num_min", pgCount); err != nil {
|
|
return errors.Wrapf(err, "failed to set pg_num_min on pool %q to %q", pool.Name, pgCount)
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func poolName(poolPrefix, poolName string) string {
|
|
if strings.HasPrefix(poolName, ".") {
|
|
return poolName
|
|
}
|
|
// the name of the pool is <instance>.<name>, except for the pool ".rgw.root" that spans object stores
|
|
return fmt.Sprintf("%s.%s", poolPrefix, poolName)
|
|
}
|
|
|
|
// GetObjectBucketProvisioner returns the bucket provisioner name appended with operator namespace if OBC is watching on it
|
|
func GetObjectBucketProvisioner(namespace string) (string, error) {
|
|
provName := bucketProvisionerName
|
|
obcWatchOnNamespace := k8sutil.GetOperatorSetting("ROOK_OBC_WATCH_OPERATOR_NAMESPACE", "false")
|
|
obcProvisionerNamePrefix := k8sutil.GetOperatorSetting("ROOK_OBC_PROVISIONER_NAME_PREFIX", "")
|
|
if obcProvisionerNamePrefix != "" {
|
|
errList := validation.IsDNS1123Label(obcProvisionerNamePrefix)
|
|
if len(errList) > 0 {
|
|
return "", errors.Errorf("invalid OBC provisioner name prefix %q. %v", obcProvisionerNamePrefix, errList)
|
|
}
|
|
provName = fmt.Sprintf("%s.%s", obcProvisionerNamePrefix, bucketProvisionerName)
|
|
} else if obcWatchOnNamespace == "true" {
|
|
provName = fmt.Sprintf("%s.%s", namespace, bucketProvisionerName)
|
|
}
|
|
return provName, nil
|
|
}
|
|
|
|
// checkDashboardUser returns true if the dashboard user exists and has the same credentials as the given user, else return false
|
|
func checkDashboardUser(context *Context, user ObjectUser) (bool, error) {
|
|
dUser, errId, err := GetUser(context, DashboardUser)
|
|
|
|
// If not found or "none" error, all is good to not return the error
|
|
switch errId {
|
|
case RGWErrorNone:
|
|
// If the access key or secret key is not the same as the given user, return false
|
|
if user.AccessKey != nil && *user.AccessKey != *dUser.AccessKey {
|
|
return false, nil
|
|
}
|
|
if user.SecretKey != nil && *user.SecretKey != *dUser.SecretKey {
|
|
return false, nil
|
|
}
|
|
return true, nil
|
|
case RGWErrorNotFound:
|
|
return false, nil
|
|
}
|
|
return false, err
|
|
}
|
|
|
|
// retrieveDashboardAPICredentials Retrieves the dashboard's access and secret key and set it on the given ObjectUser
|
|
func retrieveDashboardAPICredentials(context *Context, user *ObjectUser) error {
|
|
args := []string{"dashboard", "get-rgw-api-access-key"}
|
|
cephCmd := cephclient.NewCephCommand(context.Context, context.clusterInfo, args)
|
|
out, err := cephCmd.Run()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
if string(out) != "" {
|
|
accessKey := string(out)
|
|
user.AccessKey = &accessKey
|
|
}
|
|
|
|
args = []string{"dashboard", "get-rgw-api-secret-key"}
|
|
cephCmd = cephclient.NewCephCommand(context.Context, context.clusterInfo, args)
|
|
out, err = cephCmd.Run()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
if string(out) != "" {
|
|
secretKey := string(out)
|
|
user.SecretKey = &secretKey
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func getDashboardUser(context *Context) (ObjectUser, error) {
|
|
user := ObjectUser{
|
|
UserID: DashboardUser,
|
|
DisplayName: &DashboardUser,
|
|
SystemUser: true,
|
|
}
|
|
|
|
// Retrieve RGW Dashboard credentials if some are already set
|
|
if err := retrieveDashboardAPICredentials(context, &user); err != nil {
|
|
return user, errors.Wrapf(err, "failed to retrieve RGW Dashboard credentials for %q user", DashboardUser)
|
|
}
|
|
|
|
return user, nil
|
|
}
|
|
|
|
func enableRGWDashboard(context *Context) error {
|
|
log.NamedInfo(context.NsName(), logger, "enabling rgw dashboard")
|
|
|
|
user, err := getDashboardUser(context)
|
|
if err != nil {
|
|
log.NamedDebug(context.NsName(), logger, "failed to get current dashboard user")
|
|
return err
|
|
}
|
|
|
|
checkDashboard, err := checkDashboardUser(context, user)
|
|
if err != nil {
|
|
log.NamedDebug(context.NsName(), logger, "Unable to fetch dashboard user key for RGW, hence skipping")
|
|
return nil
|
|
}
|
|
if checkDashboard {
|
|
log.NamedDebug(context.NsName(), logger, "RGW Dashboard is already enabled")
|
|
return nil
|
|
}
|
|
|
|
// TODO:
|
|
// Use admin ops user instead!
|
|
// It's safe to create the user with the force flag regardless if the cluster's dashboard is
|
|
// configured as a secondary rgw site. The creation will return the user already exists and we
|
|
// will just fetch it (it has been created by the primary cluster)
|
|
u, errCode, err := CreateOrRecreateUserIfExists(context, user, true)
|
|
if err != nil || errCode != RGWErrorNone {
|
|
// Handle already exists ErrorCodeFileExists
|
|
return errors.Wrapf(err, "failed to create/ re-create user %q", DashboardUser)
|
|
}
|
|
|
|
var accessArgs, secretArgs []string
|
|
var secretFile *os.File
|
|
|
|
accessFile, err := util.CreateTempFile(*u.AccessKey)
|
|
if err != nil {
|
|
return errors.Wrap(err, "failed to create a temporary dashboard access-key file")
|
|
}
|
|
|
|
accessArgs = []string{"dashboard", "set-rgw-api-access-key", "-i", accessFile.Name()}
|
|
defer func() {
|
|
if err := os.Remove(accessFile.Name()); err != nil {
|
|
log.NamedError(context.NsName(), logger, "failed to clean up dashboard access-key file. %v", err)
|
|
}
|
|
}()
|
|
|
|
secretFile, err = util.CreateTempFile(*u.SecretKey)
|
|
if err != nil {
|
|
return errors.Wrap(err, "failed to create a temporary dashboard secret-key file")
|
|
}
|
|
|
|
secretArgs = []string{"dashboard", "set-rgw-api-secret-key", "-i", secretFile.Name()}
|
|
|
|
cephCmd := cephclient.NewCephCommand(context.Context, context.clusterInfo, accessArgs)
|
|
_, err = cephCmd.Run()
|
|
if err != nil {
|
|
return errors.Wrapf(err, "failed to set user %q accesskey", DashboardUser)
|
|
}
|
|
|
|
cephCmd = cephclient.NewCephCommand(context.Context, context.clusterInfo, secretArgs)
|
|
go func() {
|
|
// Setting the dashboard api secret started hanging in some clusters
|
|
// starting in ceph v15.2.8. We run it in a goroutine until the fix
|
|
// is found. We expect the ceph command to timeout so at least the goroutine exits.
|
|
log.NamedInfo(context.NsName(), logger, "setting the dashboard api secret key")
|
|
_, err = cephCmd.RunWithTimeout(exec.CephCommandsTimeout)
|
|
if err != nil {
|
|
log.NamedError(context.NsName(), logger, "failed to set user %q secretkey. %v", DashboardUser, err)
|
|
}
|
|
if err := os.Remove(secretFile.Name()); err != nil {
|
|
log.NamedError(context.NsName(), logger, "failed to clean up dashboard secret-key file. %v", err)
|
|
}
|
|
|
|
log.NamedInfo(context.NsName(), logger, "done setting the dashboard api secret key")
|
|
}()
|
|
|
|
return nil
|
|
}
|
|
|
|
func disableRGWDashboard(context *Context) {
|
|
log.NamedInfo(context.NsName(), logger, "disabling the dashboard api user and secret key")
|
|
|
|
// GetUser and DeleteUser are radosgw-admin non-raw commands that trigger full RGW
|
|
// initialization and auto-create pools. Guard them so we don't create ghost pools
|
|
// when called during object store deletion after pools are already removed.
|
|
if poolsExistForNonRawOps(context) {
|
|
_, _, err := GetUser(context, DashboardUser)
|
|
if err != nil {
|
|
log.NamedInfo(context.NsName(), logger, "unable to fetch the user %q details from this objectstore %q", DashboardUser, context.Name)
|
|
} else {
|
|
log.NamedInfo(context.NsName(), logger, "deleting rgw dashboard user")
|
|
_, err = DeleteUser(context, DashboardUser)
|
|
if err != nil {
|
|
log.NamedWarning(context.NsName(), logger, "failed to delete ceph user %q. %v", DashboardUser, err)
|
|
}
|
|
}
|
|
}
|
|
|
|
args := []string{"dashboard", "reset-rgw-api-access-key"}
|
|
cephCmd := cephclient.NewCephCommand(context.Context, context.clusterInfo, args)
|
|
_, err := cephCmd.RunWithTimeout(exec.CephCommandsTimeout)
|
|
if err != nil {
|
|
log.NamedWarning(context.NsName(), logger, "failed to reset user accesskey for user %q. %v", DashboardUser, err)
|
|
}
|
|
|
|
args = []string{"dashboard", "reset-rgw-api-secret-key"}
|
|
cephCmd = cephclient.NewCephCommand(context.Context, context.clusterInfo, args)
|
|
_, err = cephCmd.RunWithTimeout(exec.CephCommandsTimeout)
|
|
if err != nil {
|
|
log.NamedWarning(context.NsName(), logger, "failed to reset user secretkey for user %q. %v", DashboardUser, err)
|
|
}
|
|
log.NamedInfo(context.NsName(), logger, "done disabling the dashboard api secret key")
|
|
}
|
|
|
|
func errorOrIsNotFound(err error, msg string, args ...string) error {
|
|
// This handles the case where the pod we use to exec command (act as a proxy) is not found/ready yet
|
|
// The caller can nicely handle the error and not overflow the op logs with misleading error messages
|
|
if kerrors.IsNotFound(err) {
|
|
return err
|
|
}
|
|
return errors.Wrapf(err, msg, args)
|
|
}
|
|
|
|
// ShouldUpdateZoneEndpointList checks whether zone endpoint list need to be updated or not
|
|
func ShouldUpdateZoneEndpointList(zones []zoneType, desiredEndpointList []string, zoneName string) (bool, error) {
|
|
if zoneName == "" {
|
|
return false, errors.Errorf("zone name can't be empty")
|
|
}
|
|
|
|
zoneExists, endpoints := findZoneEndpoints(zoneName, zones)
|
|
if !zoneExists {
|
|
return false, nil
|
|
}
|
|
|
|
return !listsAreEqual(desiredEndpointList, endpoints), nil
|
|
}
|
|
|
|
func findZoneEndpoints(targetZone string, zones []zoneType) (bool, []string) {
|
|
for _, z := range zones {
|
|
if z.Name == targetZone {
|
|
return true, z.Endpoints
|
|
}
|
|
}
|
|
return false, []string{}
|
|
}
|
|
|
|
func listsAreEqual(a, b []string) bool {
|
|
as := make([]string, len(a))
|
|
bs := make([]string, len(b))
|
|
copy(as, a)
|
|
copy(bs, b)
|
|
sort.Strings(as)
|
|
sort.Strings(bs)
|
|
|
|
if len(a) != len(b) {
|
|
return false
|
|
}
|
|
for i := 0; i < len(a); i++ {
|
|
if as[i] != bs[i] {
|
|
return false
|
|
}
|
|
}
|
|
return true
|
|
}
|
|
|
|
func CheckIfZonePresentInZoneGroup(objContext *Context) (bool, error) {
|
|
output, err := runAdminCommand(objContext, true, "zonegroup", "get")
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
zoneGroupJson, err := DecodeZoneGroupConfig(output)
|
|
if err != nil {
|
|
return false, errors.Wrap(err, "failed to parse `radosgw-admin zonegroup get` output")
|
|
}
|
|
zoneExists, _ := findZoneEndpoints(objContext.Zone, zoneGroupJson.Zones)
|
|
if zoneExists {
|
|
return true, nil
|
|
}
|
|
return false, nil
|
|
}
|
|
|
|
// ValidateObjectStorePoolsConfig returns error if given ObjectStore pool configuration is inconsistent.
|
|
func ValidateObjectStorePoolsConfig(metadataPool, dataPool cephv1.PoolSpec, sharedPools cephv1.ObjectSharedPoolsSpec) error {
|
|
if err := validatePoolPlacements(sharedPools.PoolPlacements); err != nil {
|
|
return err
|
|
}
|
|
if !EmptyPool(dataPool) && sharedPools.DataPoolName != "" {
|
|
return fmt.Errorf("invalidObjStorePoolConfig: object store dataPool and sharedPools.dataPool=%s are mutually exclusive. Only one of them can be set", sharedPools.DataPoolName)
|
|
}
|
|
if !EmptyPool(metadataPool) && sharedPools.MetadataPoolName != "" {
|
|
return fmt.Errorf("invalidObjStorePoolConfig: object store metadataPool and sharedPools.metadataPool=%s are mutually exclusive. Only one of them can be set", sharedPools.MetadataPoolName)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func SetDefaultRealm(objContext *Context, realmName string) error {
|
|
args := []string{"realm", "default", fmt.Sprintf("--rgw-realm=%s", realmName)}
|
|
output, err := RunAdminCommandNoMultisite(objContext, false, args...)
|
|
if err != nil {
|
|
return errors.Wrapf(err, "failed to set realm %q as default, reason: %q", realmName, output)
|
|
}
|
|
|
|
log.NamedInfo(objContext.NsName(), logger, "successfully set realm %q as default", realmName+"/"+objContext.clusterInfo.Namespace)
|
|
return nil
|
|
}
|
|
|
|
// InitializeObjectStoreContext finds the CephObjectStore by storeName, verifies it
|
|
// is initialized (has running RGW pods or is external), creates a multisite context,
|
|
// and initializes the admin ops API client.
|
|
func InitializeObjectStoreContext(
|
|
clusterdContext *clusterd.Context,
|
|
clusterInfo *cephclient.ClusterInfo,
|
|
k8sClient client.Client,
|
|
opManagerContext context.Context,
|
|
storeName string,
|
|
newAdminOpsCtxFunc func(*Context, *cephv1.ObjectStoreSpec) (*AdminOpsContext, error),
|
|
) (*AdminOpsContext, *cephv1.CephObjectStore, error) {
|
|
store, err := getObjectStore(k8sClient, opManagerContext, storeName, clusterInfo.Namespace)
|
|
if err != nil {
|
|
return nil, nil, errors.Wrapf(err, "failed to get object store %q", storeName)
|
|
}
|
|
|
|
if !store.Spec.IsExternal() {
|
|
if err := checkRGWPodsRunning(k8sClient, opManagerContext, clusterInfo.Namespace, storeName); err != nil {
|
|
return nil, nil, errors.Wrapf(err, "failed to detect if object store %q is initialized", storeName)
|
|
}
|
|
}
|
|
|
|
objContext, err := NewMultisiteContext(clusterdContext, clusterInfo, store)
|
|
if err != nil {
|
|
return nil, nil, errors.Wrapf(err, "failed to create multisite context for object store %q", store.Name)
|
|
}
|
|
|
|
opsContext, err := newAdminOpsCtxFunc(objContext, &store.Spec)
|
|
if err != nil {
|
|
return nil, nil, errors.Wrap(err, "failed to initialize rgw admin ops client api")
|
|
}
|
|
|
|
return opsContext, store, nil
|
|
}
|
|
|
|
func getObjectStore(k8sClient client.Client, ctx context.Context, storeName, namespace string) (*cephv1.CephObjectStore, error) {
|
|
store := &cephv1.CephObjectStore{}
|
|
if err := k8sClient.Get(ctx, types.NamespacedName{Name: storeName, Namespace: namespace}, store); err != nil {
|
|
return nil, errors.Wrapf(err, "could not find CephObjectStore %q", storeName)
|
|
}
|
|
return store, nil
|
|
}
|
|
|
|
func checkRGWPodsRunning(k8sClient client.Client, ctx context.Context, namespace, storeName string) error {
|
|
pods := &v1.PodList{}
|
|
listOpts := []client.ListOption{
|
|
client.InNamespace(namespace),
|
|
client.MatchingLabels(labelsForRgw(storeName)),
|
|
}
|
|
|
|
if err := k8sClient.List(ctx, pods, listOpts...); err != nil {
|
|
if kerrors.IsNotFound(err) {
|
|
return errors.Wrap(err, "no rgw pod could not be found")
|
|
}
|
|
return errors.Wrap(err, "failed to list rgw pods")
|
|
}
|
|
|
|
if len(pods.Items) > 0 {
|
|
return nil
|
|
}
|
|
|
|
return errors.New("no rgw pod found")
|
|
}
|
|
|
|
func labelsForRgw(name string) map[string]string {
|
|
return map[string]string{"rgw": name, k8sutil.AppAttr: AppName}
|
|
}
|