Files
Michael Adam c90913051b rgw: fix a comment typo
This fixes a comment typo caught bt the codespell CI check.

Signed-off-by: Michael Adam <obnox@samba.org>
2026-07-15 18:25:53 +02:00

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}
}