org.apache.cloudstack
cloud-plugin-storage-object-simulator
diff --git a/plugins/pom.xml b/plugins/pom.xml
index 92768827f658..e99459c887c2 100755
--- a/plugins/pom.xml
+++ b/plugins/pom.xml
@@ -142,6 +142,7 @@
storage/object/minio
storage/object/ceph
storage/object/cloudian
+ storage/object/seaweedfs
storage/object/simulator
diff --git a/plugins/storage/object/seaweedfs/pom.xml b/plugins/storage/object/seaweedfs/pom.xml
new file mode 100644
index 000000000000..815359a23277
--- /dev/null
+++ b/plugins/storage/object/seaweedfs/pom.xml
@@ -0,0 +1,70 @@
+
+
+ 4.0.0
+ cloud-plugin-storage-object-seaweedfs
+ Apache CloudStack Plugin - SeaweedFS object storage provider
+
+ org.apache.cloudstack
+ cloudstack-plugins
+ 24.0.0-SNAPSHOT
+ ../../../pom.xml
+
+
+
+ org.apache.cloudstack
+ cloud-engine-storage
+ ${project.version}
+
+
+ org.apache.cloudstack
+ cloud-engine-storage-object
+ ${project.version}
+
+
+ org.apache.cloudstack
+ cloud-engine-schema
+ ${project.version}
+
+
+ com.amazonaws
+ aws-java-sdk-core
+
+
+ com.amazonaws
+ aws-java-sdk-iam
+
+
+ com.amazonaws
+ aws-java-sdk-s3
+
+
+ com.fasterxml.jackson.core
+ jackson-databind
+ ${cs.jackson.version}
+
+
+ com.github.tomakehurst
+ wiremock-standalone
+ ${cs.wiremock.version}
+ test
+
+
+
diff --git a/plugins/storage/object/seaweedfs/src/main/java/org/apache/cloudstack/storage/datastore/driver/SeaweedFSObjectStoreDriverImpl.java b/plugins/storage/object/seaweedfs/src/main/java/org/apache/cloudstack/storage/datastore/driver/SeaweedFSObjectStoreDriverImpl.java
new file mode 100644
index 000000000000..48884610a127
--- /dev/null
+++ b/plugins/storage/object/seaweedfs/src/main/java/org/apache/cloudstack/storage/datastore/driver/SeaweedFSObjectStoreDriverImpl.java
@@ -0,0 +1,1072 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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.
+ */
+// SPDX-License-Identifier: Apache-2.0
+package org.apache.cloudstack.storage.datastore.driver;
+
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+import javax.inject.Inject;
+
+import org.apache.cloudstack.engine.subsystem.api.storage.DataStore;
+import org.apache.cloudstack.storage.datastore.db.ObjectStoreDao;
+import org.apache.cloudstack.storage.datastore.db.ObjectStoreDetailsDao;
+import org.apache.cloudstack.storage.datastore.db.ObjectStoreVO;
+import org.apache.cloudstack.storage.datastore.util.SeaweedFSObjectStoreUtil;
+import org.apache.cloudstack.storage.object.BaseObjectStoreDriverImpl;
+import org.apache.cloudstack.storage.object.Bucket;
+import org.apache.cloudstack.storage.object.BucketObject;
+
+import com.amazonaws.AmazonClientException;
+import com.amazonaws.services.identitymanagement.AmazonIdentityManagement;
+import com.amazonaws.services.identitymanagement.model.AccessKey;
+import com.amazonaws.services.identitymanagement.model.AccessKeyMetadata;
+import com.amazonaws.services.identitymanagement.model.CreateAccessKeyRequest;
+import com.amazonaws.services.identitymanagement.model.CreateAccessKeyResult;
+import com.amazonaws.services.identitymanagement.model.CreateUserRequest;
+import com.amazonaws.services.identitymanagement.model.DeleteAccessKeyRequest;
+import com.amazonaws.services.identitymanagement.model.EntityAlreadyExistsException;
+import com.amazonaws.services.identitymanagement.model.ListAccessKeysRequest;
+import com.amazonaws.services.identitymanagement.model.PutUserPolicyRequest;
+import com.amazonaws.services.s3.AmazonS3;
+import com.amazonaws.services.s3.model.AccessControlList;
+import com.amazonaws.services.s3.model.BucketPolicy;
+import com.amazonaws.services.s3.model.BucketVersioningConfiguration;
+import com.amazonaws.services.s3.model.CreateBucketRequest;
+import com.amazonaws.services.s3.model.DeleteBucketPolicyRequest;
+import com.amazonaws.services.s3.model.BucketCrossOriginConfiguration;
+import com.amazonaws.services.s3.model.CORSRule;
+import com.amazonaws.services.s3.model.GetBucketPolicyRequest;
+import com.amazonaws.services.s3.model.SSEAlgorithm;
+import com.amazonaws.services.s3.model.ServerSideEncryptionByDefault;
+import com.amazonaws.services.s3.model.ServerSideEncryptionConfiguration;
+import com.amazonaws.services.s3.model.ServerSideEncryptionRule;
+import com.amazonaws.services.s3.model.SetBucketCrossOriginConfigurationRequest;
+import com.amazonaws.services.s3.model.SetBucketEncryptionRequest;
+import com.amazonaws.services.s3.model.SetBucketVersioningConfigurationRequest;
+import com.cloud.agent.api.to.BucketTO;
+import com.cloud.agent.api.to.DataStoreTO;
+import com.cloud.storage.BucketVO;
+import com.cloud.storage.dao.BucketDao;
+import com.cloud.user.Account;
+import com.cloud.user.AccountDetailsDao;
+import com.cloud.user.dao.AccountDao;
+import com.cloud.utils.db.GlobalLock;
+import com.cloud.utils.exception.CloudRuntimeException;
+
+/**
+ * SeaweedFS object store driver.
+ *
+ * Bucket operations use the AWS S3 SDK v1 (path-style access, endpoint-pinned).
+ * User/credential management uses the AWS IAM SDK v1, since SeaweedFS exposes a
+ * standard AWS IAM-compatible API. No proprietary admin client is needed.
+ *
+ * Modeled on CloudianHyperStoreObjectStoreDriverImpl, which uses the same
+ * S3 + IAM SDK pair.
+ */
+public class SeaweedFSObjectStoreDriverImpl extends BaseObjectStoreDriverImpl {
+
+ @Inject
+ AccountDao _accountDao;
+
+ @Inject
+ AccountDetailsDao _accountDetailsDao;
+
+ @Inject
+ ObjectStoreDao _storeDao;
+
+ @Inject
+ BucketDao _bucketDao;
+
+ @Inject
+ ObjectStoreDetailsDao _storeDetailsDao;
+
+ private static final String ACS_PREFIX = "acs";
+
+ /**
+ * DB-backed global lock name prefix for serializing IAM provisioning and
+ * policy refreshes per store+account. Uses {@link GlobalLock} so the
+ * critical section is serialized across management servers in a
+ * clustered deployment, not just within a single JVM.
+ */
+ private static final String IAM_LOCK_PREFIX = "seaweedfs.iam.";
+ private static final String BUCKET_NAME_LOCK_PREFIX = "seaweedfs.bucket.";
+
+ private static String getIamLockName(long storeId, long accountId) {
+ return IAM_LOCK_PREFIX + storeId + "." + accountId;
+ }
+
+ /**
+ * Acquire a DB-backed global lock for IAM operations on the given
+ * store+account. Returns a {@link GlobalLock} that the caller must
+ * {@link GlobalLock#unlock()} in a {@code finally} block, or {@code null}
+ * if the lock could not be acquired within the timeout.
+ *
+ * Protected so tests can override with a no-op lock (the DB-backed
+ * {@link GlobalLock} requires a real transaction context).
+ */
+ protected GlobalLock acquireIamLock(long storeId, long accountId) {
+ GlobalLock lock = GlobalLock.getInternLock(getIamLockName(storeId, accountId));
+ if (!lock.lock(300)) {
+ logger.warn("Failed to acquire IAM lock for store {} account {}", storeId, accountId);
+ lock.releaseRef();
+ return null;
+ }
+ return lock;
+ }
+
+ private static String getBucketNameLockName(long storeId, String bucketName) {
+ return BUCKET_NAME_LOCK_PREFIX + storeId + "." + bucketName;
+ }
+
+ /**
+ * Acquire a DB-backed global lock covering a bucket name on a store,
+ * independent of the owning account.
+ *
+ * S3 bucket names are globally unique within a store and are reusable after
+ * deletion, while the IAM lock is scoped to (store, account). Without this
+ * lock, once a delete removes the remote bucket but before the old owner's
+ * IAM policy is refreshed, a different account can recreate the same name and
+ * the old owner's credentials would still grant that ARN. Both createBucket
+ * and deleteBucket take this lock so the two never interleave for a name.
+ *
+ *
Protected so tests can override with a no-op lock (the DB-backed
+ * {@link GlobalLock} requires a real transaction context).
+ *
+ * @return the held lock, which the caller must {@link GlobalLock#unlock()}
+ * and {@link GlobalLock#releaseRef()}, or {@code null} on timeout
+ */
+ protected GlobalLock acquireBucketNameLock(long storeId, String bucketName) {
+ GlobalLock lock = GlobalLock.getInternLock(getBucketNameLockName(storeId, bucketName));
+ if (!lock.lock(300)) {
+ logger.warn("Failed to acquire bucket name lock for store {} bucket {}", storeId, bucketName);
+ lock.releaseRef();
+ return null;
+ }
+ return lock;
+ }
+
+ @Override
+ public DataStoreTO getStoreTO(DataStore store) {
+ return null;
+ }
+
+ /**
+ * Get the SeaweedFS IAM user name for the given CloudStack account and
+ * store. The store ID is included so that two CloudStack pools pointing
+ * at the same SeaweedFS IAM service do not collide on the same
+ * {@code acs-} user and overwrite each other's policy and access
+ * keys.
+ */
+ protected String getUserNameForAccount(Account account, long storeId) {
+ return String.format("%s-%d-%s", ACS_PREFIX, storeId, account.getUuid());
+ }
+
+ /**
+ * Create the IAM user for the CloudStack account if it doesn't exist,
+ * attach the restricted S3 policy, and ensure the account has a usable
+ * IAM access key persisted in its account details.
+ *
+ * If a previously stored access key is still present in IAM, it is
+ * reused rather than rotated. A new key is only created when no stored
+ * key exists or the stored key is no longer found in IAM; in the latter
+ * case any unmanaged (leftover) keys for the user are deleted first to
+ * avoid hitting IAM access-key limits. This keeps bucket records that
+ * reference the stored credentials valid across repeated calls.
+ *
+ * @return true if the user exists or was created, false on failure.
+ */
+ @Override
+ public boolean createUser(long accountId, long storeId) {
+ Account account = _accountDao.findById(accountId);
+ if (account == null) {
+ logger.error("Account {} not found", accountId);
+ return false;
+ }
+ String userName = getUserNameForAccount(account, storeId);
+ AmazonIdentityManagement iamClient = getIAMClient(storeId);
+
+ // Serialize per store+account across management servers so two
+ // concurrent bucket requests do not both rotate credentials and leave
+ // bucket rows with mismatched key pairs.
+ GlobalLock lock = acquireIamLock(storeId, accountId);
+ if (lock == null) {
+ return false;
+ }
+ try {
+
+ // Create the IAM user if it doesn't already exist
+ try {
+ iamClient.createUser(new CreateUserRequest(userName));
+ logger.info("Created IAM user {} for account {}", userName, account.getAccountName());
+ } catch (EntityAlreadyExistsException e) {
+ logger.debug("IAM user {} already exists", userName);
+ }
+
+ // Attach a scoped IAM policy that allows access only to this
+ // account's own buckets (the tenant boundary). Refreshed whenever
+ // buckets are created or deleted. Use the lock-free variant since
+ // createUser already holds the IAM lock.
+ updateAccountIAMPolicyLocked(iamClient, storeId, accountId, null);
+
+ // Reuse the stored access key only if both the access key id and the
+ // secret key are present and the key is still Active in IAM; otherwise
+ // create a replacement.
+ Map details = _accountDetailsDao.findDetails(accountId);
+ String accessKeyDetailKey = SeaweedFSObjectStoreUtil.keyAccessKey(storeId);
+ String secretKeyDetailKey = SeaweedFSObjectStoreUtil.keySecretKey(storeId);
+ String storedAccessKeyId = details.get(accessKeyDetailKey);
+ String storedSecretKey = details.get(secretKeyDetailKey);
+ if (storedAccessKeyId != null && storedSecretKey != null
+ && iamAccessKeyExists(iamClient, userName, storedAccessKeyId)) {
+ logger.debug("Reusing existing IAM access key {} for user {}", storedAccessKeyId, userName);
+ updateAccountBucketCredentials(storeId, accountId, storedAccessKeyId, storedSecretKey);
+ return true;
+ }
+
+ // The stored key is missing, inactive, or no longer in IAM. Clean up
+ // ALL keys (including the inactive stored one) before creating a
+ // replacement so we do not accumulate keys and hit IAM limits.
+ deleteUnmanagedAccessKeys(iamClient, userName, null);
+
+ CreateAccessKeyResult result = iamClient.createAccessKey(
+ new CreateAccessKeyRequest().withUserName(userName));
+ AccessKey key = result.getAccessKey();
+
+ // Persist the credential pair in account details (namespaced by storeId)
+ // before updating BucketVO rows. The reuse path above reconciles bucket
+ // rows every time, so a later bucket update failure is repairable on
+ // retry. AccountDetailsDao.persist(accountId, map) is deliberately not
+ // used: it expunges every detail for the account and can clobber another
+ // store's namespaced credentials.
+ try {
+ persistAccountCredentialsOrRollback(accountId, accessKeyDetailKey, secretKeyDetailKey,
+ storedAccessKeyId, storedSecretKey, key);
+ } catch (RuntimeException e) {
+ deleteAccessKeyAfterCredentialPersistenceFailure(iamClient, userName, key.getAccessKeyId(), e);
+ throw e;
+ }
+
+ updateAccountBucketCredentials(storeId, accountId, key);
+
+ logger.info("Created IAM credentials {} for user {}", key.getAccessKeyId(), userName);
+ return true;
+ } finally {
+ lock.unlock();
+ lock.releaseRef();
+ }
+ }
+
+ /**
+ * Persist a SeaweedFS IAM credential pair without using
+ * AccountDetailsDao.persist(accountId, map), and restore the previous pair if
+ * either single-key write fails so callers never observe a half-new pair.
+ */
+ private void persistAccountCredentialsOrRollback(long accountId, String accessKeyDetailKey, String secretKeyDetailKey,
+ String previousAccessKey, String previousSecretKey, AccessKey key) {
+ try {
+ _accountDetailsDao.addDetail(accountId, accessKeyDetailKey, key.getAccessKeyId(), false);
+ _accountDetailsDao.addDetail(accountId, secretKeyDetailKey, key.getSecretAccessKey(), false);
+ } catch (RuntimeException e) {
+ CloudRuntimeException wrapped = new CloudRuntimeException("Failed to persist SeaweedFS IAM credential pair for account " + accountId, e);
+ restoreAccountCredentialDetail(accountId, accessKeyDetailKey, previousAccessKey, wrapped);
+ restoreAccountCredentialDetail(accountId, secretKeyDetailKey, previousSecretKey, wrapped);
+ throw wrapped;
+ }
+ }
+
+ private void restoreAccountCredentialDetail(long accountId, String detailKey, String previousValue, RuntimeException cause) {
+ try {
+ if (previousValue == null) {
+ _accountDetailsDao.removeDetail(accountId, detailKey);
+ } else {
+ _accountDetailsDao.addDetail(accountId, detailKey, previousValue, false);
+ }
+ } catch (RuntimeException rollbackEx) {
+ logger.error("Failed to restore account detail {} for account {} after SeaweedFS credential persistence failed",
+ detailKey, accountId, rollbackEx);
+ cause.addSuppressed(rollbackEx);
+ }
+ }
+
+ private void deleteAccessKeyAfterCredentialPersistenceFailure(AmazonIdentityManagement iamClient, String userName,
+ String accessKeyId, RuntimeException cause) {
+ try {
+ iamClient.deleteAccessKey(new DeleteAccessKeyRequest()
+ .withUserName(userName)
+ .withAccessKeyId(accessKeyId));
+ } catch (AmazonClientException cleanupEx) {
+ logger.error("Failed to delete IAM access key {} for user {} after account credential persistence failed",
+ accessKeyId, userName, cleanupEx);
+ cause.addSuppressed(cleanupEx);
+ }
+ }
+
+ private void updateAccountBucketCredentials(long storeId, long accountId, AccessKey iamCredential) {
+ updateAccountBucketCredentials(storeId, accountId, iamCredential.getAccessKeyId(), iamCredential.getSecretAccessKey());
+ }
+
+ /**
+ * Update the IAM credentials on all BucketVO rows for this store/account so
+ * previously created buckets reflect the current key pair.
+ */
+ private void updateAccountBucketCredentials(long storeId, long accountId, String accessKeyId, String secretAccessKey) {
+ List bucketList = _bucketDao.listByObjectStoreIdAndAccountId(storeId, accountId);
+ for (BucketVO bucketVO : bucketList) {
+ if (accessKeyId.equals(bucketVO.getAccessKey()) && secretAccessKey.equals(bucketVO.getSecretKey())) {
+ continue;
+ }
+ logger.info("Updating accountId={} bucket {} with new IAM credentials", accountId, bucketVO.getName());
+ bucketVO.setAccessKey(accessKeyId);
+ bucketVO.setSecretKey(secretAccessKey);
+ if (!_bucketDao.update(bucketVO.getId(), bucketVO)) {
+ throw new CloudRuntimeException("Failed to update IAM credentials on bucket " + bucketVO.getName());
+ }
+ }
+ }
+
+ /**
+ * Refresh the per-account IAM user policy so it grants S3 access only to
+ * the account's current buckets (optionally excluding one, e.g. a bucket
+ * being deleted). This is the tenant boundary: each account's IAM
+ * credentials can only operate on that account's own buckets.
+ *
+ * Acquires the per-store/account IAM lock. Callers that already hold the
+ * lock (e.g. createUser, createBucket post-create) should call
+ * {@link #updateAccountIAMPolicyLocked} instead to avoid re-entrant lock
+ * acquisition warnings from GlobalLock.
+ *
+ * @param iamClient the IAM client
+ * @param storeId the object store
+ * @param accountId the CloudStack account
+ * @param excludeBucket a bucket name to omit (e.g. a bucket being deleted),
+ * or null to include all of the account's buckets
+ */
+ protected void updateAccountIAMPolicy(AmazonIdentityManagement iamClient, long storeId, long accountId, String excludeBucket) {
+ GlobalLock lock = acquireIamLock(storeId, accountId);
+ if (lock == null) {
+ throw new CloudRuntimeException("Failed to acquire IAM lock for store " + storeId + " account " + accountId);
+ }
+ try {
+ updateAccountIAMPolicyLocked(iamClient, storeId, accountId, excludeBucket);
+ } finally {
+ lock.unlock();
+ lock.releaseRef();
+ }
+ }
+
+ /**
+ * Lock-free variant of {@link #updateAccountIAMPolicy} for callers that
+ * already hold the per-store/account IAM lock. Performs the policy refresh
+ * without reacquiring the lock, avoiding the GlobalLock re-entrant
+ * acquisition warning.
+ */
+ protected void updateAccountIAMPolicyLocked(AmazonIdentityManagement iamClient, long storeId, long accountId, String excludeBucket) {
+ Account account = _accountDao.findById(accountId);
+ if (account == null) {
+ return;
+ }
+ String userName = getUserNameForAccount(account, storeId);
+ List buckets = _bucketDao.listByObjectStoreIdAndAccountId(storeId, accountId);
+ List bucketNames = new ArrayList<>();
+ for (BucketVO bvo : buckets) {
+ if (excludeBucket != null && excludeBucket.equals(bvo.getName())) {
+ continue;
+ }
+ // Skip buckets whose remote counterpart has been deleted but whose
+ // row is still present for resource accounting. Including them
+ // would re-grant a bucket name that is now free for another account
+ // to claim.
+ if (Bucket.State.Destroyed.equals(bvo.getState())) {
+ continue;
+ }
+ bucketNames.add(bvo.getName());
+ }
+ String policy = SeaweedFSObjectStoreUtil.buildAccountIAMPolicy(bucketNames);
+ iamClient.putUserPolicy(new PutUserPolicyRequest(userName,
+ SeaweedFSObjectStoreUtil.IAM_USER_POLICY_NAME, policy));
+ }
+
+ /**
+ * Check whether the given access key id is still listed and Active in IAM
+ * for the user. Listing failures are propagated rather than swallowed so
+ * a transient IAM outage does not send createUser into the replacement
+ * path (which would overwrite stored credentials and invalidate bucket
+ * records).
+ */
+ private boolean iamAccessKeyExists(AmazonIdentityManagement iamClient, String userName, String accessKeyId) {
+ for (AccessKeyMetadata metadata :
+ iamClient.listAccessKeys(new ListAccessKeysRequest()
+ .withUserName(userName)).getAccessKeyMetadata()) {
+ if (accessKeyId.equals(metadata.getAccessKeyId())) {
+ return "Active".equalsIgnoreCase(metadata.getStatus());
+ }
+ }
+ return false;
+ }
+
+ /**
+ * Delete access keys for the user other than the (optionally) preserved
+ * key id. Used to clean up unmanaged leftover keys before creating a
+ * replacement so repeated calls do not hit IAM access-key limits.
+ */
+ private void deleteUnmanagedAccessKeys(AmazonIdentityManagement iamClient, String userName, String preserveAccessKeyId) {
+ try {
+ for (AccessKeyMetadata metadata :
+ iamClient.listAccessKeys(new ListAccessKeysRequest()
+ .withUserName(userName)).getAccessKeyMetadata()) {
+ String keyId = metadata.getAccessKeyId();
+ if (preserveAccessKeyId != null && preserveAccessKeyId.equals(keyId)) {
+ continue;
+ }
+ DeleteAccessKeyRequest deleteReq =
+ new DeleteAccessKeyRequest()
+ .withUserName(userName)
+ .withAccessKeyId(keyId);
+ logger.info("Deleting un-managed IAM access key {} for user {}", keyId, userName);
+ iamClient.deleteAccessKey(deleteReq);
+ }
+ } catch (AmazonClientException e) {
+ // Propagate so the caller does not proceed to create a replacement
+ // key while stale unmanaged keys remain (which could hit IAM key
+ // limits or leave orphaned credentials).
+ throw new CloudRuntimeException("Failed to clean up IAM access keys for user " + userName, e);
+ }
+ }
+
+ @Override
+ public Bucket createBucket(Bucket bucket, boolean objectLock) {
+ String bucketName = bucket.getName();
+ long storeId = bucket.getObjectStoreId();
+
+ // Serialize against a concurrent deleteBucket of the same name by any
+ // account on this store. Bucket names are globally unique per store and
+ // reusable, so without this an account could claim a name while the
+ // previous owner's IAM policy still granted that ARN.
+ GlobalLock nameLock = acquireBucketNameLock(storeId, bucketName);
+ if (nameLock == null) {
+ throw new CloudRuntimeException("Failed to acquire bucket name lock for store " + storeId + " bucket " + bucketName);
+ }
+ try {
+ return createBucketLocked(bucket, objectLock);
+ } finally {
+ nameLock.unlock();
+ nameLock.releaseRef();
+ }
+ }
+
+ private Bucket createBucketLocked(Bucket bucket, boolean objectLock) {
+ String bucketName = bucket.getName();
+ long storeId = bucket.getObjectStoreId();
+ long accountId = bucket.getAccountId();
+
+ // Use the store's admin credentials to create the bucket
+ AmazonS3 s3client = getS3ClientByStoreId(storeId);
+
+ // Check if the bucket already exists
+ try {
+ if (s3client.doesBucketExistV2(bucketName)) {
+ throw new CloudRuntimeException("Bucket already exists with name " + bucketName);
+ }
+ } catch (AmazonClientException e) {
+ throw new CloudRuntimeException(e);
+ }
+
+ // Create the bucket
+ try {
+ CreateBucketRequest request = new CreateBucketRequest(bucketName);
+ if (objectLock) {
+ request.setObjectLockEnabledForBucket(true);
+ }
+ s3client.createBucket(request);
+ } catch (AmazonClientException e) {
+ logger.error("Create bucket failed", e);
+ throw new CloudRuntimeException(e);
+ }
+
+ // Step 2: update the bucket record with the account's IAM credentials,
+ // configure CORS, and refresh the IAM policy. If any of these fail,
+ // clean up the remote bucket so a retry does not find it already
+ // existing — mirroring the Cloudian createBucket pattern.
+ //
+ // Hold the IAM lock for the account so a concurrent createUser key
+ // rotation does not change the account credentials between reading
+ // them and writing the BucketVO, which would leave the bucket with
+ // a stale key pair.
+ GlobalLock iamLock = acquireIamLock(storeId, accountId);
+ if (iamLock == null) {
+ // The S3 bucket has already been created. Clean it up so a
+ // retry does not find it already existing, then throw.
+ try {
+ s3client.deleteBucket(bucketName);
+ } catch (AmazonClientException cleanupEx) {
+ logger.error("Failed to clean up bucket {} after IAM lock timeout", bucketName, cleanupEx);
+ }
+ throw new CloudRuntimeException("Failed to acquire IAM lock for store " + storeId + " account " + accountId);
+ }
+ try {
+ // Configure permissive CORS so the CloudStack S3 bucket browser
+ // (which performs list/upload/delete from the browser) can function.
+ // SeaweedFS supports the standard PutBucketCors operation.
+ configureBucketCORS(s3client, bucketName);
+
+ Map accountDetails = _accountDetailsDao.findDetails(accountId);
+ String accessKey = accountDetails.get(SeaweedFSObjectStoreUtil.keyAccessKey(storeId));
+ String secretKey = accountDetails.get(SeaweedFSObjectStoreUtil.keySecretKey(storeId));
+ if (accessKey == null || secretKey == null) {
+ throw new CloudRuntimeException("No IAM credentials found for account " + accountId
+ + " on store " + storeId + ". Run createUser before creating a bucket.");
+ }
+
+ String s3Url = getS3Url(storeId);
+ BucketVO bucketVO = _bucketDao.findById(bucket.getId());
+ bucketVO.setAccessKey(accessKey);
+ bucketVO.setSecretKey(secretKey);
+ // Normalize the endpoint: s3Url is operator-supplied and may carry a
+ // trailing slash, which would persist a broken "...//bucket" URL
+ // that BucketResponse and the object store browser both use.
+ bucketVO.setBucketURL(SeaweedFSObjectStoreUtil.stripTrailingSlashes(s3Url) + "/" + bucketName);
+ _bucketDao.update(bucket.getId(), bucketVO);
+
+ // Refresh the account's IAM policy to include the new bucket.
+ // Use the lock-free variant since createBucket already holds the
+ // IAM lock for the post-create section.
+ AmazonIdentityManagement iamClient = getIAMClient(storeId);
+ updateAccountIAMPolicyLocked(iamClient, storeId, accountId, null);
+
+ // Return the updated BucketVO (not the stale input bucket) so
+ // BucketApiServiceImpl.createBucket does not overwrite the
+ // persisted credentials with the stale values.
+ return bucketVO;
+ } catch (Exception e) {
+ logger.error("Post-create bucket record update failed for {}; cleaning up remote bucket", bucketName, e);
+ CloudRuntimeException primary = new CloudRuntimeException(e);
+ try {
+ s3client.deleteBucket(bucketName);
+ logger.info("Cleanup of bucket {} succeeded", bucketName);
+ } catch (AmazonClientException cleanupEx) {
+ logger.error("Cleanup of bucket {} also failed", bucketName, cleanupEx);
+ primary.addSuppressed(cleanupEx);
+ }
+ // Revoke the IAM policy grant for the new bucket so the account's
+ // credentials cannot access a bucket that no longer exists. If the
+ // policy PUT succeeded before the DB update failed, the grant
+ // would otherwise persist and could be reused if another account
+ // later creates the same bucket name. Use the lock-free variant
+ // since createBucket already holds the IAM lock. Propagate
+ // failures as suppressed exceptions so they are not silently lost.
+ try {
+ AmazonIdentityManagement iamClient = getIAMClient(storeId);
+ updateAccountIAMPolicyLocked(iamClient, storeId, accountId, bucketName);
+ } catch (Exception policyEx) {
+ logger.warn("Failed to revoke IAM policy for bucket {} after cleanup: {}", bucketName, policyEx.getMessage());
+ primary.addSuppressed(policyEx);
+ }
+ throw primary;
+ } finally {
+ iamLock.unlock();
+ iamLock.releaseRef();
+ }
+ }
+
+ /**
+ * Configure a permissive CORS policy on the bucket so the CloudStack
+ * S3 bucket browser (which performs list/upload/delete from the
+ * browser) can function. Mirrors the Cloudian configureBucketCORS.
+ */
+ private void configureBucketCORS(AmazonS3 s3client, String bucketName) {
+ logger.debug("Configuring CORS for bucket {}", bucketName);
+ List corsRules = new ArrayList<>();
+ CORSRule allowAnyRule = new CORSRule().withId("AllowAny");
+ allowAnyRule.setAllowedOrigins("*");
+ allowAnyRule.setAllowedHeaders("*");
+ allowAnyRule.setAllowedMethods(
+ CORSRule.AllowedMethods.HEAD,
+ CORSRule.AllowedMethods.GET,
+ CORSRule.AllowedMethods.PUT,
+ CORSRule.AllowedMethods.POST,
+ CORSRule.AllowedMethods.DELETE);
+ corsRules.add(allowAnyRule);
+ BucketCrossOriginConfiguration corsConfig = new BucketCrossOriginConfiguration();
+ corsConfig.setRules(corsRules);
+ SetBucketCrossOriginConfigurationRequest corsRequest = new SetBucketCrossOriginConfigurationRequest(bucketName, corsConfig);
+ s3client.setBucketCrossOriginConfiguration(corsRequest);
+ logger.info("Successfully configured CORS for bucket {}", bucketName);
+ }
+
+ @Override
+ public List listBuckets(long storeId) {
+ AmazonS3 s3client = getS3ClientByStoreId(storeId);
+ List bucketsList = new ArrayList<>();
+ try {
+ List s3Buckets = s3client.listBuckets();
+ for (com.amazonaws.services.s3.model.Bucket s3Bucket : s3Buckets) {
+ Bucket bucket = new BucketObject();
+ bucket.setName(s3Bucket.getName());
+ bucketsList.add(bucket);
+ }
+ } catch (AmazonClientException e) {
+ throw new CloudRuntimeException(e);
+ }
+ return bucketsList;
+ }
+
+ @Override
+ public boolean deleteBucket(BucketTO bucket, long storeId) {
+ String bucketName = bucket.getName();
+
+ // Hold the store-wide bucket name lock across the remote delete and the
+ // IAM policy refresh. The IAM lock below is scoped to (store, account)
+ // and therefore does not serialize against createBucket for a *different*
+ // account: once the S3 delete succeeded, that account could claim the now
+ // free name while this account's policy still granted the ARN, letting the
+ // old credentials reach the new tenant's bucket.
+ GlobalLock nameLock = acquireBucketNameLock(storeId, bucketName);
+ if (nameLock == null) {
+ throw new CloudRuntimeException("Failed to acquire bucket name lock for store " + storeId + " bucket " + bucketName);
+ }
+ try {
+ return deleteBucketLocked(bucket, storeId);
+ } finally {
+ nameLock.unlock();
+ nameLock.releaseRef();
+ }
+ }
+
+ private boolean deleteBucketLocked(BucketTO bucket, long storeId) {
+ String bucketName = bucket.getName();
+ long accountId = bucket.getAccountId();
+ AmazonS3 s3client = getS3ClientByStoreId(storeId);
+
+ // Acquire the per-store/account IAM lock so the S3 delete and the
+ // subsequent policy refresh are atomic with respect to concurrent
+ // createUser operations for this account.
+ GlobalLock iamLock = acquireIamLock(storeId, accountId);
+ if (iamLock == null) {
+ throw new CloudRuntimeException("Failed to acquire IAM lock for store " + storeId + " account " + accountId);
+ }
+ try {
+ // Delete the S3 bucket first. If this fails (non-empty bucket,
+ // transient error), the IAM policy is still intact so the user
+ // can empty the bucket and retry.
+ //
+ // If the bucket is already gone (e.g. from a previous partial
+ // failure where the S3 delete succeeded but the IAM policy refresh
+ // failed), skip the S3 delete so the retry is idempotent.
+ try {
+ if (s3client.doesBucketExistV2(bucketName)) {
+ s3client.deleteBucket(bucketName);
+ }
+ } catch (AmazonClientException e) {
+ throw new CloudRuntimeException(e);
+ }
+
+ // Mark the row Destroyed inside the lock. The row itself must
+ // survive so BucketApiServiceImpl.deleteCheckedBucket can still
+ // decrement the resource counts, update allocated size, and remove
+ // the row in one transaction after this returns. If that cleanup
+ // fails, the Destroyed marker keeps retries from re-granting a
+ // bucket name whose remote bucket is already gone.
+ // Marking it (rather than relying only on excludeBucket) closes the
+ // window where a concurrent createUser/createBucket for this
+ // account rebuilds the policy from the DB and re-adds this bucket's
+ // ARN: buildAccountIAMPolicy skips Destroyed rows, so the grant
+ // cannot come back on a name that is now reusable.
+ markBucketDestroyed(storeId, accountId, bucketName);
+
+ // Refresh the account's IAM policy to drop the deleted bucket.
+ // Bucket names are reusable, so a stale grant would let the old
+ // account access a new tenant's bucket with the same name. This
+ // must succeed; if it fails the exception propagates and the
+ // caller can retry (the S3 delete is skipped when the bucket is
+ // already gone, and the Destroyed marker is idempotent). The
+ // lock-free variant is used because deleteBucket holds the lock.
+ AmazonIdentityManagement iamClient = getIAMClient(storeId);
+ updateAccountIAMPolicyLocked(iamClient, storeId, accountId, bucketName);
+
+ return true;
+ } finally {
+ iamLock.unlock();
+ iamLock.releaseRef();
+ }
+ }
+
+ @Override
+ public AccessControlList getBucketAcl(BucketTO bucket, long storeId) {
+ AmazonS3 s3client = getS3ClientByStoreId(storeId);
+ try {
+ return s3client.getBucketAcl(bucket.getName());
+ } catch (AmazonClientException e) {
+ throw new CloudRuntimeException(e);
+ }
+ }
+
+ @Override
+ public void setBucketAcl(BucketTO bucket, AccessControlList acl, long storeId) {
+ AmazonS3 s3client = getS3ClientByStoreId(storeId);
+ try {
+ s3client.setBucketAcl(bucket.getName(), acl);
+ } catch (AmazonClientException e) {
+ throw new CloudRuntimeException(e);
+ }
+ }
+
+ @Override
+ public void setBucketPolicy(BucketTO bucket, String policy, long storeId) {
+ if ("private".equalsIgnoreCase(policy)) {
+ deleteBucketPolicy(bucket, storeId);
+ return;
+ }
+
+ StringBuilder sb = new StringBuilder();
+ sb.append("{\n");
+ sb.append(" \"Version\": \"2012-10-17\",\n");
+ sb.append(" \"Statement\": [\n");
+ sb.append(" {\n");
+ sb.append(" \"Sid\": \"PublicReadForObjects\",\n");
+ sb.append(" \"Effect\": \"Allow\",\n");
+ sb.append(" \"Principal\": \"*\",\n");
+ sb.append(" \"Action\": \"s3:GetObject\",\n");
+ sb.append(" \"Resource\": \"arn:aws:s3:::%s/*\"\n");
+ sb.append(" }\n");
+ sb.append(" ]\n");
+ sb.append("}\n");
+
+ String jsonPolicy = String.format(sb.toString(), bucket.getName());
+ AmazonS3 s3client = getS3ClientByStoreId(storeId);
+ try {
+ s3client.setBucketPolicy(bucket.getName(), jsonPolicy);
+ } catch (AmazonClientException e) {
+ throw new CloudRuntimeException(e);
+ }
+ }
+
+ @Override
+ public BucketPolicy getBucketPolicy(BucketTO bucket, long storeId) {
+ AmazonS3 s3client = getS3ClientByStoreId(storeId);
+ try {
+ return s3client.getBucketPolicy(new GetBucketPolicyRequest(bucket.getName()));
+ } catch (AmazonClientException e) {
+ throw new CloudRuntimeException(e);
+ }
+ }
+
+ @Override
+ public void deleteBucketPolicy(BucketTO bucket, long storeId) {
+ AmazonS3 s3client = getS3ClientByStoreId(storeId);
+ try {
+ s3client.deleteBucketPolicy(new DeleteBucketPolicyRequest(bucket.getName()));
+ } catch (AmazonClientException e) {
+ throw new CloudRuntimeException(e);
+ }
+ }
+
+ @Override
+ public boolean setBucketEncryption(BucketTO bucket, long storeId) {
+ AmazonS3 s3client = getS3ClientByStoreId(storeId);
+ try {
+ SetBucketEncryptionRequest eRequest = new SetBucketEncryptionRequest();
+ eRequest.setBucketName(bucket.getName());
+
+ ServerSideEncryptionByDefault sseByDefault = new ServerSideEncryptionByDefault();
+ sseByDefault.setSSEAlgorithm(SSEAlgorithm.AES256.toString());
+
+ ServerSideEncryptionRule sseRule = new ServerSideEncryptionRule();
+ sseRule.setApplyServerSideEncryptionByDefault(sseByDefault);
+
+ List sseRules = new ArrayList<>();
+ sseRules.add(sseRule);
+
+ ServerSideEncryptionConfiguration sseConf = new ServerSideEncryptionConfiguration();
+ sseConf.setRules(sseRules);
+
+ eRequest.setServerSideEncryptionConfiguration(sseConf);
+ s3client.setBucketEncryption(eRequest);
+ return true;
+ } catch (AmazonClientException e) {
+ throw new CloudRuntimeException(e);
+ }
+ }
+
+ @Override
+ public boolean deleteBucketEncryption(BucketTO bucket, long storeId) {
+ AmazonS3 s3client = getS3ClientByStoreId(storeId);
+ try {
+ s3client.deleteBucketEncryption(bucket.getName());
+ return true;
+ } catch (AmazonClientException e) {
+ throw new CloudRuntimeException(e);
+ }
+ }
+
+ @Override
+ public boolean setBucketVersioning(BucketTO bucket, long storeId) {
+ AmazonS3 s3client = getS3ClientByStoreId(storeId);
+ try {
+ BucketVersioningConfiguration vConf = new BucketVersioningConfiguration(BucketVersioningConfiguration.ENABLED);
+ s3client.setBucketVersioningConfiguration(
+ new SetBucketVersioningConfigurationRequest(bucket.getName(), vConf));
+ return true;
+ } catch (AmazonClientException e) {
+ throw new CloudRuntimeException(e);
+ }
+ }
+
+ @Override
+ public boolean deleteBucketVersioning(BucketTO bucket, long storeId) {
+ AmazonS3 s3client = getS3ClientByStoreId(storeId);
+ try {
+ BucketVersioningConfiguration vConf = new BucketVersioningConfiguration(BucketVersioningConfiguration.SUSPENDED);
+ s3client.setBucketVersioningConfiguration(
+ new SetBucketVersioningConfigurationRequest(bucket.getName(), vConf));
+ return true;
+ } catch (AmazonClientException e) {
+ throw new CloudRuntimeException(e);
+ }
+ }
+
+ /**
+ * Set the bucket quota via the SeaweedFS S3 extension
+ * ({@code PUT /{bucket}?seaweedfs-quota}), signed with the store admin
+ * S3 credentials.
+ *
+ * SeaweedFS enforces bucket quota server-side by setting a read-only flag
+ * when usage exceeds the configured limit. The quota is configured via the
+ * SeaweedFS S3 extension endpoint PUT /{bucket}?seaweedfs-quota,
+ * authenticated via standard S3 SigV4 and authorized via the
+ * s3:PutBucketQuota IAM permission.
+ *
+ * @param size the GiB size to set the quota to. 0 disables quota.
+ * @throws CloudRuntimeException if the S3 endpoint or credentials are missing or the request fails.
+ */
+ @Override
+ public void setBucketQuota(BucketTO bucket, long storeId, long size) {
+ String s3Url = getS3Url(storeId);
+ String accessKey = getAccessKey(storeId);
+ String secretKey = getSecretKey(storeId);
+ if (s3Url == null || s3Url.isEmpty() || accessKey == null || accessKey.isEmpty() || secretKey == null || secretKey.isEmpty()) {
+ throw new CloudRuntimeException("SeaweedFS S3 URL and credentials are required to set bucket quota. " +
+ "Configure 's3Url', 'accesskey', and 'secretkey' in the object store details.");
+ }
+ // A 404/405 from the optional quota extension is only tolerable on the
+ // initial create, where CreateBucketCmd always calls setQuota and there
+ // is no existing quota to clear. Every other call - an update, or the
+ // rollback of a failed update - must propagate the failure so
+ // CloudStack accounting cannot diverge from SeaweedFS state.
+ //
+ // The initial create is identified by the BucketVO still being in the
+ // Allocated state: BucketApiServiceImpl.createBucket only promotes it to
+ // Created after setQuota returns. Keying off the previous quota value
+ // instead would misclassify an update from quota 0, and would wrongly
+ // tolerate a 404 while rolling such an update back even though the
+ // successful positive update had just proved the extension exists.
+ boolean initialCreate = isBucketAllocated(storeId, bucket);
+ SeaweedFSObjectStoreUtil.setBucketQuotaViaS3Extension(s3Url, accessKey, secretKey, bucket.getName(), size,
+ getS3ExtensionHttpClient(), initialCreate);
+ }
+
+ /**
+ * Mark the BucketVO for this bucket as {@link Bucket.State#Destroyed} so
+ * {@link #updateAccountIAMPolicyLocked} stops granting its ARN while the
+ * row is still present for the caller's resource accounting. Must be
+ * called while holding the store/account IAM lock.
+ */
+ protected void markBucketDestroyed(long storeId, long accountId, String bucketName) {
+ for (BucketVO bvo : _bucketDao.listByObjectStoreIdAndAccountId(storeId, accountId)) {
+ if (bucketName.equals(bvo.getName())) {
+ if (!Bucket.State.Destroyed.equals(bvo.getState())) {
+ bvo.setState(Bucket.State.Destroyed);
+ _bucketDao.update(bvo.getId(), bvo);
+ }
+ return;
+ }
+ }
+ }
+
+ /**
+ * Returns true when the persisted BucketVO for this bucket is still in the
+ * {@link Bucket.State#Allocated} state, i.e. the bucket is being created and
+ * has not yet been promoted to {@link Bucket.State#Created}. Used to
+ * recognise the initial create, which is the only call permitted to tolerate
+ * a missing optional quota extension.
+ */
+ protected boolean isBucketAllocated(long storeId, BucketTO bucket) {
+ for (BucketVO bvo : _bucketDao.listByObjectStoreIdAndAccountId(storeId, bucket.getAccountId())) {
+ if (bucket.getName().equals(bvo.getName())) {
+ return Bucket.State.Allocated.equals(bvo.getState());
+ }
+ }
+ return false;
+ }
+
+ /**
+ * Returns the HTTP client used to send SeaweedFS S3 extension requests
+ * (e.g. PUT /{bucket}?seaweedfs-quota). Exposed as a protected seam so
+ * tests can inject a mock client and assert the signed request without
+ * touching the network.
+ */
+ protected java.net.http.HttpClient getS3ExtensionHttpClient() {
+ return SeaweedFSObjectStoreUtil.newS3ExtensionHttpClient();
+ }
+
+ @Override
+ public Map getAllBucketsUsage(long storeId) {
+ Map bucketUsage = new HashMap<>();
+ List bucketList = _bucketDao.listByObjectStoreId(storeId);
+ if (bucketList.isEmpty()) {
+ return bucketUsage;
+ }
+
+ // If the operator has configured a Prometheus metricsUrl, scrape
+ // per-bucket sizes from the /metrics endpoint in a single HTTP GET.
+ // This is O(buckets) and avoids the O(total objects) ListObjectsV2
+ // scan that doesn't scale to large deployments. Falls back to
+ // ListObjectsV2 when metricsUrl is not configured.
+ String metricsUrl = getMetricsUrl(storeId);
+ if (metricsUrl != null) {
+ java.util.Set bucketNames = new java.util.HashSet<>();
+ for (BucketVO bucket : bucketList) {
+ bucketNames.add(bucket.getName());
+ }
+ try {
+ return SeaweedFSObjectStoreUtil.parseBucketUsageFromMetrics(
+ metricsUrl, bucketNames, getS3ExtensionHttpClient());
+ } catch (CloudRuntimeException e) {
+ logger.warn("Prometheus metrics scrape failed for store {}; falling back to ListObjectsV2", storeId, e);
+ }
+ }
+
+ // Fallback: list objects per bucket via S3. This is O(total objects)
+ // and does not scale to large deployments; configure metricsUrl for
+ // production usage reporting. ListObjectsV2 only counts current
+ // object versions; noncurrent versions and delete markers are omitted.
+ AmazonS3 s3client = getS3ClientByStoreId(storeId);
+ for (BucketVO bucket : bucketList) {
+ try {
+ long size = 0L;
+ com.amazonaws.services.s3.model.ListObjectsV2Result result;
+ String continuationToken = null;
+ do {
+ com.amazonaws.services.s3.model.ListObjectsV2Request req =
+ new com.amazonaws.services.s3.model.ListObjectsV2Request()
+ .withBucketName(bucket.getName())
+ .withMaxKeys(1000);
+ if (continuationToken != null) {
+ req.setContinuationToken(continuationToken);
+ }
+ result = s3client.listObjectsV2(req);
+ for (com.amazonaws.services.s3.model.S3ObjectSummary summary : result.getObjectSummaries()) {
+ size += summary.getSize();
+ }
+ continuationToken = result.getNextContinuationToken();
+ } while (result.isTruncated());
+ bucketUsage.put(bucket.getName(), size);
+ } catch (AmazonClientException e) {
+ // Propagate the failure so BucketApiServiceImpl does not
+ // overwrite objectStoreVO.usedSize with a partial total.
+ // Returning only the successful buckets would under-report
+ // store usage and trigger false capacity alerts.
+ throw new CloudRuntimeException("Failed to get usage for bucket " + bucket.getName(), e);
+ }
+ }
+ return bucketUsage;
+ }
+
+ // ---- Client builders ----
+
+ protected String getS3Url(long storeId) {
+ // Read the configured S3 endpoint from the persisted details first
+ // (it may differ from the generic ObjectStoreVO.url), falling back to
+ // the store URL only if the detail is missing. This matches the
+ // Cloudian HyperStore pattern.
+ Map storeDetails = _storeDetailsDao.getDetails(storeId);
+ String s3Url = storeDetails.get(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_S3_URL);
+ if (s3Url == null || s3Url.isEmpty()) {
+ ObjectStoreVO store = _storeDao.findById(storeId);
+ if (store != null) {
+ s3Url = store.getUrl();
+ }
+ }
+ return s3Url;
+ }
+
+ protected String getIAMUrl(long storeId) {
+ Map storeDetails = _storeDetailsDao.getDetails(storeId);
+ String iamUrl = storeDetails.get(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_IAM_URL);
+ if (iamUrl == null || iamUrl.isEmpty()) {
+ // iamUrl was not explicitly configured; SeaweedFS serves the IAM
+ // API from the same endpoint as S3 by default, so fall back to
+ // the current S3 endpoint. This also keeps IAM provisioning on the
+ // live endpoint after updateObjectStore changes the store URL.
+ iamUrl = getS3Url(storeId);
+ }
+ return iamUrl;
+ }
+
+ protected String getAccessKey(long storeId) {
+ Map storeDetails = _storeDetailsDao.getDetails(storeId);
+ return storeDetails.get(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_ACCESS_KEY);
+ }
+
+ protected String getSecretKey(long storeId) {
+ Map storeDetails = _storeDetailsDao.getDetails(storeId);
+ return storeDetails.get(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_SECRET_KEY);
+ }
+
+ /**
+ * Returns the configured Prometheus metrics endpoint URL for the store,
+ * or {@code null} if not configured. When set, {@link #getAllBucketsUsage}
+ * scrapes per-bucket sizes from this endpoint instead of listing every
+ * object via S3 ListObjectsV2.
+ *
+ * The URL must point at a single SeaweedFS S3 server's Prometheus exporter
+ * (the address configured with {@code -metricsPort}), not at a Prometheus
+ * server and not at a load-balanced S3 service: SeaweedFS only refreshes
+ * the bucket-size gauges on the instance holding the {@code s3.leader}
+ * lock. A wrong endpoint is detected by the scrape validation and falls
+ * back to the S3 listing. See
+ * {@link SeaweedFSObjectStoreUtil#parseBucketUsageFromMetrics}.
+ */
+ protected String getMetricsUrl(long storeId) {
+ Map storeDetails = _storeDetailsDao.getDetails(storeId);
+ String metricsUrl = storeDetails.get(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_METRICS_URL);
+ if (metricsUrl == null || metricsUrl.isEmpty()) {
+ return null;
+ }
+ return metricsUrl;
+ }
+
+ protected AmazonS3 getS3ClientByStoreId(long storeId) {
+ String s3Url = getS3Url(storeId);
+ Map storeDetails = _storeDetailsDao.getDetails(storeId);
+ String accessKey = storeDetails.get(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_ACCESS_KEY);
+ String secretKey = storeDetails.get(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_SECRET_KEY);
+ return SeaweedFSObjectStoreUtil.getS3Client(s3Url, accessKey, secretKey);
+ }
+
+ protected AmazonIdentityManagement getIAMClient(long storeId) {
+ String iamUrl = getIAMUrl(storeId);
+ Map storeDetails = _storeDetailsDao.getDetails(storeId);
+ String accessKey = storeDetails.get(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_ACCESS_KEY);
+ String secretKey = storeDetails.get(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_SECRET_KEY);
+ return SeaweedFSObjectStoreUtil.getIAMClient(iamUrl, accessKey, secretKey);
+ }
+}
diff --git a/plugins/storage/object/seaweedfs/src/main/java/org/apache/cloudstack/storage/datastore/lifecycle/SeaweedFSObjectStoreLifeCycleImpl.java b/plugins/storage/object/seaweedfs/src/main/java/org/apache/cloudstack/storage/datastore/lifecycle/SeaweedFSObjectStoreLifeCycleImpl.java
new file mode 100644
index 000000000000..02937194568f
--- /dev/null
+++ b/plugins/storage/object/seaweedfs/src/main/java/org/apache/cloudstack/storage/datastore/lifecycle/SeaweedFSObjectStoreLifeCycleImpl.java
@@ -0,0 +1,198 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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.
+ */
+// SPDX-License-Identifier: Apache-2.0
+package org.apache.cloudstack.storage.datastore.lifecycle;
+
+import com.cloud.agent.api.StoragePoolInfo;
+import com.cloud.hypervisor.Hypervisor.HypervisorType;
+import com.cloud.utils.exception.CloudRuntimeException;
+
+import org.apache.cloudstack.engine.subsystem.api.storage.ClusterScope;
+import org.apache.cloudstack.engine.subsystem.api.storage.DataStore;
+import org.apache.cloudstack.engine.subsystem.api.storage.HostScope;
+import org.apache.cloudstack.engine.subsystem.api.storage.ZoneScope;
+import org.apache.cloudstack.storage.datastore.db.ObjectStoreVO;
+import org.apache.cloudstack.storage.datastore.util.SeaweedFSObjectStoreUtil;
+import org.apache.cloudstack.storage.object.datastore.ObjectStoreHelper;
+import org.apache.cloudstack.storage.object.datastore.ObjectStoreProviderManager;
+import org.apache.cloudstack.storage.object.store.lifecycle.ObjectStoreLifeCycle;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import javax.inject.Inject;
+
+import java.util.HashMap;
+import java.util.Map;
+
+public class SeaweedFSObjectStoreLifeCycleImpl implements ObjectStoreLifeCycle {
+
+ protected Logger logger = LogManager.getLogger(SeaweedFSObjectStoreLifeCycleImpl.class);
+
+ @Inject
+ ObjectStoreHelper objectStoreHelper;
+ @Inject
+ ObjectStoreProviderManager objectStoreMgr;
+
+ public SeaweedFSObjectStoreLifeCycleImpl() {
+ }
+
+ @Override
+ public DataStore initialize(Map dsInfos) {
+
+ String name = (String)dsInfos.get(SeaweedFSObjectStoreUtil.STORE_KEY_NAME);
+ String url = (String)dsInfos.get(SeaweedFSObjectStoreUtil.STORE_KEY_URL);
+ String providerName = (String)dsInfos.get(SeaweedFSObjectStoreUtil.STORE_KEY_PROVIDER_NAME);
+ Long size = (Long)dsInfos.get(SeaweedFSObjectStoreUtil.STORE_KEY_SIZE);
+
+ // Check the providerName is what we expect
+ if (! StringUtils.equalsIgnoreCase(providerName, SeaweedFSObjectStoreUtil.OBJECT_STORE_PROVIDER_NAME)) {
+ String msg = String.format("Unexpected providerName \"%s\". Expected \"%s\"", providerName, SeaweedFSObjectStoreUtil.OBJECT_STORE_PROVIDER_NAME);
+ logger.error(msg);
+ throw new CloudRuntimeException(msg);
+ }
+
+ Map objectStoreParameters = new HashMap();
+ objectStoreParameters.put(SeaweedFSObjectStoreUtil.STORE_KEY_NAME, name);
+ objectStoreParameters.put(SeaweedFSObjectStoreUtil.STORE_KEY_URL, url);
+ objectStoreParameters.put(SeaweedFSObjectStoreUtil.STORE_KEY_PROVIDER_NAME, providerName);
+ objectStoreParameters.put(SeaweedFSObjectStoreUtil.STORE_KEY_SIZE, size);
+
+ // Pull out the details map
+ @SuppressWarnings("unchecked")
+ Map details = (Map) dsInfos.get(SeaweedFSObjectStoreUtil.STORE_KEY_DETAILS);
+ if (details == null) {
+ String msg = String.format("Unexpected null receiving Object Store initialization \"%s\"", SeaweedFSObjectStoreUtil.STORE_KEY_DETAILS);
+ logger.error(msg);
+ throw new CloudRuntimeException(msg);
+ }
+
+ // The admin/root access key and secret key are available as accesskey/secretkey
+ String accessKey = details.get(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_ACCESS_KEY);
+ String secretKey = details.get(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_SECRET_KEY);
+ String s3Url = details.get(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_S3_URL);
+ String iamUrl = details.get(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_IAM_URL);
+ String metricsUrl = details.get(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_METRICS_URL);
+
+ // Track whether the endpoints were explicitly supplied. Only
+ // explicitly-supplied values are persisted in the details map; the
+ // driver falls back to ObjectStoreVO.url / getS3Url() when the
+ // details are absent, so updateObjectStore can change the store URL
+ // without a stale persisted s3Url/iamUrl overriding it.
+ // StorageManagerImpl.updateObjectStore rewrites the stored
+ // BucketVO.bucketURL values on a URL change so the object-store
+ // browser does not keep targeting the old endpoint.
+ boolean s3UrlExplicit = StringUtils.isNotBlank(s3Url);
+ boolean iamUrlExplicit = StringUtils.isNotBlank(iamUrl);
+
+ // Resolve the endpoints for validation, defaulting to the store URL
+ // (and s3Url) as needed.
+ if (StringUtils.isBlank(s3Url)) {
+ s3Url = url;
+ }
+ // If iamUrl is not provided, default it to the s3Url.
+ // SeaweedFS registers its embedded IAM API at POST / on the same S3
+ // endpoint (UnifiedPostHandler), so the IAM endpoint is the same as
+ // the S3 endpoint unless the deployment runs a separate weed iam server.
+ if (StringUtils.isBlank(iamUrl)) {
+ iamUrl = s3Url;
+ }
+
+ if (StringUtils.isAnyBlank(accessKey, secretKey, s3Url, iamUrl)) {
+ final String asteriskPassword = (secretKey == null) ? null : "*".repeat(secretKey.length());
+ logger.error("Required parameters are missing; accessKey={} secretKey={} s3Url={} iamUrl={}",
+ accessKey, asteriskPassword, s3Url, iamUrl);
+ throw new CloudRuntimeException("Required SeaweedFS configuration parameters are missing/empty.");
+ }
+
+ // Persist only explicitly-supplied endpoint overrides. Defaulted
+ // values are not written so the driver resolves them from the current
+ // store URL at runtime.
+ if (s3UrlExplicit) {
+ details.put(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_S3_URL, s3Url);
+ } else {
+ details.remove(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_S3_URL);
+ }
+ if (iamUrlExplicit) {
+ details.put(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_IAM_URL, iamUrl);
+ } else {
+ details.remove(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_IAM_URL);
+ }
+ // Preserve the optional Prometheus metrics endpoint so getMetricsUrl()
+ // can find it; without this the scalable usage-reporting path is never
+ // used and every poll falls back to the ListObjectsV2 scan.
+ if (StringUtils.isNotBlank(metricsUrl)) {
+ details.put(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_METRICS_URL, metricsUrl);
+ } else {
+ details.remove(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_METRICS_URL);
+ }
+
+ // Validate the S3 and IAM endpoints behave like the respective
+ // services, then authenticate with the supplied admin credentials so a
+ // store with bad credentials is rejected here rather than failing on
+ // the first bucket or IAM operation.
+ logger.info("Validating SeaweedFS S3 endpoint: {}", s3Url);
+ SeaweedFSObjectStoreUtil.validateS3Url(s3Url);
+ logger.info("Validating SeaweedFS IAM endpoint: {}", iamUrl);
+ SeaweedFSObjectStoreUtil.validateIAMUrl(iamUrl);
+ logger.info("Validating SeaweedFS admin credentials");
+ SeaweedFSObjectStoreUtil.validateCredentials(s3Url, iamUrl, accessKey, secretKey);
+
+ logger.info("Successfully validated SeaweedFS object store: {} (quota management via S3 ?seaweedfs-quota extension)", name);
+
+ ObjectStoreVO objectStore = objectStoreHelper.createObjectStore(objectStoreParameters, details);
+ return objectStoreMgr.getObjectStore(objectStore.getId());
+ }
+
+ @Override
+ public boolean attachCluster(DataStore store, ClusterScope scope) {
+ return false;
+ }
+
+ @Override
+ public boolean attachHost(DataStore store, HostScope scope, StoragePoolInfo existingInfo) {
+ return false;
+ }
+
+ @Override
+ public boolean attachZone(DataStore dataStore, ZoneScope scope, HypervisorType hypervisorType) {
+ return false;
+ }
+
+ @Override
+ public boolean maintain(DataStore store) {
+ return false;
+ }
+
+ @Override
+ public boolean cancelMaintain(DataStore store) {
+ return false;
+ }
+
+ @Override
+ public boolean deleteDataStore(DataStore store) {
+ return false;
+ }
+
+ @Override
+ public boolean migrateToObjectStore(DataStore store) {
+ return false;
+ }
+
+}
diff --git a/plugins/storage/object/seaweedfs/src/main/java/org/apache/cloudstack/storage/datastore/provider/SeaweedFSObjectStoreProviderImpl.java b/plugins/storage/object/seaweedfs/src/main/java/org/apache/cloudstack/storage/datastore/provider/SeaweedFSObjectStoreProviderImpl.java
new file mode 100644
index 000000000000..dac7adb38b6e
--- /dev/null
+++ b/plugins/storage/object/seaweedfs/src/main/java/org/apache/cloudstack/storage/datastore/provider/SeaweedFSObjectStoreProviderImpl.java
@@ -0,0 +1,87 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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.
+ */
+// SPDX-License-Identifier: Apache-2.0
+package org.apache.cloudstack.storage.datastore.provider;
+
+import com.cloud.utils.component.ComponentContext;
+import org.apache.cloudstack.engine.subsystem.api.storage.DataStoreDriver;
+import org.apache.cloudstack.engine.subsystem.api.storage.DataStoreLifeCycle;
+import org.apache.cloudstack.engine.subsystem.api.storage.HypervisorHostListener;
+import org.apache.cloudstack.engine.subsystem.api.storage.ObjectStoreProvider;
+import org.apache.cloudstack.storage.datastore.driver.SeaweedFSObjectStoreDriverImpl;
+import org.apache.cloudstack.storage.datastore.lifecycle.SeaweedFSObjectStoreLifeCycleImpl;
+import org.apache.cloudstack.storage.datastore.util.SeaweedFSObjectStoreUtil;
+import org.apache.cloudstack.storage.object.ObjectStoreDriver;
+import org.apache.cloudstack.storage.object.datastore.ObjectStoreHelper;
+import org.apache.cloudstack.storage.object.datastore.ObjectStoreProviderManager;
+import org.apache.cloudstack.storage.object.store.lifecycle.ObjectStoreLifeCycle;
+import org.springframework.stereotype.Component;
+
+import javax.inject.Inject;
+import java.util.HashSet;
+import java.util.Map;
+import java.util.Set;
+
+@Component
+public class SeaweedFSObjectStoreProviderImpl implements ObjectStoreProvider {
+
+ @Inject
+ ObjectStoreProviderManager storeMgr;
+ @Inject
+ ObjectStoreHelper helper;
+
+ private final String providerName = SeaweedFSObjectStoreUtil.OBJECT_STORE_PROVIDER_NAME;
+ protected ObjectStoreLifeCycle lifeCycle;
+ protected ObjectStoreDriver driver;
+
+ @Override
+ public DataStoreLifeCycle getDataStoreLifeCycle() {
+ return lifeCycle;
+ }
+
+ @Override
+ public String getName() {
+ return this.providerName;
+ }
+
+ @Override
+ public boolean configure(Map params) {
+ lifeCycle = ComponentContext.inject(SeaweedFSObjectStoreLifeCycleImpl.class);
+ driver = ComponentContext.inject(SeaweedFSObjectStoreDriverImpl.class);
+ storeMgr.registerDriver(this.getName(), driver);
+ return true;
+ }
+
+ @Override
+ public DataStoreDriver getDataStoreDriver() {
+ return this.driver;
+ }
+
+ @Override
+ public HypervisorHostListener getHostListener() {
+ return null;
+ }
+
+ @Override
+ public Set getTypes() {
+ Set types = new HashSet();
+ types.add(DataStoreProviderType.OBJECT);
+ return types;
+ }
+}
diff --git a/plugins/storage/object/seaweedfs/src/main/java/org/apache/cloudstack/storage/datastore/util/SeaweedFSObjectStoreUtil.java b/plugins/storage/object/seaweedfs/src/main/java/org/apache/cloudstack/storage/datastore/util/SeaweedFSObjectStoreUtil.java
new file mode 100644
index 000000000000..50aa16a65365
--- /dev/null
+++ b/plugins/storage/object/seaweedfs/src/main/java/org/apache/cloudstack/storage/datastore/util/SeaweedFSObjectStoreUtil.java
@@ -0,0 +1,686 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you 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.
+// SPDX-License-Identifier: Apache-2.0
+package org.apache.cloudstack.storage.datastore.util;
+
+import org.apache.commons.lang3.StringUtils;
+
+import com.amazonaws.AmazonServiceException;
+import com.amazonaws.auth.AWSStaticCredentialsProvider;
+import com.amazonaws.auth.BasicAWSCredentials;
+import com.amazonaws.client.builder.AwsClientBuilder;
+import com.amazonaws.services.identitymanagement.AmazonIdentityManagement;
+import com.amazonaws.services.identitymanagement.AmazonIdentityManagementClientBuilder;
+import com.amazonaws.services.s3.AmazonS3;
+import com.amazonaws.services.s3.AmazonS3ClientBuilder;
+import com.cloud.utils.exception.CloudRuntimeException;
+
+/**
+ * Utility class for the SeaweedFS object storage provider.
+ *
+ * SeaweedFS exposes both an S3-compatible API and an AWS IAM-compatible API,
+ * so this provider needs no proprietary admin client — only the AWS S3 and IAM
+ * SDKs, the same pair Cloudian HyperStore already uses in this tree.
+ */
+public class SeaweedFSObjectStoreUtil {
+
+ /** The name of our Object Store Provider */
+ public static final String OBJECT_STORE_PROVIDER_NAME = "SeaweedFS";
+
+ public static final String STORE_KEY_PROVIDER_NAME = "providerName";
+ public static final String STORE_KEY_URL = "url";
+ public static final String STORE_KEY_NAME = "name";
+ public static final String STORE_KEY_SIZE = "size";
+ public static final String STORE_KEY_DETAILS = "details";
+
+ // Store Details Map key names - managed outside of plugin
+ public static final String STORE_DETAILS_KEY_ACCESS_KEY = "accesskey"; // admin/root access key
+ public static final String STORE_DETAILS_KEY_SECRET_KEY = "secretkey"; // admin/root secret key
+ public static final String STORE_DETAILS_KEY_S3_URL = "s3Url"; // S3 endpoint URL
+ public static final String STORE_DETAILS_KEY_IAM_URL = "iamUrl"; // IAM endpoint URL
+ public static final String STORE_DETAILS_KEY_METRICS_URL = "metricsUrl"; // Prometheus metrics endpoint URL (optional, for scalable usage reporting)
+
+ // Account Detail Map key names - credentials created per CloudStack account.
+ // Namespaced by store ID so one account can use multiple SeaweedFS pools
+ // without the second pool overwriting the first pool's credentials.
+ public static final String KEY_ACCESS_KEY_PREFIX = "swfs_AccessKey_";
+ public static final String KEY_SECRET_KEY_PREFIX = "swfs_SecretKey_";
+
+ /**
+ * Build the account-detail key for the IAM access key of a given store.
+ */
+ public static String keyAccessKey(long storeId) {
+ return KEY_ACCESS_KEY_PREFIX + storeId;
+ }
+
+ /**
+ * Build the account-detail key for the IAM secret key of a given store.
+ */
+ public static String keySecretKey(long storeId) {
+ return KEY_SECRET_KEY_PREFIX + storeId;
+ }
+
+ /**
+ * Strip trailing slashes from a configured endpoint URL so callers can
+ * append a {@code "/" + path} suffix without producing a double slash.
+ * The S3 and IAM endpoint URLs are operator-supplied and are accepted with
+ * or without a trailing slash.
+ *
+ * @param url the URL to normalize, may be null
+ * @return the URL without trailing slashes, or the input if null/empty
+ */
+ public static String stripTrailingSlashes(String url) {
+ if (url == null) {
+ return null;
+ }
+ String normalized = url;
+ while (normalized.endsWith("/")) {
+ normalized = normalized.substring(0, normalized.length() - 1);
+ }
+ return normalized;
+ }
+
+ /**
+ * Connect timeout for the S3 extension HTTP client, in seconds.
+ */
+ public static final int S3_EXTENSION_CONNECT_TIMEOUT_SECONDS = 10;
+ /**
+ * Per-request timeout for the S3 extension HTTP request, in seconds.
+ */
+ public static final int S3_EXTENSION_REQUEST_TIMEOUT_SECONDS = 30;
+
+ /**
+ * IAM user policy name applied to each per-account IAM user.
+ */
+ public static final String IAM_USER_POLICY_NAME = "CloudStackPolicy";
+
+ /**
+ * Build an IAM user policy that grants full S3 access only to the given
+ * buckets (both the bucket and its contents), while denying bucket
+ * creation and deletion everywhere so CloudStack retains control of the
+ * bucket lifecycle. When no buckets are provided, all S3 access is denied.
+ *
+ * This is the tenant boundary: each account's IAM credentials can only
+ * operate on that account's own buckets, not on every bucket in the
+ * SeaweedFS pool. The policy is refreshed whenever buckets are created or
+ * deleted (see
+ * {@code SeaweedFSObjectStoreDriverImpl.updateAccountIAMPolicy}).
+ *
+ * @param bucketNames the bucket names the account is allowed to access
+ * @return a JSON IAM policy document
+ */
+ public static String buildAccountIAMPolicy(java.util.List bucketNames) {
+ StringBuilder sb = new StringBuilder();
+ sb.append("{\n");
+ sb.append(" \"Version\": \"2012-10-17\",\n");
+ sb.append(" \"Statement\": [\n");
+ if (bucketNames == null || bucketNames.isEmpty()) {
+ // No buckets: deny all S3 access. A Resource cannot be empty in
+ // an IAM policy, so deny everything explicitly.
+ sb.append(" {\n");
+ sb.append(" \"Sid\": \"DenyAllS3\",\n");
+ sb.append(" \"Effect\": \"Deny\",\n");
+ sb.append(" \"Action\": [\"s3:*\"],\n");
+ sb.append(" \"Resource\": [\"arn:aws:s3:::*\", \"arn:aws:s3:::*/*\"]\n");
+ sb.append(" }\n");
+ } else {
+ sb.append(" {\n");
+ sb.append(" \"Sid\": \"AllowAccountBuckets\",\n");
+ sb.append(" \"Effect\": \"Allow\",\n");
+ sb.append(" \"Action\": [\"s3:*\"],\n");
+ sb.append(" \"Resource\": [\n");
+ for (int i = 0; i < bucketNames.size(); i++) {
+ String name = bucketNames.get(i);
+ sb.append(" \"arn:aws:s3:::").append(name).append("\",\n");
+ sb.append(" \"arn:aws:s3:::").append(name).append("/*\"");
+ if (i < bucketNames.size() - 1) {
+ sb.append(",");
+ }
+ sb.append("\n");
+ }
+ sb.append(" ]\n");
+ sb.append(" }\n");
+ }
+ // Always deny bucket creation/deletion and quota mutation —
+ // CloudStack controls lifecycle and resource accounting. Denying
+ // s3:PutBucketQuota prevents a tenant from using the credentials
+ // returned in BucketResponse to call the SeaweedFS quota extension
+ // directly and bypass CloudStack's resource accounting.
+ sb.append(" ,{\n");
+ sb.append(" \"Sid\": \"DenyBucketLifecycleAndQuota\",\n");
+ sb.append(" \"Effect\": \"Deny\",\n");
+ sb.append(" \"Action\": [\"s3:CreateBucket\", \"s3:DeleteBucket\", \"s3:PutBucketQuota\"],\n");
+ sb.append(" \"Resource\": \"*\"\n");
+ sb.append(" }\n");
+ sb.append(" ]\n");
+ sb.append("}\n");
+ return sb.toString();
+ }
+
+ /**
+ * Returns an S3 connection for the given endpoint and credentials.
+ * Uses path-style access, which SeaweedFS requires.
+ *
+ * @param url the url of the S3 service
+ * @param accessKey the credentials to use for the S3 connection.
+ * @param secretKey the matching secret key.
+ * @return an S3 connection (never null)
+ * @throws CloudRuntimeException on failure.
+ */
+ public static AmazonS3 getS3Client(String url, String accessKey, String secretKey) {
+ AmazonS3 client = AmazonS3ClientBuilder.standard()
+ .enablePathStyleAccess()
+ .withCredentials(new AWSStaticCredentialsProvider(new BasicAWSCredentials(accessKey, secretKey)))
+ .withEndpointConfiguration(new AwsClientBuilder.EndpointConfiguration(url, "us-east-1"))
+ .build();
+ if (client == null) {
+ throw new CloudRuntimeException("Error while creating SeaweedFS S3 client");
+ }
+ return client;
+ }
+
+ /**
+ * Returns an IAM connection for the given endpoint and credentials.
+ *
+ * @param url the url of the IAM service
+ * @param accessKey the credentials to use for the iam connection.
+ * @param secretKey the matching secret key.
+ * @return an IAM connection (never null)
+ * @throws CloudRuntimeException on failure.
+ */
+ public static AmazonIdentityManagement getIAMClient(String url, String accessKey, String secretKey) {
+ AmazonIdentityManagement iamClient = AmazonIdentityManagementClientBuilder.standard()
+ .withCredentials(new AWSStaticCredentialsProvider(new BasicAWSCredentials(accessKey, secretKey)))
+ .withEndpointConfiguration(new AwsClientBuilder.EndpointConfiguration(url, "us-east-1"))
+ .build();
+ if (iamClient == null) {
+ throw new CloudRuntimeException("Error while creating SeaweedFS IAM client");
+ }
+ return iamClient;
+ }
+
+ /**
+ * Test the S3Url to confirm it behaves like an S3 Service.
+ *
+ * Uses bad credentials and looks for the particular error from S3 that says
+ * InvalidAccessKeyId was used. Quietly returns if we connect and get the
+ * expected error back.
+ *
+ * @param s3Url the url to check
+ * @throws CloudRuntimeException if there is any unexpected issue.
+ */
+ public static void validateS3Url(String s3Url) {
+ try {
+ AmazonS3 s3Client = SeaweedFSObjectStoreUtil.getS3Client(s3Url, "unknown", "unknown");
+ s3Client.listBuckets();
+ } catch (AmazonServiceException e) {
+ if (StringUtils.compareIgnoreCase(e.getErrorCode(), "InvalidAccessKeyId") != 0
+ && StringUtils.compareIgnoreCase(e.getErrorCode(), "SignatureDoesNotMatch") != 0) {
+ throw new CloudRuntimeException("Unexpected response from S3 Endpoint.", e);
+ }
+ }
+ }
+
+ /**
+ * Test the IAMUrl to confirm it behaves like an IAM Service.
+ *
+ * Uses bad credentials and looks for the particular error from IAM that says
+ * InvalidAccessKeyId or InvalidClientTokenId was used. Quietly returns if we
+ * connect and get the expected error back.
+ *
+ * @param iamUrl the url to check
+ * @throws CloudRuntimeException if there is any unexpected issue.
+ */
+ public static void validateIAMUrl(String iamUrl) {
+ try {
+ AmazonIdentityManagement iamClient = SeaweedFSObjectStoreUtil.getIAMClient(iamUrl, "unknown", "unknown");
+ iamClient.listAccessKeys();
+ } catch (AmazonServiceException e) {
+ if (! StringUtils.equalsAnyIgnoreCase(e.getErrorCode(), "InvalidAccessKeyId", "InvalidClientTokenId", "SignatureDoesNotMatch")) {
+ throw new CloudRuntimeException("Unexpected response from IAM Endpoint.", e);
+ }
+ }
+ }
+
+ /**
+ * Verify the configured admin credentials actually authenticate against
+ * both the S3 and IAM endpoints.
+ *
+ * {@link #validateS3Url} and {@link #validateIAMUrl} deliberately use bad
+ * credentials to probe that the endpoint behaves like the respective
+ * service, so on their own they accept a store whose admin credentials are
+ * wrong; that only surfaces later on the first bucket or IAM operation.
+ * This performs an authenticated call with the supplied credentials so the
+ * failure happens at registration time.
+ *
+ * @param s3Url the S3 endpoint URL
+ * @param iamUrl the IAM endpoint URL
+ * @param accessKey the admin access key
+ * @param secretKey the admin secret key
+ * @throws CloudRuntimeException if the credentials are rejected
+ */
+ public static void validateCredentials(String s3Url, String iamUrl, String accessKey, String secretKey) {
+ try {
+ getS3Client(s3Url, accessKey, secretKey).listBuckets();
+ } catch (AmazonServiceException e) {
+ if (StringUtils.equalsAnyIgnoreCase(e.getErrorCode(), "InvalidAccessKeyId", "SignatureDoesNotMatch",
+ "AccessDenied", "InvalidClientTokenId")) {
+ throw new CloudRuntimeException("SeaweedFS rejected the supplied admin credentials on the S3 endpoint: "
+ + e.getErrorCode(), e);
+ }
+ throw new CloudRuntimeException("Unexpected response validating admin credentials against the S3 endpoint.", e);
+ }
+ try {
+ getIAMClient(iamUrl, accessKey, secretKey).listUsers();
+ } catch (AmazonServiceException e) {
+ if (StringUtils.equalsAnyIgnoreCase(e.getErrorCode(), "InvalidAccessKeyId", "SignatureDoesNotMatch",
+ "AccessDenied", "InvalidClientTokenId")) {
+ throw new CloudRuntimeException("SeaweedFS rejected the supplied admin credentials on the IAM endpoint: "
+ + e.getErrorCode(), e);
+ }
+ throw new CloudRuntimeException("Unexpected response validating admin credentials against the IAM endpoint.", e);
+ }
+ }
+
+ /**
+ * Set bucket quota via the SeaweedFS S3 extension endpoint.
+ *
+ * SeaweedFS exposes a custom S3 subresource at
+ * PUT /{bucket}?seaweedfs-quota
+ * authenticated via standard S3 SigV4 and authorized via the
+ * s3:PutBucketQuota IAM permission. This avoids the need for a
+ * separate admin API credential.
+ *
+ * The request body is JSON:
+ * {"quota_size": , "quota_unit": "GB", "quota_enabled": true}
+ *
+ * @param s3Url the S3 endpoint URL (e.g. http://host:8333)
+ * @param accessKey the S3 access key (must have s3:PutBucketQuota permission)
+ * @param secretKey the S3 secret key
+ * @param bucketName the bucket name
+ * @param sizeGiB the quota size in GiB (0 to disable quota)
+ * @param allowMissingExtension tolerate a 404/405 for quota 0 when the
+ * optional quota extension is not deployed (initial create only)
+ * @throws CloudRuntimeException on any failure
+ */
+ public static void setBucketQuotaViaS3Extension(String s3Url, String accessKey, String secretKey, String bucketName,
+ long sizeGiB, boolean allowMissingExtension) {
+ setBucketQuotaViaS3Extension(s3Url, accessKey, secretKey, bucketName, sizeGiB, newS3ExtensionHttpClient(),
+ allowMissingExtension);
+ }
+
+ /**
+ * Build a bounded HTTP client for SeaweedFS S3 extension requests with a
+ * connect timeout so a stalled endpoint cannot block the management-server
+ * API thread indefinitely.
+ */
+ public static java.net.http.HttpClient newS3ExtensionHttpClient() {
+ return java.net.http.HttpClient.newBuilder()
+ .connectTimeout(java.time.Duration.ofSeconds(S3_EXTENSION_CONNECT_TIMEOUT_SECONDS))
+ .build();
+ }
+
+ /**
+ * Set bucket quota via the SeaweedFS S3 extension endpoint using the
+ * supplied HTTP client. The client is injected so tests can assert the
+ * signed request without hitting the network.
+ *
+ * @param allowMissingExtension when true and {@code sizeGiB == 0}, a
+ * 404/405 response (indicating the optional SeaweedFS quota
+ * extension is not deployed) is tolerated as a no-op. This is only
+ * safe for the initial bucket create, where the bucket has no quota
+ * to clear. It must be false when clearing an existing positive
+ * quota, because reporting success would lower CloudStack's
+ * accounting while SeaweedFS retains the old quota/read-only state.
+ */
+ public static void setBucketQuotaViaS3Extension(String s3Url, String accessKey, String secretKey,
+ String bucketName, long sizeGiB, java.net.http.HttpClient httpClient,
+ boolean allowMissingExtension) {
+ if (sizeGiB < 0) {
+ // Only zero disables a quota; a negative value would corrupt
+ // resource accounting (BucketApiServiceImpl persists the requested
+ // value and computes deltas from it), so reject it outright.
+ throw new CloudRuntimeException("Bucket quota cannot be negative: " + sizeGiB);
+ }
+ String body;
+ if (sizeGiB == 0) {
+ body = "{\"quota_size\":0,\"quota_unit\":\"B\",\"quota_enabled\":false}";
+ } else {
+ body = String.format("{\"quota_size\":%d,\"quota_unit\":\"GB\",\"quota_enabled\":true}", sizeGiB);
+ }
+ try {
+ executeSignedS3Request("PUT", s3Url, "/" + bucketName + "?seaweedfs-quota", accessKey, secretKey, body, httpClient);
+ } catch (CloudRuntimeException e) {
+ // CreateBucketCmd requires a quota parameter and
+ // BucketApiServiceImpl.createBucket invokes setQuota for every
+ // create, including quota 0. On deployments without the optional
+ // quota extension that call returns 404/405 and would abort the
+ // create. Tolerate it only for the initial create (quota 0 with
+ // allowMissingExtension), where there is no existing quota to
+ // clear, so basic bucket CRUD works without the extension.
+ //
+ // A quota clear on an existing positive quota must NOT be
+ // swallowed: reporting success would lower CloudStack's DB and
+ // resource accounting while SeaweedFS retains the old quota and
+ // read-only state, leaving the two systems inconsistent.
+ //
+ // Distinguish "extension not available" from "bucket not found":
+ // SeaweedFS returns a standard S3 NoSuchBucket error (with
+ // NoSuchBucket in the body) when the bucket does not
+ // exist, which must NOT be swallowed — it indicates CloudStack and
+ // S3 are out of sync.
+ if (allowMissingExtension && sizeGiB == 0 && e.getMessage() != null
+ && (e.getMessage().contains("status 404") || e.getMessage().contains("status 405"))
+ && !e.getMessage().contains("NoSuchBucket")) {
+ org.apache.logging.log4j.LogManager.getLogger(SeaweedFSObjectStoreUtil.class)
+ .warn("SeaweedFS quota extension not available for bucket {}; skipping quota disable (quota is already off by default)", bucketName);
+ return;
+ }
+ throw e;
+ }
+ }
+
+ /**
+ * Execute a custom S3 request with SigV4 signing.
+ *
+ * Uses the AWS SDK v1 Aws4Signer to sign the request, then sends it via
+ * java.net.http.HttpClient. This allows calling SeaweedFS-specific S3
+ * extensions (like ?seaweedfs-quota) that the AWS SDK doesn't natively
+ * support.
+ *
+ * The query string portion of {@code resourcePath} (e.g.
+ * {@code /bucket?seaweedfs-quota}) is split off and added to the request
+ * via {@code addParameter(...)} before signing, so the signer includes it
+ * in the canonical query string. {@code DefaultRequest.setResourcePath}
+ * does not parse an embedded query string, so passing it verbatim would
+ * leave the subresource unsigned while the outgoing URI would still carry
+ * it, causing a signature mismatch on the server.
+ *
+ * @param method HTTP method (PUT, GET, etc.)
+ * @param s3Url the S3 endpoint base URL
+ * @param resourcePath the path + optional query string (e.g. /bucket?seaweedfs-quota)
+ * @param accessKey S3 access key
+ * @param secretKey S3 secret key
+ * @param body the request body (null for GET)
+ * @param httpClient the HTTP client used to send the request
+ * @return the response body as a string
+ * @throws CloudRuntimeException on any failure
+ */
+ protected static String executeSignedS3Request(String method, String s3Url, String resourcePath,
+ String accessKey, String secretKey, String body,
+ java.net.http.HttpClient httpClient) {
+ try {
+ java.net.URI endpointUri = java.net.URI.create(s3Url);
+
+ // Split the resource path into a path and a query string so the
+ // query parameters are signed as canonical query parameters.
+ String path = resourcePath;
+ String queryString = "";
+ int q = resourcePath.indexOf('?');
+ if (q >= 0) {
+ path = resourcePath.substring(0, q);
+ queryString = resourcePath.substring(q + 1);
+ }
+
+ // The AWS SDK v1 AWS4Signer already combines the endpoint path
+ // (request.getEndpoint().getPath()) with the resource path
+ // (request.getResourcePath()) via SdkHttpUtils.appendUri when
+ // building the canonical URI. Set the resource path to just the
+ // bucket/key path (e.g. /bucket) and let the signer prepend the
+ // endpoint path prefix (e.g. /object-s3). The outgoing URI must
+ // also include the endpoint path so the server sees the same path
+ // the signer canonicalized.
+ String endpointPath = endpointUri.getPath();
+ if (endpointPath == null) {
+ endpointPath = "";
+ }
+ if (endpointPath.endsWith("/")) {
+ endpointPath = endpointPath.substring(0, endpointPath.length() - 1);
+ }
+
+ // Build AWS SDK v1 Request for SigV4 signing
+ com.amazonaws.DefaultRequest> request = new com.amazonaws.DefaultRequest<>("s3");
+ request.setEndpoint(endpointUri);
+ request.setHttpMethod(com.amazonaws.http.HttpMethodName.valueOf(method));
+ request.setResourcePath(path);
+ if (! queryString.isEmpty()) {
+ for (String pair : queryString.split("&")) {
+ if (pair.isEmpty()) {
+ continue;
+ }
+ int eq = pair.indexOf('=');
+ if (eq >= 0) {
+ request.addParameter(pair.substring(0, eq), pair.substring(eq + 1));
+ } else {
+ request.addParameter(pair, "");
+ }
+ }
+ }
+ if (body != null) {
+ byte[] bodyBytes = body.getBytes(java.nio.charset.StandardCharsets.UTF_8);
+ request.setContent(new java.io.ByteArrayInputStream(bodyBytes));
+ request.getHeaders().put("Content-Length", String.valueOf(bodyBytes.length));
+ request.getHeaders().put("Content-Type", "application/json");
+ }
+
+ // Sign with SigV4 (AWSS3V4Signer, not the legacy S3Signer which is SigV2)
+ com.amazonaws.auth.AWSCredentials credentials = new com.amazonaws.auth.BasicAWSCredentials(accessKey, secretKey);
+ com.amazonaws.services.s3.internal.AWSS3V4Signer signer = new com.amazonaws.services.s3.internal.AWSS3V4Signer();
+ signer.setServiceName("s3");
+ signer.setRegionName("us-east-1");
+ signer.sign(request, credentials);
+
+ // Build and send the HTTP request with signed headers. The URI
+ // carries the original query string; the signed headers (including
+ // Authorization) are copied from the signed request. Restricted
+ // headers (e.g. Content-Length, Host) are set by the HTTP client /
+ // URI itself and cannot be added via HttpRequest.Builder.header(),
+ // so they are skipped here.
+ // Build the outgoing URI preserving the endpoint path prefix (e.g.
+ // https://host/object-s3) by concatenating it with the resource
+ // path. The signer internally combines the endpoint path with the
+ // resource path to form the same canonical URI, so SigV4 verifies.
+ java.net.URI fullUri = java.net.URI.create(
+ endpointUri.getScheme() + "://" + endpointUri.getRawAuthority()
+ + endpointPath + path);
+ if (! queryString.isEmpty()) {
+ fullUri = java.net.URI.create(fullUri.toString() + "?" + queryString);
+ }
+ java.net.http.HttpRequest.Builder reqBuilder = java.net.http.HttpRequest.newBuilder()
+ .uri(fullUri)
+ .timeout(java.time.Duration.ofSeconds(S3_EXTENSION_REQUEST_TIMEOUT_SECONDS));
+ for (java.util.Map.Entry entry : request.getHeaders().entrySet()) {
+ String headerName = entry.getKey();
+ if (headerName == null || entry.getValue() == null) {
+ continue;
+ }
+ if (isRestrictedHttpHeader(headerName)) {
+ continue;
+ }
+ reqBuilder.header(headerName, entry.getValue());
+ }
+ if (body != null) {
+ reqBuilder.method(method, java.net.http.HttpRequest.BodyPublishers.ofString(body));
+ } else {
+ reqBuilder.method(method, java.net.http.HttpRequest.BodyPublishers.noBody());
+ }
+
+ java.net.http.HttpResponse response = httpClient.send(reqBuilder.build(),
+ java.net.http.HttpResponse.BodyHandlers.ofString());
+
+ int statusCode = response.statusCode();
+ if (statusCode < 200 || statusCode >= 300) {
+ throw new CloudRuntimeException(String.format(
+ "S3 extension request %s %s failed with status %d: %s",
+ method, fullUri, statusCode, response.body()));
+ }
+ return response.body();
+ } catch (CloudRuntimeException e) {
+ throw e;
+ } catch (Exception e) {
+ throw new CloudRuntimeException("S3 extension request failed: " + method + " " + resourcePath, e);
+ }
+ }
+
+ /**
+ * Headers that {@code java.net.http.HttpRequest.Builder.header()} rejects
+ * because they are managed by the HTTP client itself (content length is
+ * derived from the body publisher, host from the URI, etc.). They must be
+ * skipped when copying the signed headers onto the outgoing request.
+ */
+ private static boolean isRestrictedHttpHeader(String headerName) {
+ if (headerName == null) {
+ return true;
+ }
+ switch (headerName.toLowerCase(java.util.Locale.ROOT)) {
+ case "content-length":
+ case "host":
+ case "connection":
+ case "expect":
+ case "upgrade":
+ return true;
+ default:
+ return false;
+ }
+ }
+
+ /**
+ * Prometheus metric name for per-bucket logical size. SeaweedFS publishes
+ * this gauge from the S3 API server's bucket-size metrics loop.
+ */
+ public static final String METRIC_BUCKET_SIZE_BYTES = "SeaweedFS_s3_bucket_size_bytes";
+
+ /**
+ * Scrape the SeaweedFS Prometheus {@code /metrics} endpoint and return a
+ * map of bucket name to logical size in bytes.
+ *
+ * This is a single HTTP GET that returns all bucket sizes in O(buckets)
+ * time, replacing the O(total objects) {@code ListObjectsV2} scan used as a
+ * fallback.
+ *
+ *
{@code metricsUrl} must point at a single SeaweedFS S3 server's
+ * Prometheus exporter (the address configured with {@code -metricsPort}),
+ * NOT at a Prometheus server and NOT at a load-balanced S3 service:
+ *
+ * - A Prometheus server's own {@code /metrics} endpoint exposes its
+ * internal metrics, not the scraped SeaweedFS series.
+ * - SeaweedFS refreshes the bucket-size gauges only on the S3 instance
+ * holding the distributed {@code s3.leader} lock, so a load-balanced
+ * endpoint can route to a non-leader whose gauges are empty.
+ *
+ * Both cases would return HTTP 200 with no usable samples. To detect them,
+ * this method requires a sample line for every managed bucket:
+ * SeaweedFS publishes a zero gauge for empty buckets, so the leader always
+ * exports one sample per bucket it knows about. A missing sample therefore
+ * indicates a wrong endpoint, a non-leader, or a bucket the metrics loop
+ * has not yet observed — all of which must raise a scrape failure so the
+ * caller falls back to the accurate S3 listing rather than reporting zero.
+ *
+ * @param metricsUrl the base URL of the SeaweedFS Prometheus exporter
+ * @param bucketNames the set of bucket names CloudStack manages (used to
+ * filter the scraped metrics; buckets not in this set
+ * are ignored)
+ * @param httpClient the HTTP client used to send the request
+ * @return a map of bucket name to size in bytes, containing exactly the
+ * buckets in {@code bucketNames}
+ * @throws CloudRuntimeException on any HTTP failure, if a sample is missing
+ * for any managed bucket, or if a sample value cannot be parsed.
+ * All failures cause the caller to fall back to the S3 listing
+ * rather than reporting incorrect (zero) usage.
+ */
+ public static java.util.Map parseBucketUsageFromMetrics(String metricsUrl,
+ java.util.Set bucketNames, java.net.http.HttpClient httpClient) {
+ java.util.Map result = new java.util.HashMap<>();
+ try {
+ // Normalize trailing slashes so a configured URL ending in '/'
+ // does not request '//metrics', which can redirect or 404 (the
+ // HTTP client does not follow redirects here) and would silently
+ // force the O(total objects) S3 scan on every usage poll.
+ java.net.URI uri = java.net.URI.create(stripTrailingSlashes(metricsUrl) + "/metrics");
+ java.net.http.HttpRequest request = java.net.http.HttpRequest.newBuilder()
+ .uri(uri)
+ .timeout(java.time.Duration.ofSeconds(S3_EXTENSION_REQUEST_TIMEOUT_SECONDS))
+ .GET()
+ .build();
+ java.net.http.HttpResponse response = httpClient.send(request,
+ java.net.http.HttpResponse.BodyHandlers.ofString());
+ if (response.statusCode() < 200 || response.statusCode() >= 300) {
+ throw new CloudRuntimeException("Prometheus metrics scrape failed with status " + response.statusCode());
+ }
+ // Parse Prometheus text exposition format sample lines like:
+ // SeaweedFS_s3_bucket_size_bytes{bucket="mybucket"} 12345678
+ // Comment lines (# HELP / # TYPE) are skipped: matching them would
+ // accept a response that declares the family but exports no
+ // samples, which happens on a non-leader S3 instance.
+ for (String line : response.body().split("\n")) {
+ if (line.startsWith("#") || !line.startsWith(METRIC_BUCKET_SIZE_BYTES + "{")) {
+ continue;
+ }
+ int bucketLabelStart = line.indexOf("bucket=\"");
+ if (bucketLabelStart < 0) {
+ continue;
+ }
+ int bucketLabelEnd = line.indexOf("\"", bucketLabelStart + 8);
+ if (bucketLabelEnd < 0) {
+ continue;
+ }
+ String bucket = line.substring(bucketLabelStart + 8, bucketLabelEnd);
+ if (!bucketNames.contains(bucket)) {
+ continue;
+ }
+ int valueStart = line.indexOf(' ', bucketLabelEnd + 2);
+ if (valueStart < 0) {
+ continue;
+ }
+ String rawValue = line.substring(valueStart + 1).trim();
+ // Prometheus gauge values are floating point and may use
+ // scientific notation (e.g. 1.2345678e+07). Parse as double
+ // and round, and treat an unparseable value as a scrape
+ // failure so the caller falls back to the S3 listing rather
+ // than reporting this bucket as zero.
+ try {
+ double value = Double.parseDouble(rawValue);
+ if (Double.isNaN(value) || Double.isInfinite(value) || value < 0) {
+ throw new CloudRuntimeException("Invalid " + METRIC_BUCKET_SIZE_BYTES
+ + " value for bucket " + bucket + ": " + rawValue);
+ }
+ result.put(bucket, Math.round(value));
+ } catch (NumberFormatException e) {
+ throw new CloudRuntimeException("Unparseable " + METRIC_BUCKET_SIZE_BYTES
+ + " value for bucket " + bucket + ": " + rawValue, e);
+ }
+ }
+ // Require a sample for every managed bucket. A missing sample means
+ // the endpoint is not a SeaweedFS S3 leader exporter, or the
+ // metrics loop has not yet observed the bucket. Reporting the
+ // remaining buckets as zero would under-report store usage, so
+ // fail and let the caller fall back to the S3 listing.
+ if (!result.keySet().containsAll(bucketNames)) {
+ java.util.Set missing = new java.util.HashSet<>(bucketNames);
+ missing.removeAll(result.keySet());
+ throw new CloudRuntimeException("Prometheus metrics response from " + metricsUrl
+ + " is missing " + METRIC_BUCKET_SIZE_BYTES + " samples for buckets " + missing
+ + "; metricsUrl must point at the SeaweedFS S3 leader's metrics port");
+ }
+ return result;
+ } catch (CloudRuntimeException e) {
+ throw e;
+ } catch (Exception e) {
+ throw new CloudRuntimeException("Failed to scrape Prometheus metrics from " + metricsUrl, e);
+ }
+ }
+}
diff --git a/plugins/storage/object/seaweedfs/src/main/resources/META-INF/cloudstack/storage-object-seaweedfs/module.properties b/plugins/storage/object/seaweedfs/src/main/resources/META-INF/cloudstack/storage-object-seaweedfs/module.properties
new file mode 100644
index 000000000000..94eef7f8cdaf
--- /dev/null
+++ b/plugins/storage/object/seaweedfs/src/main/resources/META-INF/cloudstack/storage-object-seaweedfs/module.properties
@@ -0,0 +1,18 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you 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.
+name=storage-object-seaweedfs
+parent=storage
diff --git a/plugins/storage/object/seaweedfs/src/main/resources/META-INF/cloudstack/storage-object-seaweedfs/spring-storage-object-seaweedfs-context.xml b/plugins/storage/object/seaweedfs/src/main/resources/META-INF/cloudstack/storage-object-seaweedfs/spring-storage-object-seaweedfs-context.xml
new file mode 100644
index 000000000000..66f39dbd8add
--- /dev/null
+++ b/plugins/storage/object/seaweedfs/src/main/resources/META-INF/cloudstack/storage-object-seaweedfs/spring-storage-object-seaweedfs-context.xml
@@ -0,0 +1,31 @@
+
+
+
+
diff --git a/plugins/storage/object/seaweedfs/src/test/java/org/apache/cloudstack/storage/datastore/driver/SeaweedFSObjectStoreDriverImplTest.java b/plugins/storage/object/seaweedfs/src/test/java/org/apache/cloudstack/storage/datastore/driver/SeaweedFSObjectStoreDriverImplTest.java
new file mode 100644
index 000000000000..9988524d8a5e
--- /dev/null
+++ b/plugins/storage/object/seaweedfs/src/test/java/org/apache/cloudstack/storage/datastore/driver/SeaweedFSObjectStoreDriverImplTest.java
@@ -0,0 +1,1184 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you 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.
+// SPDX-License-Identifier: Apache-2.0
+package org.apache.cloudstack.storage.datastore.driver;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertNull;
+import static org.junit.Assert.assertThrows;
+import static org.junit.Assert.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.anyMap;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.doReturn;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.lenient;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.Flow;
+
+import java.io.ByteArrayOutputStream;
+import java.net.http.HttpClient;
+import java.net.http.HttpRequest;
+import java.net.http.HttpResponse;
+import java.nio.ByteBuffer;
+import java.nio.charset.StandardCharsets;
+
+import org.apache.cloudstack.storage.datastore.db.ObjectStoreDao;
+import org.apache.cloudstack.storage.datastore.db.ObjectStoreDetailsDao;
+import org.apache.cloudstack.storage.datastore.db.ObjectStoreVO;
+import org.apache.cloudstack.storage.datastore.util.SeaweedFSObjectStoreUtil;
+import org.apache.cloudstack.storage.object.Bucket;
+
+import com.amazonaws.services.identitymanagement.AmazonIdentityManagement;
+import com.amazonaws.services.identitymanagement.model.AccessKey;
+import com.amazonaws.services.identitymanagement.model.AccessKeyMetadata;
+import com.amazonaws.services.identitymanagement.model.CreateAccessKeyRequest;
+import com.amazonaws.services.identitymanagement.model.CreateAccessKeyResult;
+import com.amazonaws.services.identitymanagement.model.CreateUserRequest;
+import com.amazonaws.services.identitymanagement.model.DeleteAccessKeyRequest;
+import com.amazonaws.services.identitymanagement.model.EntityAlreadyExistsException;
+import com.amazonaws.services.identitymanagement.model.ListAccessKeysRequest;
+import com.amazonaws.services.identitymanagement.model.ListAccessKeysResult;
+import com.amazonaws.services.identitymanagement.model.PutUserPolicyRequest;
+import com.amazonaws.services.s3.AmazonS3;
+import com.amazonaws.services.s3.model.BucketVersioningConfiguration;
+import com.amazonaws.services.s3.model.CreateBucketRequest;
+import com.amazonaws.services.s3.model.ListObjectsV2Request;
+import com.amazonaws.services.s3.model.ListObjectsV2Result;
+import com.amazonaws.services.s3.model.S3ObjectSummary;
+import com.amazonaws.services.s3.model.SetBucketVersioningConfigurationRequest;
+import com.cloud.agent.api.to.BucketTO;
+import com.cloud.storage.BucketVO;
+import com.cloud.storage.dao.BucketDao;
+import com.cloud.user.AccountDetailsDao;
+import com.cloud.user.AccountVO;
+import com.cloud.user.dao.AccountDao;
+import com.cloud.utils.exception.CloudRuntimeException;
+
+import org.junit.After;
+import org.junit.Before;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.mockito.ArgumentCaptor;
+import org.mockito.ArgumentMatchers;
+import org.mockito.Mock;
+import org.mockito.MockitoAnnotations;
+import org.mockito.Spy;
+import org.mockito.junit.MockitoJUnitRunner;
+
+@RunWith(MockitoJUnitRunner.Silent.class)
+public class SeaweedFSObjectStoreDriverImplTest {
+
+ @Spy
+ SeaweedFSObjectStoreDriverImpl driver = new SeaweedFSObjectStoreDriverImpl();
+
+ @Mock
+ AmazonS3 s3Client;
+ @Mock
+ AmazonIdentityManagement iamClient;
+ @Mock
+ ObjectStoreDao objectStoreDao;
+ @Mock
+ ObjectStoreVO objectStoreVO;
+ @Mock
+ ObjectStoreDetailsDao objectStoreDetailsDao;
+ @Mock
+ AccountDao accountDao;
+ @Mock
+ BucketDao bucketDao;
+ @Mock
+ AccountDetailsDao accountDetailsDao;
+ @Mock
+ AccountVO account;
+
+ BucketVO bucketVo;
+ Map storeDetailsMap;
+ Map accountDetailsMap;
+
+ static long TEST_STORE_ID = 1010L;
+ static long TEST_ACCOUNT_ID = 2010L;
+ static long TEST_DOMAIN_ID = 3010L;
+ static String TEST_ACCESS_KEY = "test_access_key";
+ static String TEST_SECRET_KEY = "test_secret_key";
+ static String TEST_BUCKET_NAME = "testbucketname";
+ static String TEST_S3_URL = "http://s3-endpoint";
+ static String TEST_IAM_URL = "http://iam-endpoint";
+ static String TEST_AK = "user_access_key";
+ static String TEST_SK = "user_secret_key";
+ static String TEST_BUCKET_URL = TEST_S3_URL + "/" + TEST_BUCKET_NAME;
+ static String TEST_ACCOUNT_UUID = "account-uuid-1234";
+
+ private AutoCloseable closeable;
+
+ @Before
+ public void setUp() {
+ closeable = MockitoAnnotations.openMocks(this);
+ driver._storeDao = objectStoreDao;
+ driver._storeDetailsDao = objectStoreDetailsDao;
+ driver._accountDao = accountDao;
+ driver._bucketDao = bucketDao;
+ driver._accountDetailsDao = accountDetailsDao;
+
+ lenient().when(objectStoreDao.findById(TEST_STORE_ID)).thenReturn(objectStoreVO);
+ lenient().when(objectStoreVO.getUrl()).thenReturn(TEST_S3_URL);
+
+ storeDetailsMap = new HashMap<>();
+ storeDetailsMap.put(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_ACCESS_KEY, TEST_ACCESS_KEY);
+ storeDetailsMap.put(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_SECRET_KEY, TEST_SECRET_KEY);
+ storeDetailsMap.put(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_S3_URL, TEST_S3_URL);
+ storeDetailsMap.put(SeaweedFSObjectStoreUtil.STORE_DETAILS_KEY_IAM_URL, TEST_IAM_URL);
+ lenient().when(objectStoreDetailsDao.getDetails(TEST_STORE_ID)).thenReturn(storeDetailsMap);
+
+ accountDetailsMap = new HashMap<>();
+ accountDetailsMap.put(SeaweedFSObjectStoreUtil.keyAccessKey(TEST_STORE_ID), TEST_AK);
+ accountDetailsMap.put(SeaweedFSObjectStoreUtil.keySecretKey(TEST_STORE_ID), TEST_SK);
+ lenient().when(accountDetailsDao.findDetails(TEST_ACCOUNT_ID)).thenReturn(accountDetailsMap);
+
+ bucketVo = new BucketVO(TEST_ACCOUNT_ID, TEST_DOMAIN_ID, TEST_STORE_ID, TEST_BUCKET_NAME, null, false, false, false, null);
+
+ // Stub the DB-backed locks with no-op mocks so tests don't
+ // require a real transaction context.
+ com.cloud.utils.db.GlobalLock mockIamLock = mock(com.cloud.utils.db.GlobalLock.class);
+ lenient().doReturn(mockIamLock).when(driver).acquireIamLock(anyLong(), anyLong());
+ lenient().when(mockIamLock.unlock()).thenReturn(true);
+ com.cloud.utils.db.GlobalLock mockNameLock = mock(com.cloud.utils.db.GlobalLock.class);
+ lenient().doReturn(mockNameLock).when(driver).acquireBucketNameLock(anyLong(), anyString());
+ lenient().when(mockNameLock.unlock()).thenReturn(true);
+ }
+
+ @After
+ public void tearDown() throws Exception {
+ closeable.close();
+ }
+
+ @Test
+ public void testGetStoreTO() {
+ assertNull(driver.getStoreTO(null));
+ }
+
+ @Test
+ public void testCreateBucket() throws Exception {
+ doReturn(s3Client).when(driver).getS3ClientByStoreId(TEST_STORE_ID);
+ when(s3Client.doesBucketExistV2(TEST_BUCKET_NAME)).thenReturn(false);
+ when(bucketDao.findById(anyLong())).thenReturn(bucketVo);
+
+ Bucket result = driver.createBucket(bucketVo, false);
+
+ assertEquals(TEST_BUCKET_NAME, result.getName());
+
+ ArgumentCaptor captor = ArgumentCaptor.forClass(BucketVO.class);
+ verify(bucketDao, times(1)).update(any(), captor.capture());
+ BucketVO updated = captor.getValue();
+ assertEquals(TEST_AK, updated.getAccessKey());
+ assertEquals(TEST_SK, updated.getSecretKey());
+ assertEquals(TEST_BUCKET_URL, updated.getBucketURL());
+
+ verify(s3Client, times(1)).createBucket(any(CreateBucketRequest.class));
+ }
+
+ @Test
+ public void testCreateBucketNormalizesTrailingSlashInUrl() throws Exception {
+ // s3Url is operator-supplied and may carry a trailing slash; the stored
+ // bucketURL must not become "...//bucket", which BucketResponse and the
+ // object store browser would both use.
+ doReturn(TEST_S3_URL + "/").when(driver).getS3Url(TEST_STORE_ID);
+ doReturn(s3Client).when(driver).getS3ClientByStoreId(TEST_STORE_ID);
+ when(s3Client.doesBucketExistV2(TEST_BUCKET_NAME)).thenReturn(false);
+ when(bucketDao.findById(anyLong())).thenReturn(bucketVo);
+
+ driver.createBucket(bucketVo, false);
+
+ ArgumentCaptor captor = ArgumentCaptor.forClass(BucketVO.class);
+ verify(bucketDao, times(1)).update(any(), captor.capture());
+ assertEquals(TEST_BUCKET_URL, captor.getValue().getBucketURL());
+ }
+
+ @Test
+ public void testCreateBucketAlreadyExists() throws Exception {
+ doReturn(s3Client).when(driver).getS3ClientByStoreId(TEST_STORE_ID);
+ when(s3Client.doesBucketExistV2(TEST_BUCKET_NAME)).thenReturn(true);
+
+ assertThrows(CloudRuntimeException.class, () -> driver.createBucket(bucketVo, false));
+ verify(s3Client, never()).createBucket(any(CreateBucketRequest.class));
+ }
+
+ @Test
+ public void testListBuckets() throws Exception {
+ doReturn(s3Client).when(driver).getS3ClientByStoreId(TEST_STORE_ID);
+ List s3Buckets = new ArrayList<>();
+ s3Buckets.add(new com.amazonaws.services.s3.model.Bucket("bucket1"));
+ s3Buckets.add(new com.amazonaws.services.s3.model.Bucket("bucket2"));
+ when(s3Client.listBuckets()).thenReturn(s3Buckets);
+
+ List result = driver.listBuckets(TEST_STORE_ID);
+
+ assertEquals(2, result.size());
+ assertEquals("bucket1", result.get(0).getName());
+ assertEquals("bucket2", result.get(1).getName());
+ }
+
+ @Test
+ public void testDeleteBucket() throws Exception {
+ doReturn(s3Client).when(driver).getS3ClientByStoreId(TEST_STORE_ID);
+ BucketTO bucketTO = mock(BucketTO.class);
+ when(bucketTO.getName()).thenReturn(TEST_BUCKET_NAME);
+ when(s3Client.doesBucketExistV2(TEST_BUCKET_NAME)).thenReturn(true);
+
+ assertTrue(driver.deleteBucket(bucketTO, TEST_STORE_ID));
+ verify(s3Client, times(1)).deleteBucket(TEST_BUCKET_NAME);
+ }
+
+ @Test
+ public void testDeleteBucketMarksRowDestroyedInsideLock() throws Exception {
+ doReturn(s3Client).when(driver).getS3ClientByStoreId(TEST_STORE_ID);
+ BucketTO bucketTO = mock(BucketTO.class);
+ when(bucketTO.getName()).thenReturn(TEST_BUCKET_NAME);
+ when(bucketTO.getAccountId()).thenReturn(TEST_ACCOUNT_ID);
+ when(s3Client.doesBucketExistV2(TEST_BUCKET_NAME)).thenReturn(true);
+
+ BucketVO existing = new BucketVO(TEST_ACCOUNT_ID, TEST_DOMAIN_ID, TEST_STORE_ID, TEST_BUCKET_NAME,
+ null, false, false, false, null);
+ List buckets = new ArrayList<>();
+ buckets.add(existing);
+ when(bucketDao.listByObjectStoreIdAndAccountId(TEST_STORE_ID, TEST_ACCOUNT_ID)).thenReturn(buckets);
+
+ assertTrue(driver.deleteBucket(bucketTO, TEST_STORE_ID));
+ // The row must be marked Destroyed inside the IAM lock so a concurrent
+ // policy rebuild cannot re-add the deleted bucket ARN, but must NOT be
+ // removed: BucketApiServiceImpl still needs it to decrement the
+ // resource counts and allocated size.
+ assertEquals(Bucket.State.Destroyed, existing.getState());
+ verify(bucketDao, times(1)).update(existing.getId(), existing);
+ verify(bucketDao, never()).remove(existing.getId());
+ }
+
+ @Test
+ public void testBuildPolicySkipsDestroyedBuckets() throws Exception {
+ BucketVO live = new BucketVO(TEST_ACCOUNT_ID, TEST_DOMAIN_ID, TEST_STORE_ID, "live-bucket",
+ null, false, false, false, null);
+ BucketVO destroyed = new BucketVO(TEST_ACCOUNT_ID, TEST_DOMAIN_ID, TEST_STORE_ID, "destroyed-bucket",
+ null, false, false, false, null);
+ destroyed.setState(Bucket.State.Destroyed);
+ List buckets = new ArrayList<>();
+ buckets.add(live);
+ buckets.add(destroyed);
+ when(bucketDao.listByObjectStoreIdAndAccountId(TEST_STORE_ID, TEST_ACCOUNT_ID)).thenReturn(buckets);
+ when(accountDao.findById(TEST_ACCOUNT_ID)).thenReturn(account);
+
+ driver.updateAccountIAMPolicyLocked(iamClient, TEST_STORE_ID, TEST_ACCOUNT_ID, null);
+
+ ArgumentCaptor captor = ArgumentCaptor.forClass(PutUserPolicyRequest.class);
+ verify(iamClient, times(1)).putUserPolicy(captor.capture());
+ String policy = captor.getValue().getPolicyDocument();
+ // A bucket whose remote counterpart is gone must not be granted: the
+ // name is reusable and another account could claim it.
+ assertTrue(policy.contains("live-bucket"));
+ assertFalse(policy.contains("destroyed-bucket"));
+ }
+
+ @Test
+ public void testDeleteBucketNotFound() throws Exception {
+ doReturn(s3Client).when(driver).getS3ClientByStoreId(TEST_STORE_ID);
+ BucketTO bucketTO = mock(BucketTO.class);
+ when(bucketTO.getName()).thenReturn(TEST_BUCKET_NAME);
+ when(bucketTO.getAccountId()).thenReturn(TEST_ACCOUNT_ID);
+ when(s3Client.doesBucketExistV2(TEST_BUCKET_NAME)).thenReturn(false);
+
+ // Idempotent: if the bucket is already gone (e.g. from a previous
+ // partial failure), skip the S3 delete and proceed to policy refresh.
+ assertTrue(driver.deleteBucket(bucketTO, TEST_STORE_ID));
+ verify(s3Client, never()).deleteBucket(TEST_BUCKET_NAME);
+ }
+
+ @Test
+ public void testSetBucketVersioning() throws Exception {
+ doReturn(s3Client).when(driver).getS3ClientByStoreId(TEST_STORE_ID);
+ BucketTO bucketTO = mock(BucketTO.class);
+ when(bucketTO.getName()).thenReturn(TEST_BUCKET_NAME);
+
+ assertTrue(driver.setBucketVersioning(bucketTO, TEST_STORE_ID));
+ verify(s3Client, times(1)).setBucketVersioningConfiguration(any(SetBucketVersioningConfigurationRequest.class));
+ }
+
+ @Test
+ public void testDeleteBucketVersioning() throws Exception {
+ doReturn(s3Client).when(driver).getS3ClientByStoreId(TEST_STORE_ID);
+ BucketTO bucketTO = mock(BucketTO.class);
+ when(bucketTO.getName()).thenReturn(TEST_BUCKET_NAME);
+
+ assertTrue(driver.deleteBucketVersioning(bucketTO, TEST_STORE_ID));
+ ArgumentCaptor captor =
+ ArgumentCaptor.forClass(SetBucketVersioningConfigurationRequest.class);
+ verify(s3Client, times(1)).setBucketVersioningConfiguration(captor.capture());
+ assertEquals(BucketVersioningConfiguration.SUSPENDED, captor.getValue().getVersioningConfiguration().getStatus());
+ }
+
+ @Test
+ public void testSetBucketQuotaZero() throws Exception {
+ BucketTO bucketTO = mock(BucketTO.class);
+ when(bucketTO.getName()).thenReturn(TEST_BUCKET_NAME);
+ doReturn(TEST_S3_URL).when(driver).getS3Url(TEST_STORE_ID);
+ doReturn("access-key").when(driver).getAccessKey(TEST_STORE_ID);
+ doReturn("secret-key").when(driver).getSecretKey(TEST_STORE_ID);
+
+ HttpClient mockHttpClient = mock(HttpClient.class);
+ HttpResponse mockResponse = mock(HttpResponse.class);
+ when(mockResponse.statusCode()).thenReturn(200);
+ when(mockResponse.body()).thenReturn("");
+ when(mockHttpClient.send(ArgumentMatchers.any(),
+ ArgumentMatchers.>any())).thenReturn(mockResponse);
+ doReturn(mockHttpClient).when(driver).getS3ExtensionHttpClient();
+
+ driver.setBucketQuota(bucketTO, TEST_STORE_ID, 0);
+
+ ArgumentCaptor reqCaptor = ArgumentCaptor.forClass(HttpRequest.class);
+ verify(mockHttpClient, times(1)).send(reqCaptor.capture(),
+ ArgumentMatchers.>any());
+ HttpRequest sent = reqCaptor.getValue();
+ assertEquals("PUT", sent.method());
+ assertEquals("/" + TEST_BUCKET_NAME, sent.uri().getPath());
+ assertTrue("query must carry the seaweedfs-quota subresource",
+ sent.uri().getQuery().contains("seaweedfs-quota"));
+ assertNotNull("request must be SigV4-signed", sent.headers().firstValue("Authorization"));
+ assertEquals("{\"quota_size\":0,\"quota_unit\":\"B\",\"quota_enabled\":false}", extractBody(sent));
+ }
+
+ @Test
+ public void testSetBucketQuotaNegativeRejected() {
+ BucketTO bucketTO = mock(BucketTO.class);
+ when(bucketTO.getName()).thenReturn(TEST_BUCKET_NAME);
+ doReturn(TEST_S3_URL).when(driver).getS3Url(TEST_STORE_ID);
+ doReturn("access-key").when(driver).getAccessKey(TEST_STORE_ID);
+ doReturn("secret-key").when(driver).getSecretKey(TEST_STORE_ID);
+ // Negative quotas must be rejected, not treated as a disable
+ assertThrows(CloudRuntimeException.class, () -> driver.setBucketQuota(bucketTO, TEST_STORE_ID, -1));
+ }
+
+ @Test
+ public void testSetBucketQuotaNonZero() throws Exception {
+ BucketTO bucketTO = mock(BucketTO.class);
+ when(bucketTO.getName()).thenReturn(TEST_BUCKET_NAME);
+ doReturn(TEST_S3_URL).when(driver).getS3Url(TEST_STORE_ID);
+ doReturn("access-key").when(driver).getAccessKey(TEST_STORE_ID);
+ doReturn("secret-key").when(driver).getSecretKey(TEST_STORE_ID);
+
+ HttpClient mockHttpClient = mock(HttpClient.class);
+ HttpResponse mockResponse = mock(HttpResponse.class);
+ when(mockResponse.statusCode()).thenReturn(200);
+ when(mockResponse.body()).thenReturn("");
+ when(mockHttpClient.send(ArgumentMatchers.any(),
+ ArgumentMatchers.>any())).thenReturn(mockResponse);
+ doReturn(mockHttpClient).when(driver).getS3ExtensionHttpClient();
+
+ driver.setBucketQuota(bucketTO, TEST_STORE_ID, 10);
+
+ ArgumentCaptor reqCaptor = ArgumentCaptor.forClass(HttpRequest.class);
+ verify(mockHttpClient, times(1)).send(reqCaptor.capture(),
+ ArgumentMatchers.>any());
+ HttpRequest sent = reqCaptor.getValue();
+ assertEquals("PUT", sent.method());
+ assertEquals("/" + TEST_BUCKET_NAME, sent.uri().getPath());
+ assertTrue(sent.uri().getQuery().contains("seaweedfs-quota"));
+ assertNotNull(sent.headers().firstValue("Authorization"));
+ assertEquals("{\"quota_size\":10,\"quota_unit\":\"GB\",\"quota_enabled\":true}", extractBody(sent));
+ }
+
+ @Test
+ public void testSetBucketQuotaPropagatesFailure() throws Exception {
+ BucketTO bucketTO = mock(BucketTO.class);
+ when(bucketTO.getName()).thenReturn(TEST_BUCKET_NAME);
+ doReturn(TEST_S3_URL).when(driver).getS3Url(TEST_STORE_ID);
+ doReturn("access-key").when(driver).getAccessKey(TEST_STORE_ID);
+ doReturn("secret-key").when(driver).getSecretKey(TEST_STORE_ID);
+
+ HttpClient mockHttpClient = mock(HttpClient.class);
+ HttpResponse mockResponse = mock(HttpResponse.class);
+ when(mockResponse.statusCode()).thenReturn(403);
+ when(mockResponse.body()).thenReturn("forbidden");
+ when(mockHttpClient.send(ArgumentMatchers.any(),
+ ArgumentMatchers.>any())).thenReturn(mockResponse);
+ doReturn(mockHttpClient).when(driver).getS3ExtensionHttpClient();
+
+ assertThrows(CloudRuntimeException.class, () -> driver.setBucketQuota(bucketTO, TEST_STORE_ID, 10));
+ }
+
+ @Test
+ public void testSetBucketQuotaZeroTolerates404() throws Exception {
+ BucketTO bucketTO = mock(BucketTO.class);
+ when(bucketTO.getName()).thenReturn(TEST_BUCKET_NAME);
+ when(bucketTO.getAccountId()).thenReturn(TEST_ACCOUNT_ID);
+ doReturn(TEST_S3_URL).when(driver).getS3Url(TEST_STORE_ID);
+ doReturn("access-key").when(driver).getAccessKey(TEST_STORE_ID);
+ doReturn("secret-key").when(driver).getSecretKey(TEST_STORE_ID);
+
+ // Still Allocated => this is the initial create, the only case allowed
+ // to tolerate a missing quota extension.
+ List allocated = new ArrayList<>();
+ allocated.add(new BucketVO(TEST_ACCOUNT_ID, TEST_DOMAIN_ID, TEST_STORE_ID, TEST_BUCKET_NAME,
+ null, false, false, false, null));
+ when(bucketDao.listByObjectStoreIdAndAccountId(TEST_STORE_ID, TEST_ACCOUNT_ID)).thenReturn(allocated);
+
+ HttpClient mockHttpClient = mock(HttpClient.class);
+ HttpResponse mockResponse = mock(HttpResponse.class);
+ when(mockResponse.statusCode()).thenReturn(404);
+ when(mockResponse.body()).thenReturn("not found");
+ when(mockHttpClient.send(ArgumentMatchers.any(),
+ ArgumentMatchers.>any())).thenReturn(mockResponse);
+ doReturn(mockHttpClient).when(driver).getS3ExtensionHttpClient();
+
+ // Quota 0 (disable) tolerates 404 so bucket creation works on
+ // deployments without the SeaweedFS quota extension.
+ driver.setBucketQuota(bucketTO, TEST_STORE_ID, 0);
+ }
+
+ @Test
+ public void testSetBucketQuotaClearExistingPropagates404() throws Exception {
+ BucketTO bucketTO = mock(BucketTO.class);
+ when(bucketTO.getName()).thenReturn(TEST_BUCKET_NAME);
+ when(bucketTO.getAccountId()).thenReturn(TEST_ACCOUNT_ID);
+ doReturn(TEST_S3_URL).when(driver).getS3Url(TEST_STORE_ID);
+ doReturn("access-key").when(driver).getAccessKey(TEST_STORE_ID);
+ doReturn("secret-key").when(driver).getSecretKey(TEST_STORE_ID);
+
+ // The bucket is already Created, so this quota 0 request is an update,
+ // not the initial create. A 404 must NOT be tolerated: reporting success
+ // would lower CloudStack accounting while SeaweedFS keeps the old quota
+ // and read-only state.
+ BucketVO created = new BucketVO(TEST_ACCOUNT_ID, TEST_DOMAIN_ID, TEST_STORE_ID, TEST_BUCKET_NAME, 100, false, false, false, null);
+ created.setState(Bucket.State.Created);
+ List buckets = new ArrayList<>();
+ buckets.add(created);
+ when(bucketDao.listByObjectStoreIdAndAccountId(TEST_STORE_ID, TEST_ACCOUNT_ID)).thenReturn(buckets);
+
+ HttpClient mockHttpClient = mock(HttpClient.class);
+ HttpResponse mockResponse = mock(HttpResponse.class);
+ when(mockResponse.statusCode()).thenReturn(404);
+ when(mockResponse.body()).thenReturn("not found");
+ when(mockHttpClient.send(ArgumentMatchers.any(),
+ ArgumentMatchers.>any())).thenReturn(mockResponse);
+ doReturn(mockHttpClient).when(driver).getS3ExtensionHttpClient();
+
+ assertThrows(CloudRuntimeException.class, () -> driver.setBucketQuota(bucketTO, TEST_STORE_ID, 0));
+ }
+
+ @Test
+ public void testSetBucketQuotaRejects3xx() throws Exception {
+ BucketTO bucketTO = mock(BucketTO.class);
+ when(bucketTO.getName()).thenReturn(TEST_BUCKET_NAME);
+ doReturn(TEST_S3_URL).when(driver).getS3Url(TEST_STORE_ID);
+ doReturn("access-key").when(driver).getAccessKey(TEST_STORE_ID);
+ doReturn("secret-key").when(driver).getSecretKey(TEST_STORE_ID);
+
+ HttpClient mockHttpClient = mock(HttpClient.class);
+ HttpResponse mockResponse = mock(HttpResponse.class);
+ // 3xx must NOT be treated as success — the mutation was not applied
+ when(mockResponse.statusCode()).thenReturn(302);
+ when(mockResponse.body()).thenReturn("redirect");
+ when(mockHttpClient.send(ArgumentMatchers.any(),
+ ArgumentMatchers.>any())).thenReturn(mockResponse);
+ doReturn(mockHttpClient).when(driver).getS3ExtensionHttpClient();
+
+ assertThrows(CloudRuntimeException.class, () -> driver.setBucketQuota(bucketTO, TEST_STORE_ID, 10));
+ }
+
+ /**
+ * Deterministic SigV4 signature-verification test.
+ *
+ * Signs the same request through the AWS SDK v1 AWSS3V4Signer (the same
+ * signer the production code uses) and asserts that the Authorization
+ * header, signed headers, x-amz-content-sha256, and x-amz-date produced
+ * by the driver's request match. This catches signing regressions (e.g.
+ * the query parameter not being in the canonical query string) that a
+ * mere "header exists" check would miss.
+ */
+ @Test
+ public void testSetBucketQuotaSigV4SignatureVerification() throws Exception {
+ String accessKey = "AKIAIOSFODNN7EXAMPLE";
+ String secretKey = "wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY";
+ String bucketName = "quota-sig-test";
+ String s3Url = "http://s3.example.com:8333";
+ long quotaGiB = 5;
+
+ BucketTO bucketTO = mock(BucketTO.class);
+ when(bucketTO.getName()).thenReturn(bucketName);
+ doReturn(s3Url).when(driver).getS3Url(TEST_STORE_ID);
+ doReturn(accessKey).when(driver).getAccessKey(TEST_STORE_ID);
+ doReturn(secretKey).when(driver).getSecretKey(TEST_STORE_ID);
+
+ HttpClient mockHttpClient = mock(HttpClient.class);
+ HttpResponse mockResponse = mock(HttpResponse.class);
+ when(mockResponse.statusCode()).thenReturn(200);
+ when(mockResponse.body()).thenReturn("");
+ when(mockHttpClient.send(ArgumentMatchers.any(),
+ ArgumentMatchers.>any())).thenReturn(mockResponse);
+ doReturn(mockHttpClient).when(driver).getS3ExtensionHttpClient();
+
+ driver.setBucketQuota(bucketTO, TEST_STORE_ID, quotaGiB);
+
+ ArgumentCaptor reqCaptor = ArgumentCaptor.forClass(HttpRequest.class);
+ verify(mockHttpClient, times(1)).send(reqCaptor.capture(),
+ ArgumentMatchers.>any());
+ HttpRequest sent = reqCaptor.getValue();
+
+ // Build the expected signed request the same way the production code does
+ String expectedBody = String.format("{\"quota_size\":%d,\"quota_unit\":\"GB\",\"quota_enabled\":true}", quotaGiB);
+ byte[] bodyBytes = expectedBody.getBytes(StandardCharsets.UTF_8);
+
+ com.amazonaws.DefaultRequest> expectedRequest = new com.amazonaws.DefaultRequest<>("s3");
+ expectedRequest.setEndpoint(java.net.URI.create(s3Url));
+ expectedRequest.setHttpMethod(com.amazonaws.http.HttpMethodName.PUT);
+ expectedRequest.setResourcePath("/" + bucketName);
+ expectedRequest.addParameter("seaweedfs-quota", "");
+ expectedRequest.setContent(new java.io.ByteArrayInputStream(bodyBytes));
+ expectedRequest.getHeaders().put("Content-Length", String.valueOf(bodyBytes.length));
+ expectedRequest.getHeaders().put("Content-Type", "application/json");
+
+ // Fix the signing timestamp to match the production request so the
+ // test is deterministic and does not intermittently fail when the two
+ // signings straddle a one-second boundary.
+ String productionDate = sent.headers().firstValue("x-amz-date").orElse(null);
+ assertNotNull("production request must carry x-amz-date", productionDate);
+ expectedRequest.getHeaders().put("x-amz-date", productionDate);
+
+ com.amazonaws.auth.AWSCredentials credentials = new com.amazonaws.auth.BasicAWSCredentials(accessKey, secretKey);
+ com.amazonaws.services.s3.internal.AWSS3V4Signer signer = new com.amazonaws.services.s3.internal.AWSS3V4Signer();
+ signer.setServiceName("s3");
+ signer.setRegionName("us-east-1");
+ signer.sign(expectedRequest, credentials);
+
+ // The Authorization header must match exactly — proves the canonical
+ // query string (including seaweedfs-quota), payload hash, and signed
+ // headers all match the independently signed reference request.
+ String expectedAuth = expectedRequest.getHeaders().get("Authorization");
+ String actualAuth = sent.headers().firstValue("Authorization").orElse(null);
+ assertNotNull("Authorization header must be present", actualAuth);
+ assertEquals("SigV4 Authorization header must match the reference signature", expectedAuth, actualAuth);
+
+ // The payload hash must be present and match
+ String expectedContentSha = expectedRequest.getHeaders().get("x-amz-content-sha256");
+ String actualContentSha = sent.headers().firstValue("x-amz-content-sha256").orElse(null);
+ assertEquals("x-amz-content-sha256 must match", expectedContentSha, actualContentSha);
+
+ // The signed headers list must include the query-signing-relevant headers
+ String expectedDate = expectedRequest.getHeaders().get("x-amz-date");
+ String actualDate = sent.headers().firstValue("x-amz-date").orElse(null);
+ assertEquals("x-amz-date must match", expectedDate, actualDate);
+
+ // The query string must carry the subresource
+ assertNotNull("URI must have a query string", sent.uri().getQuery());
+ assertTrue("query must carry seaweedfs-quota", sent.uri().getQuery().contains("seaweedfs-quota"));
+ }
+
+ @Test
+ public void testSetBucketQuotaNoS3ConfigThrows() {
+ BucketTO bucketTO = mock(BucketTO.class);
+ when(bucketTO.getName()).thenReturn(TEST_BUCKET_NAME);
+ // Clear store details so no S3 URL/credentials are configured.
+ // Without this, setUp() stubs valid values and the exception would
+ // come from a real network call rather than the missing-config check.
+ storeDetailsMap.clear();
+ lenient().when(objectStoreDao.findById(TEST_STORE_ID)).thenReturn(null);
+ assertThrows(CloudRuntimeException.class, () -> driver.setBucketQuota(bucketTO, TEST_STORE_ID, 10));
+ }
+
+ @Test
+ public void testBuildAccountIAMPolicyEmptyBuckets() throws Exception {
+ String policy = SeaweedFSObjectStoreUtil.buildAccountIAMPolicy(java.util.Collections.emptyList());
+ // Empty bucket list: deny all S3 access
+ assertTrue(policy.contains("\"Sid\": \"DenyAllS3\""));
+ assertTrue(policy.contains("\"Effect\": \"Deny\""));
+ assertTrue(policy.contains("\"Action\": [\"s3:*\"]"));
+ assertTrue(policy.contains("\"arn:aws:s3:::*\""));
+ // Must still deny bucket lifecycle and quota
+ assertTrue(policy.contains("\"s3:PutBucketQuota\""));
+ assertFalse(policy.contains("\"AllowAccountBuckets\""));
+ }
+
+ @Test
+ public void testBuildAccountIAMPolicyPopulatedBuckets() throws Exception {
+ String policy = SeaweedFSObjectStoreUtil.buildAccountIAMPolicy(
+ java.util.Arrays.asList("bucket-a", "bucket-b"));
+ // Allow access to both bucket and object ARNs
+ assertTrue(policy.contains("\"Sid\": \"AllowAccountBuckets\""));
+ assertTrue(policy.contains("\"arn:aws:s3:::bucket-a\""));
+ assertTrue(policy.contains("\"arn:aws:s3:::bucket-a/*\""));
+ assertTrue(policy.contains("\"arn:aws:s3:::bucket-b\""));
+ assertTrue(policy.contains("\"arn:aws:s3:::bucket-b/*\""));
+ // Must deny bucket lifecycle and quota
+ assertTrue(policy.contains("\"Sid\": \"DenyBucketLifecycleAndQuota\""));
+ assertTrue(policy.contains("\"s3:CreateBucket\""));
+ assertTrue(policy.contains("\"s3:DeleteBucket\""));
+ assertTrue(policy.contains("\"s3:PutBucketQuota\""));
+ assertFalse(policy.contains("\"DenyAllS3\""));
+ }
+
+ /**
+ * Extract the request body from an HttpRequest.BodyPublisher so tests can
+ * assert the JSON payload sent to the SeaweedFS S3 extension.
+ */
+ private static String extractBody(HttpRequest request) throws Exception {
+ return request.bodyPublisher()
+ .map(SeaweedFSObjectStoreDriverImplTest::readBodyPublisher)
+ .orElse(null);
+ }
+
+ private static String readBodyPublisher(HttpRequest.BodyPublisher publisher) {
+ CompletableFuture future = new CompletableFuture<>();
+ publisher.subscribe(new Flow.Subscriber() {
+ final ByteArrayOutputStream baos = new ByteArrayOutputStream();
+ @Override public void onSubscribe(Flow.Subscription s) { s.request(Long.MAX_VALUE); }
+ @Override public void onNext(ByteBuffer b) {
+ byte[] arr = new byte[b.remaining()];
+ b.get(arr);
+ baos.write(arr, 0, arr.length);
+ }
+ @Override public void onError(Throwable t) { future.completeExceptionally(t); }
+ @Override public void onComplete() { future.complete(baos.toString(StandardCharsets.UTF_8)); }
+ });
+ return future.join();
+ }
+
+ @Test
+ public void testCreateUserNew() throws Exception {
+ when(accountDao.findById(TEST_ACCOUNT_ID)).thenReturn(account);
+ when(account.getUuid()).thenReturn(TEST_ACCOUNT_UUID);
+ when(account.getAccountName()).thenReturn("testaccount");
+ doReturn(iamClient).when(driver).getIAMClient(TEST_STORE_ID);
+
+ // No stored credentials yet
+ accountDetailsMap.clear();
+ // No existing access keys to clean up
+ when(iamClient.listAccessKeys(any(ListAccessKeysRequest.class)))
+ .thenReturn(listAccessKeysResult());
+
+ // Access key creation
+ AccessKey accessKey = mock(AccessKey.class);
+ CreateAccessKeyResult accessKeyResult = mock(CreateAccessKeyResult.class);
+ when(accessKey.getAccessKeyId()).thenReturn(TEST_AK);
+ when(accessKey.getSecretAccessKey()).thenReturn(TEST_SK);
+ when(accessKeyResult.getAccessKey()).thenReturn(accessKey);
+ when(iamClient.createAccessKey(any(CreateAccessKeyRequest.class))).thenReturn(accessKeyResult);
+
+ boolean created = driver.createUser(TEST_ACCOUNT_ID, TEST_STORE_ID);
+ assertTrue(created);
+
+ verify(iamClient, times(1)).createUser(any(CreateUserRequest.class));
+ verify(iamClient, times(1)).putUserPolicy(any(PutUserPolicyRequest.class));
+ verify(iamClient, times(1)).createAccessKey(any(CreateAccessKeyRequest.class));
+
+ // Credentials must be written per key. AccountDetailsDao.persist(id, map)
+ // expunges every existing detail for the account first, which would wipe
+ // another object store's namespaced credentials.
+ verify(accountDetailsDao, times(1)).addDetail(TEST_ACCOUNT_ID,
+ SeaweedFSObjectStoreUtil.keyAccessKey(TEST_STORE_ID), TEST_AK, false);
+ verify(accountDetailsDao, times(1)).addDetail(TEST_ACCOUNT_ID,
+ SeaweedFSObjectStoreUtil.keySecretKey(TEST_STORE_ID), TEST_SK, false);
+ verify(accountDetailsDao, never()).persist(anyLong(), anyMap());
+ }
+
+ @Test
+ public void testCreateUserRollsBackCredentialPairWhenSecretPersistFails() throws Exception {
+ when(accountDao.findById(TEST_ACCOUNT_ID)).thenReturn(account);
+ when(account.getUuid()).thenReturn(TEST_ACCOUNT_UUID);
+ when(account.getAccountName()).thenReturn("testaccount");
+ doReturn(iamClient).when(driver).getIAMClient(TEST_STORE_ID);
+
+ when(iamClient.listAccessKeys(any(ListAccessKeysRequest.class)))
+ .thenReturn(listAccessKeysResult());
+
+ AccessKey accessKey = mock(AccessKey.class);
+ CreateAccessKeyResult accessKeyResult = mock(CreateAccessKeyResult.class);
+ when(accessKey.getAccessKeyId()).thenReturn("new-ak");
+ when(accessKey.getSecretAccessKey()).thenReturn("new-sk");
+ when(accessKeyResult.getAccessKey()).thenReturn(accessKey);
+ when(iamClient.createAccessKey(any(CreateAccessKeyRequest.class))).thenReturn(accessKeyResult);
+
+ doThrow(new CloudRuntimeException("secret persist failed")).when(accountDetailsDao).addDetail(TEST_ACCOUNT_ID,
+ SeaweedFSObjectStoreUtil.keySecretKey(TEST_STORE_ID), "new-sk", false);
+
+ assertThrows(CloudRuntimeException.class, () -> driver.createUser(TEST_ACCOUNT_ID, TEST_STORE_ID));
+
+ verify(accountDetailsDao, times(1)).addDetail(TEST_ACCOUNT_ID,
+ SeaweedFSObjectStoreUtil.keyAccessKey(TEST_STORE_ID), "new-ak", false);
+ verify(accountDetailsDao, times(1)).addDetail(TEST_ACCOUNT_ID,
+ SeaweedFSObjectStoreUtil.keySecretKey(TEST_STORE_ID), "new-sk", false);
+ verify(accountDetailsDao, times(1)).addDetail(TEST_ACCOUNT_ID,
+ SeaweedFSObjectStoreUtil.keyAccessKey(TEST_STORE_ID), TEST_AK, false);
+ verify(accountDetailsDao, times(1)).addDetail(TEST_ACCOUNT_ID,
+ SeaweedFSObjectStoreUtil.keySecretKey(TEST_STORE_ID), TEST_SK, false);
+
+ ArgumentCaptor deleteCaptor = ArgumentCaptor.forClass(DeleteAccessKeyRequest.class);
+ verify(iamClient, times(1)).deleteAccessKey(deleteCaptor.capture());
+ assertEquals("new-ak", deleteCaptor.getValue().getAccessKeyId());
+ }
+
+ @Test
+ public void testCreateUserReusesStoredKey() throws Exception {
+ when(accountDao.findById(TEST_ACCOUNT_ID)).thenReturn(account);
+ when(account.getUuid()).thenReturn(TEST_ACCOUNT_UUID);
+ when(account.getAccountName()).thenReturn("testaccount");
+ doReturn(iamClient).when(driver).getIAMClient(TEST_STORE_ID);
+
+ BucketVO staleBucket = new BucketVO(TEST_ACCOUNT_ID, TEST_DOMAIN_ID, TEST_STORE_ID, TEST_BUCKET_NAME, null, false, false, false, null);
+ staleBucket.setAccessKey("stale-ak");
+ staleBucket.setSecretKey("stale-sk");
+ List buckets = new ArrayList<>();
+ buckets.add(staleBucket);
+ when(bucketDao.listByObjectStoreIdAndAccountId(TEST_STORE_ID, TEST_ACCOUNT_ID)).thenReturn(buckets);
+ when(bucketDao.update(staleBucket.getId(), staleBucket)).thenReturn(true);
+
+ // Stored credential still exists in IAM -> must be reused, not rotated
+ when(iamClient.listAccessKeys(any(ListAccessKeysRequest.class)))
+ .thenReturn(listAccessKeysResult(TEST_AK));
+
+ boolean created = driver.createUser(TEST_ACCOUNT_ID, TEST_STORE_ID);
+ assertTrue(created);
+
+ verify(iamClient, times(1)).putUserPolicy(any(PutUserPolicyRequest.class));
+ verify(iamClient, never()).createAccessKey(any(CreateAccessKeyRequest.class));
+ verify(iamClient, never()).deleteAccessKey(any(DeleteAccessKeyRequest.class));
+ verify(accountDetailsDao, never()).persist(anyLong(), ArgumentMatchers.