chrislusf commented on code in PR #14160:
URL: https://github.com/apache/cloudstack/pull/14160#discussion_r4017611100


##########
plugins/storage/object/seaweedfs/src/main/java/org/apache/cloudstack/storage/datastore/driver/SeaweedFSObjectStoreDriverImpl.java:
##########
@@ -0,0 +1,938 @@
+/*
+ * 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 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.
+     *
+     * <p>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;
+    }
+
+    @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-<uuid>} 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.
+     *
+     * <p>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<String, String> 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);
+            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();
+
+        // Update existing bucket records for this account/store with the new
+        // credentials BEFORE persisting the new key in account details. If a
+        // bucket update fails, the stored key remains the old one and a retry
+        // will re-enter the replacement path; if we persisted first, a retry
+        // would see the new stored key and return without repairing the
+        // remaining buckets.
+        updateAccountBucketCredentials(storeId, accountId, key);
+
+        // Persist the credentials in the account details (namespaced by
+        // storeId) with per-key writes. AccountDetailsDao.persist(accountId,
+        // map) expunges every existing detail for the account before inserting
+        // the supplied map, so it would clobber details this snapshot never 
saw
+        // — including another object store's namespaced credentials, since the
+        // IAM lock is keyed by storeId+accountId and two pools can provision
+        // the same account concurrently. addDetail touches only the named key.
+        details.put(accessKeyDetailKey, key.getAccessKeyId());
+        details.put(secretKeyDetailKey, key.getSecretAccessKey());
+        _accountDetailsDao.addDetail(accountId, accessKeyDetailKey, 
key.getAccessKeyId(), false);
+        _accountDetailsDao.addDetail(accountId, secretKeyDetailKey, 
key.getSecretAccessKey(), false);

Review Comment:
   Fixed in `932cc3c`. `createUser` now persists the access/secret pair via 
`persistAccountCredentialsOrRollback`; if either DAO write fails it restores 
the prior details and deletes the newly-created IAM key, so a mixed credential 
pair is not left behind. I also added coverage for the failed-secret-write path 
and for reconciling bucket rows when reusing a stored key.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to