From c9e4bb6fa216a5b50cb1f0e608a7af17b5cefffa Mon Sep 17 00:00:00 2001 From: sylvain senechal Date: Mon, 17 Aug 2026 18:13:10 +0200 Subject: [PATCH 1/5] copy location task reads source location from s3 with STS Issue: BB-812 --- .../queueProcessor/QueueProcessor.js | 11 ++ .../replication/tasks/CopyLocationTask.js | 158 +++++++++++++-- lib/Config.js | 14 ++ lib/management/operatorBackend.js | 1 + .../unit/replication/CopyLocationTask.spec.js | 184 ++++++++++++++++++ 5 files changed, 347 insertions(+), 21 deletions(-) diff --git a/extensions/replication/queueProcessor/QueueProcessor.js b/extensions/replication/queueProcessor/QueueProcessor.js index 24edf7d054..305ec224dd 100644 --- a/extensions/replication/queueProcessor/QueueProcessor.js +++ b/extensions/replication/queueProcessor/QueueProcessor.js @@ -15,6 +15,7 @@ const RoundRobin = require('arsenal').network.RoundRobin; const BackbeatProducer = require('../../../lib/BackbeatProducer'); const BackbeatConsumer = require('../../../lib/BackbeatConsumer'); const VaultClientCache = require('../../../lib/clients/VaultClientCache'); +const CredentialsManager = require('../../../lib/credentials/CredentialsManager'); const QueueEntry = require('../../../lib/models/QueueEntry'); const TaskScheduler = require('../../../lib/tasks/TaskScheduler'); const { getTaskSchedulerQueueKey, @@ -229,6 +230,12 @@ class QueueProcessor extends EventEmitter { this.logger = new Logger( `Backbeat:Replication:QueueProcessor:${this.site}`); + this.assumedRoleCredentialsManager = new CredentialsManager( + 'replication-copy-location', this.logger); + this.assumedRoleS3Clients = {}; + this.assumedRoleHTTPAgent = new HttpAgent.Agent({ keepAlive: true }); + this.assumedRoleHTTPSAgent = new HttpsAgent.Agent({ keepAlive: true }); + // global variables if (sourceConfig.transport === 'https') { this.sourceHTTPAgent = new HttpsAgent.Agent({ @@ -690,6 +697,10 @@ class QueueProcessor extends EventEmitter { destHTTPAgent: this.destHTTPAgent, vaultclientCache: this.vaultclientCache, accountCredsCache: this.accountCredsCache, + assumedRoleCredentialsManager: this.assumedRoleCredentialsManager, + assumedRoleS3Clients: this.assumedRoleS3Clients, + assumedRoleHTTPAgent: this.assumedRoleHTTPAgent, + assumedRoleHTTPSAgent: this.assumedRoleHTTPSAgent, replicationStatusProducer: this.replicationStatusProducer, mProducer: this._mProducer, logger: this.logger, diff --git a/extensions/replication/tasks/CopyLocationTask.js b/extensions/replication/tasks/CopyLocationTask.js index b7d455e909..0fcb68890a 100644 --- a/extensions/replication/tasks/CopyLocationTask.js +++ b/extensions/replication/tasks/CopyLocationTask.js @@ -3,12 +3,14 @@ const { v4: uuid } = require('uuid'); const { errors, jsutil, models } = require('arsenal'); const { ObjectMD } = models; +const { S3Client: AwsS3Client, GetObjectCommand: AwsGetObjectCommand } = + require('@aws-sdk/client-s3'); const BackbeatMetadataProxy = require('../../../lib/BackbeatMetadataProxy'); const BackbeatTask = require('../../../lib/tasks/BackbeatTask'); -const { +const { BackbeatRoutesClient, - GetObjectCommand, + GetObjectCommand: BackbeatRoutesGetObjectCommand, MultipleBackendPutObjectCommand, MultipleBackendInitiateMPUCommand, MultipleBackendPutMPUPartCommand, @@ -24,6 +26,8 @@ const { getAccountCredentials } = require('../../../lib/credentials/AccountCredentials'); const RoleCredentials = require('../../../lib/credentials/RoleCredentials'); +const config = require('../../../lib/Config'); +const { authTypeAssumeRole } = require('../../../lib/constants'); const { metricsExtension, metricsTypeQueued, metricsTypeCompleted } = require('../constants'); @@ -101,6 +105,117 @@ class CopyLocationTask extends BackbeatTask { .setSourceClient(log); } + /** + * Get a cached S3 client, authenticated with the assumed-role + * credentials for the role carried on a location part. + * @param {Object} locationConfig - the location's config + * @param {String} roleArn - the role ARN carried by the location part + * @param {Werelogs} log - the logger instance + * @return {AwsS3Client} the client + * @throws {ArsenalError} AccessDenied (retryable) if credentials + * could not be obtained for the role + */ + _getAssumedRoleS3Client(locationConfig, roleArn, log) { + const { details } = locationConfig; + const s3Endpoint = `${details.transport}://${details.servers[0]}`; + const cacheKey = `${s3Endpoint}::${roleArn}`; + if (this.assumedRoleS3Clients[cacheKey]) { + return this.assumedRoleS3Clients[cacheKey]; + } + const accountId = roleArn.split(':')[4]; + const roleName = roleArn.split(':role/')[1]; + const credentials = this.assumedRoleCredentialsManager.getCredentials({ + id: roleArn, + accountId, + authConfig: { + type: authTypeAssumeRole, + roleName, + }, + stsConfig: { + endpoint: `${details.transport}://${details.sts.host}:${details.sts.port}`, + credentials: { + accessKeyId: details.sts.accessKey, + secretAccessKey: details.sts.secretKey, + }, + }, + }); + if (!credentials) { + log.error('unable to obtain assumed-role credentials for source location', { + method: 'CopyLocationTask._getAssumedRoleS3Client', + roleArn, + endpoint: s3Endpoint, + }); + const err = errors.AccessDenied.customizeDescription( + `unable to assume role ${roleArn} for isCRR source location`); + err.retryable = true; + throw err; + } + const isHttps = details.transport === 'https'; + const client = new AwsS3Client({ + endpoint: s3Endpoint, + credentials: credentials.getCredentialsProvider(), + region: 'us-east-1', + forcePathStyle: true, + requestHandler: { + [isHttps ? 'httpsAgent' : 'httpAgent']: + isHttps ? this.assumedRoleHTTPSAgent : this.assumedRoleHTTPAgent, + requestTimeout: TIMEOUT_MS, + }, + maxAttempts: 1, + }); + client.middlewareStack.add(isRetryableMiddleware(), { + step: 'deserialize', + priority: 'high', + }); + this.assumedRoleS3Clients[cacheKey] = client; + return this.assumedRoleS3Clients[cacheKey]; + } + + /** + * Send a GetObject request for the object's data, + * reading either through Cloudserver's multiple-backend routes, + * or directly from a CRR source location's own S3 + * endpoint via an assumed role. + * @param {ActionQueueEntry} actionEntry - the action entry + * @param {ObjectMD} objMD - metadata object + * @param {Object} [range] - byte range to request, or undefined for the whole object + * @param {Werelogs} log - the logger instance + * @param {AbortController} abortController - abort controller for the GET request + * @return {Promise} resolves to the GetObject response + */ + async _sendGetObject(actionEntry, objMD, range, log, abortController) { + const locationConfig = config.getLocationConstraint(objMD.getDataStoreName()); + if (locationConfig?.isCRR === true) { + const locations = objMD.getLocation(); + const part = locations && locations[0]; + if (!part || !part.role) { + const err = errors.AccessDenied.customizeDescription( + 'missing role on location part for isCRR source location'); + err.retryable = true; + throw err; + } + const s3Client = this._getAssumedRoleS3Client(locationConfig, part.role, log); + const command = new AwsGetObjectCommand({ + Bucket: part.bucket, + Key: objMD.getKey(), + VersionId: part.dataStoreVersionId, + Range: range && `bytes=${range.start}-${range.end}`, + }); + return await s3Client.send(command, { abortSignal: abortController.signal }); + } + + const { bucket, key, version } = actionEntry.getAttribute('target'); + const command = new BackbeatRoutesGetObjectCommand({ + Bucket: bucket, + Key: key, + VersionId: version, + Range: range && `bytes=${range.start}-${range.end}`, + LocationConstraint: objMD.getDataStoreName(), + RequestUids: log.getSerializedUids(), + }); + return await this.backbeatClient.send(command, { abortSignal: abortController.signal }); + } + processQueueEntry(actionEntry, kafkaEntry, done) { const startTime = Date.now(); const log = this.logger.newRequestLogger(); @@ -237,16 +352,8 @@ class CopyLocationTask extends BackbeatTask { let sourceStreamAborted = false; let abortedByPut = false; const abortController = new AbortController(); - const { bucket, key, version } = actionEntry.getAttribute('target'); - const getObjectCommand = new GetObjectCommand({ - Bucket: bucket, - Key: key, - VersionId: version, - LocationConstraint: objMD.getDataStoreName(), - RequestUids: log.getSerializedUids(), - }); - return this.backbeatClient.send(getObjectCommand, { abortSignal: abortController.signal }) + return this._sendGetObject(actionEntry, objMD, undefined, log, abortController) .then(response => { const incomingMsg = response.Body; incomingMsg.on('error', err => { @@ -293,6 +400,14 @@ class CopyLocationTask extends BackbeatTask { actionEntry, objMD, size, incomingMsg, log, putDone); }) .catch(err => { + if (err.name === 'NoSuchVersion') { + log.info('source version no longer exists', Object.assign({ + method: 'CopyLocationTask._getAndPutObjectOnce', + error: err.message, + }, actionEntry.getLogInfo())); + return doneOnce(errors.InvalidObjectState.customizeDescription( + 'source version no longer exists')); + } if (err.$metadata?.httpStatusCode === 404) { log.error('the source object was not found', Object.assign({ method: 'CopyLocationTask._getAndPutObjectOnce', @@ -307,6 +422,7 @@ class CopyLocationTask extends BackbeatTask { method: 'CopyLocationTask._getAndPutObjectOnce', peer: this.sourceConfig.s3, error: err.message, + errorName: err.name, httpStatus: err.$metadata?.httpStatusCode, }, actionEntry.getLogInfo())); return doneOnce(err); @@ -418,20 +534,19 @@ class CopyLocationTask extends BackbeatTask { } const abortController = new AbortController(); - const { bucket, key, version } = actionEntry.getAttribute('target'); - const getObjectCommand = new GetObjectCommand({ - Bucket: bucket, - Key: key, - VersionId: version, - Range: range && `bytes=${range.start}-${range.end}`, - LocationConstraint: objMD.getDataStoreName(), - RequestUids: log.getSerializedUids(), - }); - return this.backbeatClient.send(getObjectCommand, { abortSignal: abortController.signal }) + return this._sendGetObject(actionEntry, objMD, range, log, abortController) .then(response => this._putMPUPart(actionEntry, objMD, response.Body, size, uploadId, partNumber, log, abortController, done)) .catch(err => { + if (err.name === 'NoSuchVersion') { + log.info('source version no longer exists', Object.assign({ + method: 'CopyLocationTask._getRangeAndPutMPUPartOnce', + error: err.message, + }, actionEntry.getLogInfo())); + return done(errors.InvalidObjectState.customizeDescription( + 'source version no longer exists')); + } if (err.$metadata?.httpStatusCode === 404) { return done(err); } @@ -439,6 +554,7 @@ class CopyLocationTask extends BackbeatTask { Object.assign({ method: 'CopyLocationTask._getRangeAndPutMPUPartOnce', error: err.message, + errorName: err.name, httpStatus: err.$metadata?.httpStatusCode, }, actionEntry.getLogInfo())); return done(err); diff --git a/lib/Config.js b/lib/Config.js index 57a9aaa880..58a2786a03 100644 --- a/lib/Config.js +++ b/lib/Config.js @@ -178,6 +178,7 @@ class Config extends EventEmitter { Object.assign(this, parsedConfig); this.transientLocations = {}; + this.locationConstraints = {}; this._setTimeOptions(); this._setLifecycleConductorOptions(); @@ -311,6 +312,19 @@ class Config extends EventEmitter { return this.transientLocations[locationName] || false; } + setLocationConstraints(locationConstraints) { + this.locationConstraints = locationConstraints; + } + + /** + * Get the raw location constraint config for a given location name + * @param {String} locationName - the location constraint name + * @return {Object|undefined} the location config, or undefined if unknown + */ + getLocationConstraint(locationName) { + return this.locationConstraints[locationName]; + } + getPublicInstanceId() { return this.publicInstanceId; } diff --git a/lib/management/operatorBackend.js b/lib/management/operatorBackend.js index f886899081..3273407220 100644 --- a/lib/management/operatorBackend.js +++ b/lib/management/operatorBackend.js @@ -78,6 +78,7 @@ function initManagement(params, done) { })); const locations = require('../../conf/locationConfig.json') || {}; + config.setLocationConstraints(locations); Object.keys(locations).forEach(locName => { config.setIsTransientLocation( locName, locations[locName].isTransient); diff --git a/tests/unit/replication/CopyLocationTask.spec.js b/tests/unit/replication/CopyLocationTask.spec.js index dcf9d84818..2c282406a8 100644 --- a/tests/unit/replication/CopyLocationTask.spec.js +++ b/tests/unit/replication/CopyLocationTask.spec.js @@ -299,4 +299,188 @@ describe('CopyLocationTask', () => { assert.strictEqual(task.retryParams.maxRetries, 13); }); }); + + describe('_sendGetObject', () => { + let task; + let config; + + beforeEach(() => { + config = require('../../../lib/Config'); + task = new CopyLocationTask({ + getStateVars: () => ({ + mProducer: { getProducer: () => {} }, + sourceConfig: { transport: 'http' }, + }), + }); + }); + + afterEach(() => { + sinon.restore(); + }); + + it('should read through Cloudserver when the location is not isCRR', () => { + sinon.stub(config, 'getLocationConstraint').returns({ locationType: 'location-aws-s3-v1', isCRR: false }); + task.backbeatClient = { send: sinon.stub().resolves({ Body: 'stream' }) }; + + const entry = new ActionQueueEntry({ + target: { bucket: 'bucket', key: 'key', version: 'v1' }, + }); + const objMd = new ObjectMD(); + objMd.setDataStoreName('some-location'); + + return task._sendGetObject(entry, objMd, undefined, fakeLogger, new AbortController()) + .then(response => { + assert.deepStrictEqual(response, { Body: 'stream' }); + assert(task.backbeatClient.send.calledOnce); + const command = task.backbeatClient.send.firstCall.args[0]; + assert.strictEqual(command.input.Bucket, 'bucket'); + assert.strictEqual(command.input.Key, 'key'); + assert.strictEqual(command.input.VersionId, 'v1'); + assert.strictEqual(command.input.LocationConstraint, 'some-location'); + }); + }); + + it('should read directly from the CRR source location when isCRR', () => { + sinon.stub(config, 'getLocationConstraint').returns({ + locationType: 'location-scality-crr-v1', + isCRR: true, + details: { + servers: ['production.example.com:443'], + transport: 'https', + sts: { + host: 'sts.production.example.com', + port: '443', + accessKey: 'AK', + secretKey: 'SK', + }, + }, + }); + + const fakeS3Client = { send: sinon.stub().resolves({ Body: 'remote-stream' }) }; + sinon.stub(task, '_getAssumedRoleS3Client').returns(fakeS3Client); + + const entry = new ActionQueueEntry({ + target: { bucket: 'local-bucket', key: 'key', version: 'v1' }, + }); + const objMd = new ObjectMD(); + objMd.setDataStoreName('source-site'); + objMd.setKey('backups/vm001.vbk'); + objMd.setLocation([{ + key: 'backups/vm001.vbk', + size: 1048576, + start: 0, + dataStoreName: 'source-site', + dataStoreType: 'aws_s3', + dataStoreETag: '1:9b2cf535f27731c974343645a3985328', + dataStoreVersionId: 'aJdO95zrzY5BKLXf9GHFItC0d1CkQ0Ei', + bucket: 'backup-repo-01', + role: 'arn:aws:iam::123456789012:role/clean-room-read', + }]); + + return task._sendGetObject(entry, objMd, undefined, fakeLogger, new AbortController()) + .then(response => { + assert.deepStrictEqual(response, { Body: 'remote-stream' }); + assert(task._getAssumedRoleS3Client.calledOnce); + const [locationConfig, roleArn] = task._getAssumedRoleS3Client.firstCall.args; + assert.strictEqual(locationConfig.isCRR, true); + assert.strictEqual(roleArn, 'arn:aws:iam::123456789012:role/clean-room-read'); + assert(fakeS3Client.send.calledOnce); + const command = fakeS3Client.send.firstCall.args[0]; + assert.strictEqual(command.input.Bucket, 'backup-repo-01'); + assert.strictEqual(command.input.Key, 'backups/vm001.vbk'); + assert.strictEqual(command.input.VersionId, 'aJdO95zrzY5BKLXf9GHFItC0d1CkQ0Ei'); + }); + }); + + it('should reject without calling Cloudserver or the remote site when the role is missing', () => { + sinon.stub(config, 'getLocationConstraint').returns({ + locationType: 'location-scality-crr-v1', + isCRR: true, + details: {}, + }); + task.backbeatClient = { send: sinon.stub() }; + sinon.stub(task, '_getAssumedRoleS3Client'); + + const entry = new ActionQueueEntry({ target: {} }); + const objMd = new ObjectMD(); + objMd.setDataStoreName('source-site'); + objMd.setLocation([{ + key: 'k', + bucket: 'b', + dataStoreName: 'source-site', + // no role: owner absent from the ownerId->role map + }]); + + return task._sendGetObject(entry, objMd, undefined, fakeLogger, new AbortController()) + .then(() => assert.fail('expected rejection')) + .catch(err => { + assert(err.AccessDenied); + assert.strictEqual(err.retryable, true); + assert(task.backbeatClient.send.notCalled); + assert(task._getAssumedRoleS3Client.notCalled); + }); + }); + }); + + describe('_getAssumedRoleS3Client', () => { + let task; + const locationConfig = { + details: { + servers: ['production.example.com:443'], + transport: 'https', + sts: { + host: 'sts.production.example.com', + port: '443', + accessKey: 'AK', + secretKey: 'SK', + }, + }, + }; + const roleArn = 'arn:aws:iam::123456789012:role/clean-room-read'; + + beforeEach(() => { + task = new CopyLocationTask({ + getStateVars: () => ({ + mProducer: { getProducer: () => {} }, + sourceConfig: { transport: 'http' }, + assumedRoleCredentialsManager: { + getCredentials: sinon.stub().returns({ + getCredentialsProvider: () => async () => ({}), + }), + }, + assumedRoleS3Clients: {}, + }), + }); + }); + + afterEach(() => { + sinon.restore(); + }); + + it('should cache and reuse the S3 client for the same endpoint and role', () => { + const client1 = task._getAssumedRoleS3Client(locationConfig, roleArn, fakeLogger); + const client2 = task._getAssumedRoleS3Client(locationConfig, roleArn, fakeLogger); + assert.strictEqual(client1, client2); + assert(task.assumedRoleCredentialsManager.getCredentials.calledOnce); + }); + + it('should log and throw a retryable AccessDenied when credentials cannot be obtained', () => { + task.assumedRoleCredentialsManager.getCredentials.returns(null); + const logSpy = sinon.spy(fakeLogger, 'error'); + + assert.throws( + () => task._getAssumedRoleS3Client(locationConfig, roleArn, fakeLogger), + err => err.AccessDenied && err.retryable === true); + assert(logSpy.calledOnce); + }); + + it('should keep the full role name, including any path, when the role ARN has one', () => { + const pathedRoleArn = 'arn:aws:iam::123456789012:role/service-role/clean-room-read'; + + task._getAssumedRoleS3Client(locationConfig, pathedRoleArn, fakeLogger); + + const params = task.assumedRoleCredentialsManager.getCredentials.firstCall.args[0]; + assert.strictEqual(params.authConfig.roleName, 'service-role/clean-room-read'); + }); + }); }); From 43bcdcb42af98626cd746968917f51f47c07e27b Mon Sep 17 00:00:00 2001 From: Francois Ferrand Date: Tue, 25 Aug 2026 21:33:43 +0200 Subject: [PATCH 2/5] copy location task reads the source key from the location part The bucket and version id of the direct read already come from the location part, but the key was still taken from the object metadata. Those are not the same thing: the part describes where the data sits on the source site, which is a different bucket and can be a different key (bucketMatch, key prefixes), so the read has to be fully described by the part. Issue: BB-812 --- extensions/replication/tasks/CopyLocationTask.js | 2 +- tests/unit/replication/CopyLocationTask.spec.js | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/extensions/replication/tasks/CopyLocationTask.js b/extensions/replication/tasks/CopyLocationTask.js index 0fcb68890a..306536d363 100644 --- a/extensions/replication/tasks/CopyLocationTask.js +++ b/extensions/replication/tasks/CopyLocationTask.js @@ -197,7 +197,7 @@ class CopyLocationTask extends BackbeatTask { const s3Client = this._getAssumedRoleS3Client(locationConfig, part.role, log); const command = new AwsGetObjectCommand({ Bucket: part.bucket, - Key: objMD.getKey(), + Key: part.key, VersionId: part.dataStoreVersionId, Range: range && `bytes=${range.start}-${range.end}`, }); diff --git a/tests/unit/replication/CopyLocationTask.spec.js b/tests/unit/replication/CopyLocationTask.spec.js index 2c282406a8..42381658dc 100644 --- a/tests/unit/replication/CopyLocationTask.spec.js +++ b/tests/unit/replication/CopyLocationTask.spec.js @@ -364,7 +364,7 @@ describe('CopyLocationTask', () => { }); const objMd = new ObjectMD(); objMd.setDataStoreName('source-site'); - objMd.setKey('backups/vm001.vbk'); + objMd.setKey('vm001.vbk'); objMd.setLocation([{ key: 'backups/vm001.vbk', size: 1048576, From 16ddfb2f1d3a39a3336c018b7dac840fddee3fb1 Mon Sep 17 00:00:00 2001 From: Francois Ferrand Date: Tue, 8 Sep 2026 19:57:09 +0200 Subject: [PATCH 3/5] copy location task reads the source location data over s3 The source of a clean-room copy is a CRR location, which Cloudserver has no data client for: the data has to be read over S3 on the remote site itself, with a role assumed there for the account owning the object. The location already describes how to get there -- its servers, and the STS to assume roles on -- so it is read directly, the way isCold and isTransient are read elsewhere, rather than through the replication destination config: a clean-room source happens to be a replication destination today, but nothing guarantees it, and falling back to the global destination auth would silently pick the wrong credentials. The location constraints that were carried through Config go away with it. Clients come from the replication ClientManager rather than being built here: it already does this for replication, and on top of a plain client it refreshes credentials, forgets the accounts that went quiet, and supports an STS key read from a file. A location flagged isCRR but with nothing to reach it, or holding more than a single part, now fails rather than falling through to a Cloudserver that cannot serve it or reading from the wrong offsets. Issue: BB-812 --- .../queueProcessor/QueueProcessor.js | 14 +- .../replication/tasks/CopyLocationTask.js | 188 ++++++++------ lib/Config.js | 14 -- lib/management/operatorBackend.js | 1 - .../unit/replication/CopyLocationTask.spec.js | 233 ++++++++++++------ 5 files changed, 270 insertions(+), 180 deletions(-) diff --git a/extensions/replication/queueProcessor/QueueProcessor.js b/extensions/replication/queueProcessor/QueueProcessor.js index 305ec224dd..5b2fddd585 100644 --- a/extensions/replication/queueProcessor/QueueProcessor.js +++ b/extensions/replication/queueProcessor/QueueProcessor.js @@ -15,7 +15,6 @@ const RoundRobin = require('arsenal').network.RoundRobin; const BackbeatProducer = require('../../../lib/BackbeatProducer'); const BackbeatConsumer = require('../../../lib/BackbeatConsumer'); const VaultClientCache = require('../../../lib/clients/VaultClientCache'); -const CredentialsManager = require('../../../lib/credentials/CredentialsManager'); const QueueEntry = require('../../../lib/models/QueueEntry'); const TaskScheduler = require('../../../lib/tasks/TaskScheduler'); const { getTaskSchedulerQueueKey, @@ -230,11 +229,9 @@ class QueueProcessor extends EventEmitter { this.logger = new Logger( `Backbeat:Replication:QueueProcessor:${this.site}`); - this.assumedRoleCredentialsManager = new CredentialsManager( - 'replication-copy-location', this.logger); - this.assumedRoleS3Clients = {}; - this.assumedRoleHTTPAgent = new HttpAgent.Agent({ keepAlive: true }); - this.assumedRoleHTTPSAgent = new HttpsAgent.Agent({ keepAlive: true }); + // clients to read data straight from the sites we replicate to, + // keyed by endpoint and role, shared by all copy location tasks + this.sourceClientManagers = {}; // global variables if (sourceConfig.transport === 'https') { @@ -697,10 +694,7 @@ class QueueProcessor extends EventEmitter { destHTTPAgent: this.destHTTPAgent, vaultclientCache: this.vaultclientCache, accountCredsCache: this.accountCredsCache, - assumedRoleCredentialsManager: this.assumedRoleCredentialsManager, - assumedRoleS3Clients: this.assumedRoleS3Clients, - assumedRoleHTTPAgent: this.assumedRoleHTTPAgent, - assumedRoleHTTPSAgent: this.assumedRoleHTTPSAgent, + sourceClientManagers: this.sourceClientManagers, replicationStatusProducer: this.replicationStatusProducer, mProducer: this._mProducer, logger: this.logger, diff --git a/extensions/replication/tasks/CopyLocationTask.js b/extensions/replication/tasks/CopyLocationTask.js index 306536d363..7ced332945 100644 --- a/extensions/replication/tasks/CopyLocationTask.js +++ b/extensions/replication/tasks/CopyLocationTask.js @@ -3,14 +3,12 @@ const { v4: uuid } = require('uuid'); const { errors, jsutil, models } = require('arsenal'); const { ObjectMD } = models; -const { S3Client: AwsS3Client, GetObjectCommand: AwsGetObjectCommand } = - require('@aws-sdk/client-s3'); const BackbeatMetadataProxy = require('../../../lib/BackbeatMetadataProxy'); const BackbeatTask = require('../../../lib/tasks/BackbeatTask'); const { BackbeatRoutesClient, - GetObjectCommand: BackbeatRoutesGetObjectCommand, + GetObjectCommand, MultipleBackendPutObjectCommand, MultipleBackendInitiateMPUCommand, MultipleBackendPutMPUPartCommand, @@ -26,10 +24,11 @@ const { getAccountCredentials } = require('../../../lib/credentials/AccountCredentials'); const RoleCredentials = require('../../../lib/credentials/RoleCredentials'); -const config = require('../../../lib/Config'); +const ClientManager = require('../../../lib/clients/ClientManager'); const { authTypeAssumeRole } = require('../../../lib/constants'); const { metricsExtension, metricsTypeQueued, metricsTypeCompleted } = require('../constants'); +const locations = require('../../../conf/locationConfig.json') || {}; const MPU_GCP_MAX_PARTS = 1024; @@ -106,76 +105,125 @@ class CopyLocationTask extends BackbeatTask { } /** - * Get a cached S3 client, authenticated with the assumed-role - * credentials for the role carried on a location part. - * @param {Object} locationConfig - the location's config + * Get a cached client on a remote site, authenticated with the + * assumed-role credentials for the role carried on a location part. + * @param {Object} siteConfig - transport, endpoint and STS of the source site * @param {String} roleArn - the role ARN carried by the location part * @param {Werelogs} log - the logger instance - * @return {AwsS3Client} the client + * @return {BackbeatRoutesClient} the client * @throws {ArsenalError} AccessDenied (retryable) if credentials * could not be obtained for the role */ - _getAssumedRoleS3Client(locationConfig, roleArn, log) { - const { details } = locationConfig; - const s3Endpoint = `${details.transport}://${details.servers[0]}`; - const cacheKey = `${s3Endpoint}::${roleArn}`; - if (this.assumedRoleS3Clients[cacheKey]) { - return this.assumedRoleS3Clients[cacheKey]; - } + _getAssumedRoleS3Client(siteConfig, roleArn, log) { + const { transport, endpoint, sts } = siteConfig; + // a location may carry its servers without an explicit port + const [host, port = transport === 'https' ? 443 : 80] = endpoint.split(':'); + const s3Endpoint = `${transport}://${host}:${port}`; const accountId = roleArn.split(':')[4]; const roleName = roleArn.split(':role/')[1]; - const credentials = this.assumedRoleCredentialsManager.getCredentials({ - id: roleArn, - accountId, - authConfig: { - type: authTypeAssumeRole, - roleName, - }, - stsConfig: { - endpoint: `${details.transport}://${details.sts.host}:${details.sts.port}`, - credentials: { - accessKeyId: details.sts.accessKey, - secretAccessKey: details.sts.secretKey, + // one manager per endpoint and role: it holds the credentials and + // clients of every account we assume that role on, and expires them + const cacheKey = `${s3Endpoint}::${roleName}`; + let clientManager = this.sourceClientManagers[cacheKey]; + if (!clientManager) { + clientManager = new ClientManager({ + id: 'replication-copy-location', + authConfig: { + type: authTypeAssumeRole, + roleName, + sts, }, - }, - }); - if (!credentials) { + s3Config: { host, port }, + transport, + }, this.logger); + clientManager.initSTSConfig(); + clientManager.initCredentialsManager(); + this.sourceClientManagers[cacheKey] = clientManager; + } + const client = clientManager.getBackbeatClient(accountId); + if (!client) { log.error('unable to obtain assumed-role credentials for source location', { method: 'CopyLocationTask._getAssumedRoleS3Client', roleArn, endpoint: s3Endpoint, }); const err = errors.AccessDenied.customizeDescription( - `unable to assume role ${roleArn} for isCRR source location`); + `unable to assume role ${roleArn} on source location`); err.retryable = true; throw err; } - const isHttps = details.transport === 'https'; - const client = new AwsS3Client({ - endpoint: s3Endpoint, - credentials: credentials.getCredentialsProvider(), - region: 'us-east-1', - forcePathStyle: true, - requestHandler: { - [isHttps ? 'httpsAgent' : 'httpAgent']: - isHttps ? this.assumedRoleHTTPSAgent : this.assumedRoleHTTPAgent, - requestTimeout: TIMEOUT_MS, - }, - maxAttempts: 1, - }); - client.middlewareStack.add(isRetryableMiddleware(), { - step: 'deserialize', - priority: 'high', - }); - this.assumedRoleS3Clients[cacheKey] = client; - return this.assumedRoleS3Clients[cacheKey]; + return client; + } + + /** + * Get a client to read the object data straight from the site holding + * it, for the locations Cloudserver has no data client for: the remote + * sites we replicate to over CRR, which are exactly the locations a + * clean room reads its source data from. + * @param {ObjectMD} objMD - metadata object + * @param {Werelogs} log - the logger instance + * @return {Object|null} the client and the location part to read, or + * null if the data is reachable through Cloudserver + * @throws {ArsenalError} if the site cannot be read from: no replication + * config to reach it, no single part to read, no role to assume, or no + * credentials for that role + */ + _getSourceLocationClient(objMD, log) { + const site = objMD.getDataStoreName(); + if (!locations[site]?.isCRR) { + return null; + } + // The location itself carries how to reach the site holding the data: + // the servers to read from, and the STS to assume roles on. + const { servers, transport, sts } = locations[site].details || {}; + if (!servers?.length || !sts) { + log.error('no endpoint to reach source location', { + method: 'CopyLocationTask._getSourceLocationClient', + site, + }); + // retryable: a config fix should be picked up without losing the action + const err = errors.InternalError.customizeDescription( + `no endpoint to reach source location ${site}`); + err.retryable = true; + throw err; + } + // A CRR location holds the whole object in a single part, written by + // the source site's own replication: it carries where the data landed + // there, and the role to assume to read it back. Ranged reads below + // assume that, and would read the wrong bytes if it did not hold. + const parts = objMD.getLocation(); + if (parts?.length !== 1) { + log.error('unexpected number of location parts on source location', { + method: 'CopyLocationTask._getSourceLocationClient', + site, + parts: parts?.length ?? 0, + }); + throw errors.InternalError.customizeDescription( + `expected a single location part for source location ${site}`); + } + const part = parts[0]; + if (!part.role) { + log.error('missing role on location part for source location', { + method: 'CopyLocationTask._getSourceLocationClient', + site, + }); + const err = errors.AccessDenied.customizeDescription( + `missing role on location part for source location ${site}`); + err.retryable = true; + throw err; + } + const client = this._getAssumedRoleS3Client( + { transport, endpoint: servers[0], sts }, part.role, log); + return { client, part }; } /** - * Send a GetObject request for the object's data, - * reading either through Cloudserver's multiple-backend routes, - * or directly from a CRR source location's own S3 - * endpoint via an assumed role. + * Send a GetObject request for the object's data, either through the + * local Cloudserver, or straight to the remote site holding the data + * via an assumed role. A CRR location is Scality on both ends, so the + * same command is used either way: only the target differs, and the + * location constraint, which is meaningless on the site that owns the + * data. * @param {ActionQueueEntry} actionEntry - the action entry * @param {ObjectMD} objMD - metadata object * @param {Object} [range] - byte range to request, or undefined for the whole object @@ -184,28 +232,15 @@ class CopyLocationTask extends BackbeatTask { * @return {Promise} resolves to the GetObject response */ async _sendGetObject(actionEntry, objMD, range, log, abortController) { - const locationConfig = config.getLocationConstraint(objMD.getDataStoreName()); - if (locationConfig?.isCRR === true) { - const locations = objMD.getLocation(); - const part = locations && locations[0]; - if (!part || !part.role) { - const err = errors.AccessDenied.customizeDescription( - 'missing role on location part for isCRR source location'); - err.retryable = true; - throw err; - } - const s3Client = this._getAssumedRoleS3Client(locationConfig, part.role, log); - const command = new AwsGetObjectCommand({ - Bucket: part.bucket, - Key: part.key, - VersionId: part.dataStoreVersionId, - Range: range && `bytes=${range.start}-${range.end}`, - }); - return await s3Client.send(command, { abortSignal: abortController.signal }); - } - + const source = this._getSourceLocationClient(objMD, log); const { bucket, key, version } = actionEntry.getAttribute('target'); - const command = new BackbeatRoutesGetObjectCommand({ + const command = new GetObjectCommand(source ? { + Bucket: source.part.bucket, + Key: source.part.key, + VersionId: source.part.dataStoreVersionId, + Range: range && `bytes=${range.start}-${range.end}`, + RequestUids: log.getSerializedUids(), + } : { Bucket: bucket, Key: key, VersionId: version, @@ -213,7 +248,8 @@ class CopyLocationTask extends BackbeatTask { LocationConstraint: objMD.getDataStoreName(), RequestUids: log.getSerializedUids(), }); - return await this.backbeatClient.send(command, { abortSignal: abortController.signal }); + const client = source ? source.client : this.backbeatClient; + return await client.send(command, { abortSignal: abortController.signal }); } processQueueEntry(actionEntry, kafkaEntry, done) { diff --git a/lib/Config.js b/lib/Config.js index 58a2786a03..57a9aaa880 100644 --- a/lib/Config.js +++ b/lib/Config.js @@ -178,7 +178,6 @@ class Config extends EventEmitter { Object.assign(this, parsedConfig); this.transientLocations = {}; - this.locationConstraints = {}; this._setTimeOptions(); this._setLifecycleConductorOptions(); @@ -312,19 +311,6 @@ class Config extends EventEmitter { return this.transientLocations[locationName] || false; } - setLocationConstraints(locationConstraints) { - this.locationConstraints = locationConstraints; - } - - /** - * Get the raw location constraint config for a given location name - * @param {String} locationName - the location constraint name - * @return {Object|undefined} the location config, or undefined if unknown - */ - getLocationConstraint(locationName) { - return this.locationConstraints[locationName]; - } - getPublicInstanceId() { return this.publicInstanceId; } diff --git a/lib/management/operatorBackend.js b/lib/management/operatorBackend.js index 3273407220..f886899081 100644 --- a/lib/management/operatorBackend.js +++ b/lib/management/operatorBackend.js @@ -78,7 +78,6 @@ function initManagement(params, done) { })); const locations = require('../../conf/locationConfig.json') || {}; - config.setLocationConstraints(locations); Object.keys(locations).forEach(locName => { config.setIsTransientLocation( locName, locations[locName].isTransient); diff --git a/tests/unit/replication/CopyLocationTask.spec.js b/tests/unit/replication/CopyLocationTask.spec.js index 42381658dc..f9bc359d19 100644 --- a/tests/unit/replication/CopyLocationTask.spec.js +++ b/tests/unit/replication/CopyLocationTask.spec.js @@ -2,9 +2,11 @@ const assert = require('assert'); const sinon = require('sinon'); const CopyLocationTask = require('../../../extensions/replication/tasks/CopyLocationTask'); +const ClientManager = require('../../../lib/clients/ClientManager'); const ActionQueueEntry = require('../../../lib/models/ActionQueueEntry'); const { errors } = require('arsenal'); const { ObjectMD } = require('arsenal').models; +const locationConfig = require('../../../conf/locationConfig.json'); const fakeLogger = require('../../utils/fakeLogger'); @@ -302,10 +304,44 @@ describe('CopyLocationTask', () => { describe('_sendGetObject', () => { let task; - let config; + + // what the operator writes for a location we replicate to over CRR: + // the servers to reach it, and an STS to assume roles on it + const remoteSiteDetails = { + transport: 'https', + servers: ['production.example.com:443'], + sts: { + host: 'sts.production.example.com', + port: '443', + accessKey: 'AK', + secretKey: 'SK', + }, + }; + + const sourcePart = { + key: 'backups/vm001.vbk', + size: 1048576, + start: 0, + dataStoreName: 'source-site', + dataStoreType: 'aws_s3', + dataStoreETag: '1:9b2cf535f27731c974343645a3985328', + dataStoreVersionId: 'aJdO95zrzY5BKLXf9GHFItC0d1CkQ0Ei', + bucket: 'backup-repo-01', + role: 'arn:aws:iam::123456789012:role/clean-room-read', + }; + + function sourceObjectMD(location = [sourcePart]) { + const objMd = new ObjectMD(); + objMd.setDataStoreName('source-site'); + objMd.setKey('vm001.vbk'); + objMd.setLocation(location); + return objMd; + } beforeEach(() => { - config = require('../../../lib/Config'); + locationConfig['source-site'] = + { type: 'aws_s3', isCRR: true, details: remoteSiteDetails }; + sinon.stub(fakeLogger, 'getSerializedUids').returns('req-uid-1'); task = new CopyLocationTask({ getStateVars: () => ({ mProducer: { getProducer: () => {} }, @@ -315,11 +351,11 @@ describe('CopyLocationTask', () => { }); afterEach(() => { + delete locationConfig['source-site']; sinon.restore(); }); - it('should read through Cloudserver when the location is not isCRR', () => { - sinon.stub(config, 'getLocationConstraint').returns({ locationType: 'location-aws-s3-v1', isCRR: false }); + it('should read through Cloudserver when the location is not a CRR location', () => { task.backbeatClient = { send: sinon.stub().resolves({ Body: 'stream' }) }; const entry = new ActionQueueEntry({ @@ -337,79 +373,102 @@ describe('CopyLocationTask', () => { assert.strictEqual(command.input.Key, 'key'); assert.strictEqual(command.input.VersionId, 'v1'); assert.strictEqual(command.input.LocationConstraint, 'some-location'); + assert.strictEqual(command.input.RequestUids, 'req-uid-1'); }); }); - it('should read directly from the CRR source location when isCRR', () => { - sinon.stub(config, 'getLocationConstraint').returns({ - locationType: 'location-scality-crr-v1', - isCRR: true, - details: { - servers: ['production.example.com:443'], - transport: 'https', - sts: { - host: 'sts.production.example.com', - port: '443', - accessKey: 'AK', - secretKey: 'SK', - }, - }, - }); - - const fakeS3Client = { send: sinon.stub().resolves({ Body: 'remote-stream' }) }; - sinon.stub(task, '_getAssumedRoleS3Client').returns(fakeS3Client); + it('should read directly from the source location when it is a CRR location', () => { + const remoteClient = { send: sinon.stub().resolves({ Body: 'remote-stream' }) }; + sinon.stub(task, '_getAssumedRoleS3Client').returns(remoteClient); const entry = new ActionQueueEntry({ target: { bucket: 'local-bucket', key: 'key', version: 'v1' }, }); - const objMd = new ObjectMD(); - objMd.setDataStoreName('source-site'); - objMd.setKey('vm001.vbk'); - objMd.setLocation([{ - key: 'backups/vm001.vbk', - size: 1048576, - start: 0, - dataStoreName: 'source-site', - dataStoreType: 'aws_s3', - dataStoreETag: '1:9b2cf535f27731c974343645a3985328', - dataStoreVersionId: 'aJdO95zrzY5BKLXf9GHFItC0d1CkQ0Ei', - bucket: 'backup-repo-01', - role: 'arn:aws:iam::123456789012:role/clean-room-read', - }]); - return task._sendGetObject(entry, objMd, undefined, fakeLogger, new AbortController()) + return task._sendGetObject(entry, sourceObjectMD(), undefined, fakeLogger, + new AbortController()) .then(response => { assert.deepStrictEqual(response, { Body: 'remote-stream' }); assert(task._getAssumedRoleS3Client.calledOnce); - const [locationConfig, roleArn] = task._getAssumedRoleS3Client.firstCall.args; - assert.strictEqual(locationConfig.isCRR, true); + const [siteConfig, roleArn] = task._getAssumedRoleS3Client.firstCall.args; + assert.deepStrictEqual(siteConfig, { + transport: 'https', + endpoint: 'production.example.com:443', + sts: remoteSiteDetails.sts, + }); assert.strictEqual(roleArn, 'arn:aws:iam::123456789012:role/clean-room-read'); - assert(fakeS3Client.send.calledOnce); - const command = fakeS3Client.send.firstCall.args[0]; + assert(remoteClient.send.calledOnce); + const command = remoteClient.send.firstCall.args[0]; + // bucket, key and version all describe the data on the + // source site, not the object we are copying assert.strictEqual(command.input.Bucket, 'backup-repo-01'); assert.strictEqual(command.input.Key, 'backups/vm001.vbk'); assert.strictEqual(command.input.VersionId, 'aJdO95zrzY5BKLXf9GHFItC0d1CkQ0Ei'); + // the source site is Scality too: same command, and the + // request uids let us follow the read across both sites + assert.strictEqual(command.input.RequestUids, 'req-uid-1'); + // the data is native there, it must not be redirected + assert.strictEqual(command.input.LocationConstraint, undefined); + }); + }); + + it('should pass the range on to the source location', () => { + const remoteClient = { send: sinon.stub().resolves({}) }; + sinon.stub(task, '_getAssumedRoleS3Client').returns(remoteClient); + + const entry = new ActionQueueEntry({ target: {} }); + + return task._sendGetObject(entry, sourceObjectMD(), { start: 0, end: 99 }, + fakeLogger, new AbortController()) + .then(() => { + const command = remoteClient.send.firstCall.args[0]; + assert.strictEqual(command.input.Range, 'bytes=0-99'); + }); + }); + + it('should reject when the CRR location has no endpoint to reach it', () => { + // flagged isCRR, but carrying none of the details telling us + // where to read the data from + locationConfig['source-site'] = { type: 'aws_s3', isCRR: true, details: {} }; + task.backbeatClient = { send: sinon.stub() }; + + const entry = new ActionQueueEntry({ target: {} }); + + return task._sendGetObject(entry, sourceObjectMD(), undefined, fakeLogger, + new AbortController()) + .then(() => assert.fail('expected rejection')) + .catch(err => { + assert(err.InternalError); + assert.strictEqual(err.retryable, true); + assert(task.backbeatClient.send.notCalled); + }); + }); + + it('should reject when the source location holds more than one part', () => { + task.backbeatClient = { send: sinon.stub() }; + sinon.stub(task, '_getAssumedRoleS3Client'); + + const entry = new ActionQueueEntry({ target: {} }); + const objMd = sourceObjectMD([sourcePart, { ...sourcePart, start: 1048576 }]); + + return task._sendGetObject(entry, objMd, undefined, fakeLogger, new AbortController()) + .then(() => assert.fail('expected rejection')) + .catch(err => { + assert(err.InternalError); + // the metadata will not change: retrying cannot help + assert.notStrictEqual(err.retryable, true); + assert(task.backbeatClient.send.notCalled); + assert(task._getAssumedRoleS3Client.notCalled); }); }); it('should reject without calling Cloudserver or the remote site when the role is missing', () => { - sinon.stub(config, 'getLocationConstraint').returns({ - locationType: 'location-scality-crr-v1', - isCRR: true, - details: {}, - }); task.backbeatClient = { send: sinon.stub() }; sinon.stub(task, '_getAssumedRoleS3Client'); const entry = new ActionQueueEntry({ target: {} }); - const objMd = new ObjectMD(); - objMd.setDataStoreName('source-site'); - objMd.setLocation([{ - key: 'k', - bucket: 'b', - dataStoreName: 'source-site', - // no role: owner absent from the ownerId->role map - }]); + // no role: owner absent from the ownerId->role map + const objMd = sourceObjectMD([{ key: 'k', bucket: 'b', dataStoreName: 'source-site' }]); return task._sendGetObject(entry, objMd, undefined, fakeLogger, new AbortController()) .then(() => assert.fail('expected rejection')) @@ -424,31 +483,31 @@ describe('CopyLocationTask', () => { describe('_getAssumedRoleS3Client', () => { let task; - const locationConfig = { - details: { - servers: ['production.example.com:443'], - transport: 'https', - sts: { - host: 'sts.production.example.com', - port: '443', - accessKey: 'AK', - secretKey: 'SK', - }, + const siteConfig = { + transport: 'https', + endpoint: 'production.example.com:443', + sts: { + host: 'sts.production.example.com', + port: '443', + accessKey: 'AK', + secretKey: 'SK', }, }; const roleArn = 'arn:aws:iam::123456789012:role/clean-room-read'; + let getBackbeatClient; + beforeEach(() => { + sinon.stub(ClientManager.prototype, 'initSTSConfig'); + sinon.stub(ClientManager.prototype, 'initCredentialsManager'); + getBackbeatClient = sinon.stub(ClientManager.prototype, 'getBackbeatClient') + .returns({ send: () => {} }); task = new CopyLocationTask({ getStateVars: () => ({ mProducer: { getProducer: () => {} }, sourceConfig: { transport: 'http' }, - assumedRoleCredentialsManager: { - getCredentials: sinon.stub().returns({ - getCredentialsProvider: () => async () => ({}), - }), - }, - assumedRoleS3Clients: {}, + logger: fakeLogger, + sourceClientManagers: {}, }), }); }); @@ -457,19 +516,35 @@ describe('CopyLocationTask', () => { sinon.restore(); }); - it('should cache and reuse the S3 client for the same endpoint and role', () => { - const client1 = task._getAssumedRoleS3Client(locationConfig, roleArn, fakeLogger); - const client2 = task._getAssumedRoleS3Client(locationConfig, roleArn, fakeLogger); + it('should reuse the client manager for the same endpoint and role', () => { + const client1 = task._getAssumedRoleS3Client(siteConfig, roleArn, fakeLogger); + const client2 = task._getAssumedRoleS3Client(siteConfig, roleArn, fakeLogger); + assert.strictEqual(client1, client2); - assert(task.assumedRoleCredentialsManager.getCredentials.calledOnce); + assert.strictEqual(Object.keys(task.sourceClientManagers).length, 1); + const clientManager = Object.values(task.sourceClientManagers)[0]; + assert.strictEqual(clientManager._transport, 'https'); + assert.deepStrictEqual(clientManager._s3Config, + { host: 'production.example.com', port: '443' }); + assert.deepStrictEqual(clientManager._authConfig.sts, siteConfig.sts); + assert(getBackbeatClient.alwaysCalledWith('123456789012')); + }); + + it('should default the port when the location carries none', () => { + task._getAssumedRoleS3Client( + { ...siteConfig, endpoint: 'production.example.com' }, roleArn, fakeLogger); + + const clientManager = Object.values(task.sourceClientManagers)[0]; + assert.deepStrictEqual(clientManager._s3Config, + { host: 'production.example.com', port: 443 }); }); it('should log and throw a retryable AccessDenied when credentials cannot be obtained', () => { - task.assumedRoleCredentialsManager.getCredentials.returns(null); + getBackbeatClient.returns(null); const logSpy = sinon.spy(fakeLogger, 'error'); assert.throws( - () => task._getAssumedRoleS3Client(locationConfig, roleArn, fakeLogger), + () => task._getAssumedRoleS3Client(siteConfig, roleArn, fakeLogger), err => err.AccessDenied && err.retryable === true); assert(logSpy.calledOnce); }); @@ -477,10 +552,10 @@ describe('CopyLocationTask', () => { it('should keep the full role name, including any path, when the role ARN has one', () => { const pathedRoleArn = 'arn:aws:iam::123456789012:role/service-role/clean-room-read'; - task._getAssumedRoleS3Client(locationConfig, pathedRoleArn, fakeLogger); + task._getAssumedRoleS3Client(siteConfig, pathedRoleArn, fakeLogger); - const params = task.assumedRoleCredentialsManager.getCredentials.firstCall.args[0]; - assert.strictEqual(params.authConfig.roleName, 'service-role/clean-room-read'); + const clientManager = Object.values(task.sourceClientManagers)[0]; + assert.strictEqual(clientManager._authConfig.roleName, 'service-role/clean-room-read'); }); }); }); From 219ed1e9cb3d579e34891ab244dcb0b0b4cb7e55 Mon Sep 17 00:00:00 2001 From: Francois Ferrand Date: Tue, 8 Sep 2026 19:57:28 +0200 Subject: [PATCH 4/5] copy location task reports a missing source version as a failure A source version that no longer exists was turned into an InvalidObjectState, which is committed as a skip: the action never reached the results topic, so the lifecycle side never got to reset transitionInProgress nor count the attempt, and the object stayed flagged in transition for good. Failing to copy the data is an error, and what it means for a transition is for the lifecycle side to decide once it is told. The check on the object state keeps skipping the cases it owns, where the metadata was already updated by someone else and nothing is left hanging. This matters more now that the data may come from a remote site: there, a missing version means the production data is gone while the local object is perfectly alive, which is not something to pass over quietly. Issue: BB-812 --- .../replication/tasks/CopyLocationTask.js | 22 +----- .../unit/replication/CopyLocationTask.spec.js | 76 +++++++++++++++++++ 2 files changed, 79 insertions(+), 19 deletions(-) diff --git a/extensions/replication/tasks/CopyLocationTask.js b/extensions/replication/tasks/CopyLocationTask.js index 7ced332945..7439f70727 100644 --- a/extensions/replication/tasks/CopyLocationTask.js +++ b/extensions/replication/tasks/CopyLocationTask.js @@ -164,9 +164,9 @@ class CopyLocationTask extends BackbeatTask { * @param {Werelogs} log - the logger instance * @return {Object|null} the client and the location part to read, or * null if the data is reachable through Cloudserver - * @throws {ArsenalError} if the site cannot be read from: no replication - * config to reach it, no single part to read, no role to assume, or no - * credentials for that role + * @throws {ArsenalError} if the site cannot be read from: no endpoint to + * reach it, no single part to read, no role to assume, or no credentials + * for that role */ _getSourceLocationClient(objMD, log) { const site = objMD.getDataStoreName(); @@ -436,14 +436,6 @@ class CopyLocationTask extends BackbeatTask { actionEntry, objMD, size, incomingMsg, log, putDone); }) .catch(err => { - if (err.name === 'NoSuchVersion') { - log.info('source version no longer exists', Object.assign({ - method: 'CopyLocationTask._getAndPutObjectOnce', - error: err.message, - }, actionEntry.getLogInfo())); - return doneOnce(errors.InvalidObjectState.customizeDescription( - 'source version no longer exists')); - } if (err.$metadata?.httpStatusCode === 404) { log.error('the source object was not found', Object.assign({ method: 'CopyLocationTask._getAndPutObjectOnce', @@ -575,14 +567,6 @@ class CopyLocationTask extends BackbeatTask { .then(response => this._putMPUPart(actionEntry, objMD, response.Body, size, uploadId, partNumber, log, abortController, done)) .catch(err => { - if (err.name === 'NoSuchVersion') { - log.info('source version no longer exists', Object.assign({ - method: 'CopyLocationTask._getRangeAndPutMPUPartOnce', - error: err.message, - }, actionEntry.getLogInfo())); - return done(errors.InvalidObjectState.customizeDescription( - 'source version no longer exists')); - } if (err.$metadata?.httpStatusCode === 404) { return done(err); } diff --git a/tests/unit/replication/CopyLocationTask.spec.js b/tests/unit/replication/CopyLocationTask.spec.js index f9bc359d19..9044889180 100644 --- a/tests/unit/replication/CopyLocationTask.spec.js +++ b/tests/unit/replication/CopyLocationTask.spec.js @@ -481,6 +481,82 @@ describe('CopyLocationTask', () => { }); }); + describe('source version disappearing mid-copy', () => { + let task; + + function noSuchVersion() { + const err = new Error('The version does not exist.'); + err.name = 'NoSuchVersion'; + err.$metadata = { httpStatusCode: 404 }; + return err; + } + + beforeEach(() => { + task = new CopyLocationTask({ + getStateVars: () => ({ + site: 'test-site', + mProducer: { getProducer: () => {} }, + sourceConfig: { transport: 'http', s3: {} }, + }), + }); + }); + + afterEach(() => { + sinon.restore(); + }); + + it('should fail the copy, not skip it', done => { + sinon.stub(task, '_sendGetObject').rejects(noSuchVersion()); + const put = sinon.stub(task, '_sendMultipleBackendPutObject'); + + const entry = new ActionQueueEntry({ target: {} }); + const objMd = new ObjectMD(); + objMd.setContentLength(200); + + task._getAndPutObjectOnce(entry, objMd, fakeLogger, err => { + assert(err); + // an InvalidObjectState would be committed as a skip, leaving + // the object flagged in transition for good + assert(!err.InvalidObjectState); + assert(put.notCalled); + done(); + }); + }); + + it('should fail the copy, not skip it, while streaming an MPU part', done => { + sinon.stub(task, '_sendGetObject').rejects(noSuchVersion()); + const putPart = sinon.stub(task, '_putMPUPart'); + + const entry = new ActionQueueEntry({ target: {} }); + const objMd = new ObjectMD(); + + task._getRangeAndPutMPUPartOnce(entry, objMd, { start: 0, end: 99 }, 1, + 'upload-id', fakeLogger, err => { + assert(err); + assert(!err.InvalidObjectState); + assert(putPart.notCalled); + done(); + }); + }); + + it('should report the failure so the transition can be reset', () => { + const entry = new ActionQueueEntry({ target: {} }); + entry.setResultsTopic('backbeat-lifecycle-transition-tasks'); + task.replicationStatusProducer = { sendToTopic: sinon.stub() }; + task.dataMoverConsumer = { onEntryCommittable: sinon.stub() }; + + const res = task._publishCopyLocationStatus( + noSuchVersion(), entry, null, fakeLogger); + + // committing here would drop the entry before the lifecycle side + // resets transitionInProgress and counts the attempt + assert.strictEqual(res.committable, false); + assert(task.replicationStatusProducer.sendToTopic.calledOnce); + assert.strictEqual(task.replicationStatusProducer.sendToTopic.firstCall.args[0], + 'backbeat-lifecycle-transition-tasks'); + }); + }); + describe('_getAssumedRoleS3Client', () => { let task; const siteConfig = { From 89deeebd6fc478f9e30727d20d74075448575d1e Mon Sep 17 00:00:00 2001 From: Francois Ferrand Date: Tue, 8 Sep 2026 19:57:42 +0200 Subject: [PATCH 5/5] copy location task builds the get object command once Both reads are the same Scality GetObject and differ only in where they point, so the source location hands back the same shape the action entry already carries and the command is built once from whichever applies. The wire names stay in the command, so what the code reads from is described in the terms the rest of the task uses. Issue: BB-812 --- .../replication/tasks/CopyLocationTask.js | 24 +++++++++---------- 1 file changed, 12 insertions(+), 12 deletions(-) diff --git a/extensions/replication/tasks/CopyLocationTask.js b/extensions/replication/tasks/CopyLocationTask.js index 7439f70727..37af875196 100644 --- a/extensions/replication/tasks/CopyLocationTask.js +++ b/extensions/replication/tasks/CopyLocationTask.js @@ -214,7 +214,8 @@ class CopyLocationTask extends BackbeatTask { } const client = this._getAssumedRoleS3Client( { transport, endpoint: servers[0], sts }, part.role, log); - return { client, part }; + const { bucket, key, dataStoreVersionId: version } = part; + return { client, target: { bucket, key, version }, dataStoreName: undefined }; } /** @@ -232,23 +233,22 @@ class CopyLocationTask extends BackbeatTask { * @return {Promise} resolves to the GetObject response */ async _sendGetObject(actionEntry, objMD, range, log, abortController) { - const source = this._getSourceLocationClient(objMD, log); - const { bucket, key, version } = actionEntry.getAttribute('target'); - const command = new GetObjectCommand(source ? { - Bucket: source.part.bucket, - Key: source.part.key, - VersionId: source.part.dataStoreVersionId, - Range: range && `bytes=${range.start}-${range.end}`, - RequestUids: log.getSerializedUids(), - } : { + // read from the site holding the data, where it is native and there + // is no location to redirect to, or locally through Cloudserver + const { client, target: { bucket, key, version }, dataStoreName } = + this._getSourceLocationClient(objMD, log) ?? { + client: this.backbeatClient, + target: actionEntry.getAttribute('target'), + dataStoreName: objMD.getDataStoreName(), + }; + const command = new GetObjectCommand({ Bucket: bucket, Key: key, VersionId: version, + LocationConstraint: dataStoreName, Range: range && `bytes=${range.start}-${range.end}`, - LocationConstraint: objMD.getDataStoreName(), RequestUids: log.getSerializedUids(), }); - const client = source ? source.client : this.backbeatClient; return await client.send(command, { abortSignal: abortController.signal }); }