From f6762f2203b1746b9f924e3643b0b9d5843b14ac Mon Sep 17 00:00:00 2001 From: Thomas Flament Date: Thu, 27 Aug 2026 17:51:57 +0200 Subject: [PATCH 1/8] Whitelist health check addresses only when an API is served The `server` section is optional, but the whitelist of addresses allowed to reach the health checks was extended unconditionally, so a process serving no API exited before starting: a D/R sink runs the mongo-processor alone, and configuring a section it never reads to get past this is no answer. Issue: BB-811 --- lib/Config.js | 8 +++++--- tests/unit/lib/config/Config.spec.js | 6 ++++++ 2 files changed, 11 insertions(+), 3 deletions(-) diff --git a/lib/Config.js b/lib/Config.js index 8b13bf8c1..fa783caf7 100644 --- a/lib/Config.js +++ b/lib/Config.js @@ -118,9 +118,11 @@ class Config extends EventEmitter { // whitelist IP, CIDR for health checks const defaultHealthChecks = ['127.0.0.1/8', '::1']; - const healthChecks = parsedConfig.server.healthChecks; - healthChecks.allowFrom = - healthChecks.allowFrom.concat(defaultHealthChecks); + const healthChecks = parsedConfig.server?.healthChecks; + if (healthChecks) { + healthChecks.allowFrom = + healthChecks.allowFrom.concat(defaultHealthChecks); + } // additional certs checks if (parsedConfig.certFilePaths) { diff --git a/tests/unit/lib/config/Config.spec.js b/tests/unit/lib/config/Config.spec.js index c26ee1d6a..ee333128d 100644 --- a/tests/unit/lib/config/Config.spec.js +++ b/tests/unit/lib/config/Config.spec.js @@ -33,6 +33,12 @@ describe('Config', () => { assert.doesNotThrow(() => config._parseConfig(testConfig)); }); + it('should accept a config serving no API, as a D/R sink does', () => { + delete testConfig.server; + testConfig.extensions = { mongoProcessor: testConfig.extensions.mongoProcessor }; + assert.doesNotThrow(() => config._parseConfig(testConfig)); + }); + it('should accept a config publishing no metrics', () => { delete testConfig.metrics; assert.doesNotThrow(() => config._parseConfig(testConfig)); From 367dc4f89679729f65250d01e3ecb689f1cc7d1a Mon Sep 17 00:00:00 2001 From: Thomas Flament Date: Wed, 23 Sep 2026 15:27:52 +0200 Subject: [PATCH 2/8] Treat a missing object as the expected outcome it is The mongo-processor reads the object it is about to write, and a first write finds nothing there. Ingestion reaches that path only for a restored object or a bucket carrying replication rules, so both the error log and the error test it leans on went unnoticed; a D/R sink reads the stored document for every entry, which makes a miss the steady state. Log it as the outcome it is, and test it the way the delete path a few lines below already does. `err.NoSuchKey` is arsenal's deprecated comparison, set only while `allowUnsafeErrComp` is on and slated for removal with ARSN-176 -- with it off, every first delivery of every object would be an error. Issue: BB-811 --- extensions/mongoProcessor/MongoQueueProcessor.js | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/extensions/mongoProcessor/MongoQueueProcessor.js b/extensions/mongoProcessor/MongoQueueProcessor.js index 5e797a5c9..41220db35 100644 --- a/extensions/mongoProcessor/MongoQueueProcessor.js +++ b/extensions/mongoProcessor/MongoQueueProcessor.js @@ -238,6 +238,14 @@ class MongoQueueProcessor { log.debug('getting zenko object metadata', { bucket, key, versionId, params }); return this._mongoClient.getObject(bucket, key, params, log, (err, data) => { + if (err?.is.NoSuchKey) { + log.debug('no object metadata stored yet', { + method: 'MongoQueueProcessor._getZenkoObjectMetadata', + entry: entry.getLogInfo(), + }); + return done(err); + } + if (err) { log.error('error getting zenko object metadata', { method: 'MongoQueueProcessor._getZenkoObjectMetadata', @@ -499,7 +507,7 @@ class MongoQueueProcessor { }; maybeGetZenkoObjectMetadata((err, zenkoObjMd) => { - if (err && !err.NoSuchKey) { + if (err && !err.is.NoSuchKey) { this._normalizePendingMetric(location); log.end().error('error processing object queue entry', { method: 'MongoQueueProcessor._processObjectQueueEntry', From 1bb5996144397daf3e6831f8bacfab0fe67cc95c Mon Sep 17 00:00:00 2001 From: Thomas Flament Date: Wed, 23 Sep 2026 19:16:28 +0200 Subject: [PATCH 3/8] Derive the circuit breaker expectations from the location config The bucket processor circuit breaker expands a `${location}` template over every location in the config, and its spec restated the expansion by hand, one entry per location: adding a location to the sample config broke it. Build the expectation from the config instead. Issue: BB-811 --- .../lifecycle/CircuitBreakerGroup.spec.js | 194 ++---------------- 1 file changed, 20 insertions(+), 174 deletions(-) diff --git a/tests/unit/lifecycle/CircuitBreakerGroup.spec.js b/tests/unit/lifecycle/CircuitBreakerGroup.spec.js index 6047ea97a..0d3e1265d 100644 --- a/tests/unit/lifecycle/CircuitBreakerGroup.spec.js +++ b/tests/unit/lifecycle/CircuitBreakerGroup.spec.js @@ -10,6 +10,18 @@ const logger = require('../../utils/fakeLogger'); const { BreakerState } = require('@scality/breakbeat').CircuitBreaker; describe('extractBucketProcessorCircuitBreakerConfigs', () => { + // every location in the config gets the templated probe, so the + // expectation follows the config rather than restating it + function formatLocationProbeConfigList(probe) { + return Object.keys(locations).map(location => + formatProbeConfig(probe, '${location}', location)); + } + + function formatLocationProbeConfigs(probe) { + return Object.fromEntries(Object.keys(locations).map(location => + [location, [formatProbeConfig(probe, '${location}', location)]])); + } + function formatProbeConfig(probe, template, value) { const withClause = probe.query.match(/^when\s?\(\{(.*?)\}\)\sand\s/); let query = withClause ? probe.query.replace(withClause[0], '') : probe.query; @@ -384,126 +396,18 @@ describe('extractBucketProcessorCircuitBreakerConfigs', () => { }, circuitBreakers: { transition: { - location: { - 'us-east-1': [ - formatProbeConfig( - topicSpecificLocationTemplateProbe, - '${location}', - 'us-east-1', - ), - ], - 'us-east-2': [ - formatProbeConfig( - topicSpecificLocationTemplateProbe, - '${location}', - 'us-east-2', - ), - ], - 'wontwork-location': [ - formatProbeConfig( - topicSpecificLocationTemplateProbe, - '${location}', - 'wontwork-location', - ), - ], - 'location-dmf-v1': [ - formatProbeConfig( - topicSpecificLocationTemplateProbe, - '${location}', - 'location-dmf-v1', - ), - ], - }, + location: formatLocationProbeConfigs(topicSpecificLocationTemplateProbe), topic: { - 'cold-archive-req-location-dmf-v1': [ - formatProbeConfig( - topicSpecificLocationTemplateProbe, - '${location}', - 'us-east-1', - ), - formatProbeConfig( - topicSpecificLocationTemplateProbe, - '${location}', - 'us-east-2', - ), - formatProbeConfig( - topicSpecificLocationTemplateProbe, - '${location}', - 'wontwork-location', - ), - formatProbeConfig( - topicSpecificLocationTemplateProbe, - '${location}', - 'location-dmf-v1', - ), - formatProbeConfig( - topicSpecificLocationTemplateProbe, - '${location}', - 'location-crr-source', - ), - ], + 'cold-archive-req-location-dmf-v1': + formatLocationProbeConfigList(topicSpecificLocationTemplateProbe), }, global: [], }, expiration: { - location: { - 'us-east-1': [ - formatProbeConfig( - topicSpecificLocationTemplateProbe, - '${location}', - 'us-east-1', - ), - ], - 'us-east-2': [ - formatProbeConfig( - topicSpecificLocationTemplateProbe, - '${location}', - 'us-east-2', - ), - ], - 'wontwork-location': [ - formatProbeConfig( - topicSpecificLocationTemplateProbe, - '${location}', - 'wontwork-location', - ), - ], - 'location-dmf-v1': [ - formatProbeConfig( - topicSpecificLocationTemplateProbe, - '${location}', - 'location-dmf-v1', - ), - ], - }, + location: formatLocationProbeConfigs(topicSpecificLocationTemplateProbe), topic: { - 'cold-archive-req-location-dmf-v1': [ - formatProbeConfig( - topicSpecificLocationTemplateProbe, - '${location}', - 'us-east-1', - ), - formatProbeConfig( - topicSpecificLocationTemplateProbe, - '${location}', - 'us-east-2', - ), - formatProbeConfig( - topicSpecificLocationTemplateProbe, - '${location}', - 'wontwork-location', - ), - formatProbeConfig( - topicSpecificLocationTemplateProbe, - '${location}', - 'location-dmf-v1', - ), - formatProbeConfig( - topicSpecificLocationTemplateProbe, - '${location}', - 'location-crr-source', - ), - ], + 'cold-archive-req-location-dmf-v1': + formatLocationProbeConfigList(topicSpecificLocationTemplateProbe), }, global: [], }, @@ -524,70 +428,12 @@ describe('extractBucketProcessorCircuitBreakerConfigs', () => { }, circuitBreakers: { transition: { - location: { - 'us-east-1': [ - formatProbeConfig( - locationTemplatedProbe, - '${location}', - 'us-east-1', - ), - ], - 'us-east-2': [ - formatProbeConfig( - locationTemplatedProbe, - '${location}', - 'us-east-2', - ), - ], - 'wontwork-location': [ - formatProbeConfig( - locationTemplatedProbe, - '${location}', - 'wontwork-location', - ), - ], - 'location-dmf-v1': [ - formatProbeConfig( - locationTemplatedProbe, - '${location}', - 'location-dmf-v1', - ), - ], - }, + location: formatLocationProbeConfigs(locationTemplatedProbe), topic: {}, global: [], }, expiration: { - location: { - 'us-east-1': [ - formatProbeConfig( - locationTemplatedProbe, - '${location}', - 'us-east-1', - ), - ], - 'us-east-2': [ - formatProbeConfig( - locationTemplatedProbe, - '${location}', - 'us-east-2', - ), - ], - 'wontwork-location': [ - formatProbeConfig( - locationTemplatedProbe, - '${location}', - 'wontwork-location', - ), - ], - 'location-dmf-v1': [ - formatProbeConfig( - locationTemplatedProbe, - '${location}', - 'location-dmf-v1', - ), - ], - }, + location: formatLocationProbeConfigs(locationTemplatedProbe), topic: {}, global: [], }, From 5b2ee28b84ce114e74035bfdf23d39d6ec4d5add Mon Sep 17 00:00:00 2001 From: Thomas Flament Date: Tue, 29 Sep 2026 12:29:14 +0200 Subject: [PATCH 4/8] Introduce a metadata policy, defaulting to ingestion The mongo-processor was written as the out-of-band ingestion consumer, and the D/R metadata sink is to reuse it. The two disagree about most of what it does to an object's metadata, so put those decisions behind a policy: an abstract MetadataPolicy whose methods assert, an implementation per mode, and an index mapping the configured mode to the class, as the notification extension does for its destinations. The processor hands the policy entries and asks it whether an entry needs the metadata already stored, which version on this site the entry acts on, how the entry applies onto what is stored -- including the version the result is written under, or nothing if the entry changes nothing -- and whether a delete still applies. IngestionMetadataPolicy carries today's behaviour verbatim, so this commit changes nothing. The `x-amz-meta-scal-version-id` header of a restored object is ingestion's alone, so the policy reads it, where the processor did. The version an entry is written under stays apart from the one it is read from: an entry naming a version that is not stored is ingested under its own id. The default lives beside the policy map, so a processor built programmatically gets the same mode as one built from a config file. Issue: BB-811 --- .../MongoProcessorConfigValidator.js | 2 + .../mongoProcessor/MongoQueueProcessor.js | 155 ++------------ .../metadataPolicy/IngestionMetadataPolicy.js | 202 ++++++++++++++++++ .../metadataPolicy/MetadataPolicy.js | 84 ++++++++ .../mongoProcessor/metadataPolicy/index.js | 10 + .../ingestion/MongoQueueProcessor.js | 3 +- .../mongoProcessor/MetadataPolicy.spec.js | 22 ++ .../MongoProcessorConfigValidator.spec.js | 20 +- 8 files changed, 358 insertions(+), 140 deletions(-) create mode 100644 extensions/mongoProcessor/metadataPolicy/IngestionMetadataPolicy.js create mode 100644 extensions/mongoProcessor/metadataPolicy/MetadataPolicy.js create mode 100644 extensions/mongoProcessor/metadataPolicy/index.js create mode 100644 tests/unit/mongoProcessor/MetadataPolicy.spec.js diff --git a/extensions/mongoProcessor/MongoProcessorConfigValidator.js b/extensions/mongoProcessor/MongoProcessorConfigValidator.js index b5ae49183..f4f6e11f7 100644 --- a/extensions/mongoProcessor/MongoProcessorConfigValidator.js +++ b/extensions/mongoProcessor/MongoProcessorConfigValidator.js @@ -1,4 +1,5 @@ const joi = require('joi'); +const { metadataPolicies, defaultMode } = require('./metadataPolicy'); const { retryParamsJoi, probeServerJoi, logJoiOptional } = require('../../lib/config/configItems.joi'); const { extensionConfigValidator } = require('../../lib/config/extensionConfigValidator'); @@ -7,6 +8,7 @@ const { MAX_QUEUED_DEFAULT } = require('../../lib/constants').backbeatConsumer; const joiSchema = joi.object({ topic: joi.string().required(), groupId: joi.string().required(), + mode: joi.string().valid(...Object.keys(metadataPolicies)).default(defaultMode), retry: retryParamsJoi, concurrency: joi.number().greater(0).default(1), maxQueued: joi.number().greater(0).default(MAX_QUEUED_DEFAULT), diff --git a/extensions/mongoProcessor/MongoQueueProcessor.js b/extensions/mongoProcessor/MongoQueueProcessor.js index 41220db35..90304f1d1 100644 --- a/extensions/mongoProcessor/MongoQueueProcessor.js +++ b/extensions/mongoProcessor/MongoQueueProcessor.js @@ -4,12 +4,10 @@ const async = require('async'); const Logger = require('werelogs').Logger; const errors = require('@scality/arsenal').errors; -const { replicationBackends, emptyFileMd5 } = require('@scality/arsenal').constants; +const { replicationBackends } = require('@scality/arsenal').constants; const MongoClient = require('@scality/arsenal').storage .metadata.mongoclient.MongoClientInterface; const { ObjectMD, ReplicationConfiguration } = require('@scality/arsenal').models; -const { VersionID } = require('@scality/arsenal').versioning; -const { extractVersionId } = require('../../lib/util/versioning'); const Config = require('../../lib/Config'); const BackbeatConsumer = require('../../lib/BackbeatConsumer'); @@ -19,7 +17,7 @@ const ObjectQueueEntry = require('../../lib/models/ObjectQueueEntry'); const MetricsProducer = require('../../lib/MetricsProducer'); const { metricsExtension, metricsTypeCompleted, metricsTypePendingOnly } = require('../ingestion/constants'); -const getContentType = require('./utils/contentTypeHelper'); +const { metadataPolicies, defaultMode } = require('./metadataPolicy'); const BucketMemState = require('./utils/BucketMemState'); const MongoProcessorMetrics = require('./MongoProcessorMetrics'); @@ -76,6 +74,8 @@ class MongoQueueProcessor { this._bootstrapList = null; this.logger = new Logger('Backbeat:Ingestion:MongoProcessor'); this.mongoClientConfig.logger = this.logger; + this._policy = + new metadataPolicies[mongoProcessorConfig.mode ?? defaultMode](); this._mongoClient = new MongoClient(this.mongoClientConfig); this._bucketMemState = new BucketMemState(Config); @@ -226,14 +226,7 @@ class MongoQueueProcessor { _getZenkoObjectMetadata(log, entry, versionId, done) { const bucket = entry.getBucket(); const key = entry.getObjectKey(); - const params = {}; - - // master keys with a 'null' version id comming from - // a versioning suspended bucket are considered a version - // we should not specify the version id in this case - if (versionId && !(entry.getIsNull && entry.getIsNull())) { - params.versionId = versionId; - } + const params = { versionId }; log.debug('getting zenko object metadata', { bucket, key, versionId, params }); @@ -259,78 +252,6 @@ class MongoQueueProcessor { }); } - /** - * Update ingested entry metadata fields: owner-id, owner-display-name - * @param {ObjectQueueEntry} entry - object queue entry object - * @param {BucketInfo} bucketInfo - bucket info object - * @return {undefined} - */ - _updateOwnerMD(entry, bucketInfo) { - // zenko bucket owner information is being set on ingested md - entry.setOwnerDisplayName(bucketInfo.getOwnerDisplayName()); - entry.setOwnerId(bucketInfo.getOwner()); - } - - /** - * Update ingested entry metadata fields: dataStoreName - * @param {ObjectQueueEntry} entry - object queue entry object - * @param {string} location - owner details - * @return {undefined} - */ - _updateObjectDataStoreName(entry, location) { - entry.setDataStoreName(location); - } - - /** - * Update ingested entry metadata location field. Each location change - * includes: key, dataStoreName, dataStoreType, dataStoreVersionId - * @param {ObjectQueueEntry} entry - object queue entry object - * @param {string} zenkoLocation - zenko storage location name - * @return {undefined} - */ - _updateLocations(entry, zenkoLocation) { - const locations = entry.getLocation(); - // if version id is undefined, we have a single null object. - // To hold reference to this null object, we need to encode "null" - // as its dataStoreVersionId - const dataStoreVersionId = entry.getVersionId() ? - entry.getEncodedVersionId() : 'null'; - let zenkoDataLocations; - if (!locations || locations.length === 0) { - zenkoDataLocations = [{ - key: entry.getObjectKey(), - size: 0, - start: 0, - dataStoreName: zenkoLocation, - dataStoreType: 'aws_s3', - dataStoreETag: `1:${emptyFileMd5}`, - dataStoreVersionId, - }]; - } else { - zenkoDataLocations = [{ - key: entry.getObjectKey(), - size: entry.getContentLength(), - start: 0, - dataStoreName: zenkoLocation, - dataStoreType: 'aws_s3', - dataStoreETag: `1:${entry.getContentMd5()}`, - dataStoreVersionId, - }]; - } - entry.setLocation(zenkoDataLocations); - } - - /** - * Update acl info on ingested object MD - * @param {ObjectQueueEntry} entry - object queue entry object - * @return {undefined} - */ - _updateAcl(entry) { - // reset acl info - const objectMDModel = new ObjectMD(); - entry.setAcl(objectMDModel.getAcl()); - } - /** * Update replication info on ingested object MD to match Zenko defined * replication info. @@ -388,31 +309,14 @@ class MongoQueueProcessor { _processDeleteOpQueueEntry(log, sourceEntry, location, bucketInfo, done) { const bucket = sourceEntry.getBucket(); const key = sourceEntry.getObjectKey(); - const entryVersionId = extractVersionId(sourceEntry.getObjectVersionedKey()); - - // Use x-amz-meta-scal-version-id if provided, instead of the actual versionId of the object. - // This should happen only for restored objects : in all other situations, both the source - // and ingested objects should have the same version id (and no x-amz-meta-scal-version-id - // metadata). - const scalVersionId = sourceEntry.getOverheadField('x-amz-meta-scal-version-id'); - const versionId = scalVersionId ? VersionID.decode(scalVersionId) : entryVersionId; + const versionId = this._policy.targetVersionId(sourceEntry); this.logger.debug('processing object delete', { bucket, key, versionId }); async.waterfall([ cb => this._getZenkoObjectMetadata(log, sourceEntry, versionId, cb), (zenkoObjMd, cb) => { - // Skip if the object is in a different location, i.e. when the delete was caused - // by restored-object expiration or transition. It works because the dataStoreName - // is updated before actually sending the object to GC to effectively delete the - // data. - const encode = versionId => (versionId ? VersionID.encode(versionId) : 'null'); - if (zenkoObjMd.dataStoreName !== location || - zenkoObjMd.location?.length !== 1 || - zenkoObjMd.location[0].dataStoreName !== location || - zenkoObjMd.location[0].key !== key || - (zenkoObjMd.location[0].dataStoreVersionId || 'null') !== encode(entryVersionId) - ) { + if (this._policy.skipsDelete(sourceEntry, zenkoObjMd, location)) { log.end().info('ignore delete entry, transitioned to another location', { entry: sourceEntry.getLogInfo(), location, @@ -485,24 +389,14 @@ class MongoQueueProcessor { _processObjectQueueEntry(log, sourceEntry, location, bucketInfo, done) { const bucket = sourceEntry.getBucket(); const key = sourceEntry.getObjectKey(); - const scalVersionId = sourceEntry.getValue()['x-amz-meta-scal-version-id']; - - this.logger.debug('processing object metadata', { bucket, key, scalVersionId }); + this.logger.debug('processing object metadata', { bucket, key }); const maybeGetZenkoObjectMetadata = cb => { - // NOTE: ZenkoObjMD is used for updating replication info, as well as validating the - // `x-amz-meta-scal-version-id` header of restored objects. If the Zenko bucket does - // not have repInfo set and the header is not set, then we can skip fetching. - const bucketRepInfo = bucketInfo.getReplicationConfiguration(); - if (!scalVersionId && !bucketRepInfo?.rules?.some(r => r.enabled)) { + if (this._policy.skipsMetadataFetch(sourceEntry, bucketInfo)) { return cb(); } - // Use x-amz-meta-scal-version-id if provided, instead of the actual versionId of the object. - // This should happen only for restored objects : in all other situations, both the source - // and ingested objects should have the same version id (and not x-amz-meta-scal-version-id - // metadata). - const versionId = scalVersionId ? VersionID.decode(scalVersionId) : sourceEntry.getVersionId(); + const versionId = this._policy.targetVersionId(sourceEntry); return this._getZenkoObjectMetadata(log, sourceEntry, versionId, cb); }; @@ -517,8 +411,9 @@ class MongoQueueProcessor { return done(err); } - const content = getContentType(sourceEntry, zenkoObjMd); - if (content.length === 0) { + const applied = this._policy.apply(sourceEntry, zenkoObjMd, + location, bucketInfo); + if (!applied) { this._normalizePendingMetric(location); log.end().debug('skipping duplicate entry', { method: 'MongoQueueProcessor._processObjectQueueEntry', @@ -529,22 +424,9 @@ class MongoQueueProcessor { return done(); } - if (zenkoObjMd) { - // Keep existing metadata fields, only need to update the tags - const tags = sourceEntry.getTags(); - sourceEntry._data = { ...zenkoObjMd }; // eslint-disable-line no-param-reassign - sourceEntry.setTags(tags); - } else { - // Update necessary metadata fields before saving to Zenko MongoDB - this._updateOwnerMD(sourceEntry, bucketInfo); - this._updateObjectDataStoreName(sourceEntry, location); - this._updateLocations(sourceEntry, location); - this._updateAcl(sourceEntry); - } - // Try to update replication info, if applicable - this._updateReplicationInfo(sourceEntry, bucketInfo, content, - zenkoObjMd); + this._updateReplicationInfo(sourceEntry, bucketInfo, + applied.content, zenkoObjMd); // Some object stores we are ingesting from (e.g. S3C) do // not populate the value.key property. @@ -553,11 +435,8 @@ class MongoQueueProcessor { const objVal = sourceEntry.getValue(); const params = {}; - // Versioning suspended entries will have a version id but also a isNull tag. - // These master keys are considered a version and do not have a duplicate version, - // we don't specify the version id and repairMaster in this case - if (sourceEntry.getVersionId() && !sourceEntry.getIsNull()) { - params.versionId = sourceEntry.getVersionId(); + if (applied.versionId) { + params.versionId = applied.versionId; params.repairMaster = true; } diff --git a/extensions/mongoProcessor/metadataPolicy/IngestionMetadataPolicy.js b/extensions/mongoProcessor/metadataPolicy/IngestionMetadataPolicy.js new file mode 100644 index 000000000..1b484159f --- /dev/null +++ b/extensions/mongoProcessor/metadataPolicy/IngestionMetadataPolicy.js @@ -0,0 +1,202 @@ +'use strict'; + +const { emptyFileMd5 } = require('@scality/arsenal').constants; +const { ObjectMD } = require('@scality/arsenal').models; +const { VersionID } = require('@scality/arsenal').versioning; + +const MetadataPolicy = require('./MetadataPolicy'); +const getContentType = require('../utils/contentTypeHelper'); +const DeleteOpQueueEntry = require('../../../lib/models/DeleteOpQueueEntry'); +const { extractVersionId } = require('../../../lib/util/versioning'); + +const scalVersionIdHeader = 'x-amz-meta-scal-version-id'; + +class IngestionMetadataPolicy extends MetadataPolicy { + skipsMetadataFetch(entry, bucketInfo) { + // the stored document is only read to update replication info and to + // validate the `x-amz-meta-scal-version-id` header of a restored + // object, so with neither in play the fetch is skipped + const bucketRepInfo = bucketInfo.getReplicationConfiguration(); + + return !this._scalVersionId(entry) && + !bucketRepInfo?.rules?.some(r => r.enabled); + } + + targetVersionId(entry) { + // x-amz-meta-scal-version-id wins where it is set, which happens only + // for a restored object: otherwise the source and the ingested object + // carry the same version id and the header is absent + const scalVersionId = this._scalVersionId(entry); + + if (entry instanceof DeleteOpQueueEntry) { + return scalVersionId ? + VersionID.decode(scalVersionId) : + extractVersionId(entry.getObjectVersionedKey()); + } + + // master keys with a 'null' version id comming from + // a versioning suspended bucket are considered a version + // we should not specify the version id in this case + if (entry.getIsNull()) { + return undefined; + } + return scalVersionId ? + VersionID.decode(scalVersionId) : entry.getVersionId(); + } + + /** + * The encoded `x-amz-meta-scal-version-id` of an entry, which an object + * carries in its metadata and a delete in its overhead fields. + * + * @param {ObjectQueueEntry|DeleteOpQueueEntry} entry - queue entry object + * @return {string|undefined} encoded scal version id, if any + */ + _scalVersionId(entry) { + if (entry instanceof DeleteOpQueueEntry) { + return entry.getOverheadField(scalVersionIdHeader); + } + return entry.getValue()[scalVersionIdHeader]; + } + + /** + * The version a composed entry is written under: its own -- the stored + * version's once merged into it -- or none for a master. + * + * @param {ObjectQueueEntry} entry - object queue entry object + * @return {string|undefined} version id, or undefined for the master + */ + _writtenVersionId(entry) { + // Versioning suspended entries will have a version id but also a isNull tag. + // These master keys are considered a version and do not have a duplicate version, + // we don't specify the version id and repairMaster in this case + if (entry.getVersionId() && !entry.getIsNull()) { + return entry.getVersionId(); + } + return undefined; + } + + apply(entry, objMD, location, bucketInfo) { + const content = getContentType(entry, objMD); + if (content.length === 0) { + return null; + } + + if (objMD) { + // Keep existing metadata fields, only need to update the tags + const tags = entry.getTags(); + entry._data = { ...objMD }; // eslint-disable-line no-param-reassign + entry.setTags(tags); + } else { + // Update necessary metadata fields before saving to Zenko MongoDB + this._updateOwnerMD(entry, bucketInfo); + this._updateObjectDataStoreName(entry, location); + this._updateLocations(entry, location); + this._updateAcl(entry); + } + + return { content, versionId: this._writtenVersionId(entry) }; + } + + /** + * Update ingested entry metadata fields: owner-id, owner-display-name + * @param {ObjectQueueEntry} entry - object queue entry object + * @param {BucketInfo} bucketInfo - bucket info object + * @return {undefined} + */ + _updateOwnerMD(entry, bucketInfo) { + // zenko bucket owner information is being set on ingested md + entry.setOwnerDisplayName(bucketInfo.getOwnerDisplayName()); + entry.setOwnerId(bucketInfo.getOwner()); + } + + /** + * Update ingested entry metadata fields: dataStoreName + * @param {ObjectQueueEntry} entry - object queue entry object + * @param {string} location - owner details + * @return {undefined} + */ + _updateObjectDataStoreName(entry, location) { + entry.setDataStoreName(location); + } + + /** + * Update ingested entry metadata location field. Each location change + * includes: key, dataStoreName, dataStoreType, dataStoreVersionId + * @param {ObjectQueueEntry} entry - object queue entry object + * @param {string} zenkoLocation - zenko storage location name + * @return {undefined} + */ + _updateLocations(entry, zenkoLocation) { + const locations = entry.getLocation(); + // if version id is undefined, we have a single null object. + // To hold reference to this null object, we need to encode "null" + // as its dataStoreVersionId + const dataStoreVersionId = entry.getVersionId() ? + entry.getEncodedVersionId() : 'null'; + let zenkoDataLocations; + if (!locations || locations.length === 0) { + zenkoDataLocations = [{ + key: entry.getObjectKey(), + size: 0, + start: 0, + dataStoreName: zenkoLocation, + dataStoreType: 'aws_s3', + dataStoreETag: `1:${emptyFileMd5}`, + dataStoreVersionId, + }]; + } else { + zenkoDataLocations = [{ + key: entry.getObjectKey(), + size: entry.getContentLength(), + start: 0, + dataStoreName: zenkoLocation, + dataStoreType: 'aws_s3', + dataStoreETag: `1:${entry.getContentMd5()}`, + dataStoreVersionId, + }]; + } + entry.setLocation(zenkoDataLocations); + } + + /** + * Update acl info on ingested object MD + * @param {ObjectQueueEntry} entry - object queue entry object + * @return {undefined} + */ + _updateAcl(entry) { + // reset acl info + const objectMDModel = new ObjectMD(); + entry.setAcl(objectMDModel.getAcl()); + } + + skipsDelete(entry, objMD, location) { + // an object stored elsewhere is left alone: the delete was caused by + // restored-object expiration or transition, both of which update + // dataStoreName before sending the object to GC + return !this._holdsEntryData(objMD, location, entry.getObjectKey(), + extractVersionId(entry.getObjectVersionedKey())); + } + + /** + * Whether the stored object still holds the data the delete entry names. + * + * @param {Object} objMD - metadata fetched from mongo + * @param {string} location - zenko storage location name + * @param {string} key - object key + * @param {string|null} versionId - the entry's own decoded version id, + * which the data location names + * @return {boolean} true if the object is stored where the entry says + */ + _holdsEntryData(objMD, location, key, versionId) { + const encode = vid => (vid ? VersionID.encode(vid) : 'null'); + + return objMD.dataStoreName === location && + objMD.location?.length === 1 && + objMD.location[0].dataStoreName === location && + objMD.location[0].key === key && + (objMD.location[0].dataStoreVersionId || 'null') === + encode(versionId); + } +} + +module.exports = IngestionMetadataPolicy; diff --git a/extensions/mongoProcessor/metadataPolicy/MetadataPolicy.js b/extensions/mongoProcessor/metadataPolicy/MetadataPolicy.js new file mode 100644 index 000000000..53f0b0c4f --- /dev/null +++ b/extensions/mongoProcessor/metadataPolicy/MetadataPolicy.js @@ -0,0 +1,84 @@ +'use strict'; + +const assert = require('assert'); + +/** + * A metadata policy holds the decisions that differ between the streams the + * mongo-processor writes: out-of-band ingestion, where identity and placement + * are rewritten to local values because the source system's accounts and + * locations do not exist here, and pull replication, where they are replicated + * and so are applied as they arrive. + * + * Everything else about processing an entry is shared. + */ +class MetadataPolicy { + /** + * Whether the entry can be applied without fetching the metadata already + * stored for the object. + * + * This method must be implemented by subclasses of MetadataPolicy + * @param {ObjectQueueEntry} entry - object queue entry object + * @param {BucketInfo} bucketInfo - bucket info object + * @return {boolean} true if the stored metadata is not needed + */ + skipsMetadataFetch(entry, bucketInfo) { // eslint-disable-line no-unused-vars + assert(false, + 'sub-classes of MetadataPolicy must implement ' + + 'the skipsMetadataFetch() method'); + } + + /** + * Which version of the object on this site an entry acts on: the one a + * put reads, and the one a delete reads and removes. + * + * This method must be implemented by subclasses of MetadataPolicy + * @param {ObjectQueueEntry|DeleteOpQueueEntry} entry - queue entry object + * @return {string|undefined} version id, or undefined for the master + */ + targetVersionId(entry) { // eslint-disable-line no-unused-vars + assert(false, + 'sub-classes of MetadataPolicy must implement ' + + 'the targetVersionId() method'); + } + + /** + * Apply the entry onto the metadata already stored, leaving the entry + * holding the document to write, and report what the write changes as + * replicationInfo content values, and which version it is written under. + * + * An entry that changes nothing is not written at all, so it is left + * untouched. + * + * This method must be implemented by subclasses of MetadataPolicy + * @param {ObjectQueueEntry} entry - object queue entry object + * @param {Object|undefined} objMD - metadata fetched from mongo + * @param {string} location - zenko storage location name + * @param {BucketInfo} bucketInfo - bucket info object + * @return {Object|null} null if the entry changes nothing, otherwise + * `content`, the array of ReplicationInfo Content Type, and `versionId`, + * the version the document is written under, the master following it, + * or undefined to write the master alone + */ + apply(entry, objMD, location, bucketInfo) { // eslint-disable-line no-unused-vars + assert(false, + 'sub-classes of MetadataPolicy must implement ' + + 'the apply() method'); + } + + /** + * Whether a delete entry no longer applies to the object as stored. + * + * This method must be implemented by subclasses of MetadataPolicy + * @param {DeleteOpQueueEntry} entry - delete object entry + * @param {Object} objMD - metadata fetched from mongo + * @param {string} location - zenko storage location name + * @return {boolean} true if the object should be left alone + */ + skipsDelete(entry, objMD, location) { // eslint-disable-line no-unused-vars + assert(false, + 'sub-classes of MetadataPolicy must implement ' + + 'the skipsDelete() method'); + } +} + +module.exports = MetadataPolicy; diff --git a/extensions/mongoProcessor/metadataPolicy/index.js b/extensions/mongoProcessor/metadataPolicy/index.js new file mode 100644 index 000000000..b11643cdc --- /dev/null +++ b/extensions/mongoProcessor/metadataPolicy/index.js @@ -0,0 +1,10 @@ +const IngestionMetadataPolicy = require('./IngestionMetadataPolicy'); + +const metadataPolicies = { + ingestion: IngestionMetadataPolicy, +}; + +module.exports = { + metadataPolicies, + defaultMode: 'ingestion', +}; diff --git a/tests/functional/ingestion/MongoQueueProcessor.js b/tests/functional/ingestion/MongoQueueProcessor.js index ffab98101..3cf247f2d 100644 --- a/tests/functional/ingestion/MongoQueueProcessor.js +++ b/tests/functional/ingestion/MongoQueueProcessor.js @@ -300,7 +300,8 @@ describe('MongoQueueProcessor', function mqp() { .setNullVersionId(NEW_VERSION_ID) .setIsNull(true); const entry = new ObjectQueueEntry(BUCKET, versionKey, objmd); - mqp._getZenkoObjectMetadata(fakeLogger, entry, NEW_VERSION_ID, (err, res) => { + const versionId = mqp._policy.targetVersionId(entry); + mqp._getZenkoObjectMetadata(fakeLogger, entry, versionId, (err, res) => { assert.ifError(err); assert(res); diff --git a/tests/unit/mongoProcessor/MetadataPolicy.spec.js b/tests/unit/mongoProcessor/MetadataPolicy.spec.js new file mode 100644 index 000000000..9a01483d3 --- /dev/null +++ b/tests/unit/mongoProcessor/MetadataPolicy.spec.js @@ -0,0 +1,22 @@ +'use strict'; + +const assert = require('assert'); + +const MetadataPolicy = + require('../../../extensions/mongoProcessor/metadataPolicy/MetadataPolicy'); + +describe('MetadataPolicy', () => { + const policy = new MetadataPolicy(); + + [ + 'skipsMetadataFetch', + 'targetVersionId', + 'apply', + 'skipsDelete', + ].forEach(method => it(`should refuse to ${method}() without an ` + + 'implementation', () => { + assert.throws(() => policy[method](), + new RegExp('sub-classes of MetadataPolicy must implement the ' + + `${method}\\(\\) method`)); + })); +}); diff --git a/tests/unit/mongoProcessor/MongoProcessorConfigValidator.spec.js b/tests/unit/mongoProcessor/MongoProcessorConfigValidator.spec.js index 9723c20a4..e4de5fe43 100644 --- a/tests/unit/mongoProcessor/MongoProcessorConfigValidator.spec.js +++ b/tests/unit/mongoProcessor/MongoProcessorConfigValidator.spec.js @@ -22,7 +22,7 @@ describe('MongoProcessorConfigValidator log override', () => { }); it('should leave log undefined when not set, deferring to global config.log', () => { - const validated = configValidator(globalConfig, baseExtConfig); + const validated = configValidator(globalConfig, { ...baseExtConfig }); assert.strictEqual(validated.log, undefined); }); @@ -41,3 +41,21 @@ describe('MongoProcessorConfigValidator log override', () => { } }); }); + +describe('MongoProcessorConfigValidator mode', () => { + it('should default to ingestion so an existing config is unchanged', () => { + const validated = configValidator(globalConfig, { ...baseExtConfig }); + assert.strictEqual(validated.mode, 'ingestion'); + }); + + it('should reject a mode with no implementation', () => { + let err; + try { + configValidator(globalConfig, { ...baseExtConfig, mode: 'sideways' }); + } catch (e) { + err = e; + } + assert(err, 'expected configValidator to throw on an unknown mode'); + assert.match(err.message, /mode/); + }); +}); From f64d38dcee83ba287f22dee19ce5fbccf42ee69b Mon Sep 17 00:00:00 2001 From: Thomas Flament Date: Tue, 29 Sep 2026 12:30:00 +0200 Subject: [PATCH 5/8] Add the pull replication metadata policy The D/R metadata sink replicates production's objects, accounts included, so it applies what the source-side pipeline sends rather than rewriting it into something local: a new object is written as it arrives, with its ACLs reset as they are not replicated, and an update merges the entry's tags and object-lock state into the stored document, cleared values included -- removing a legal hold is an update. Every update is written, a replay included, rather than diffed against what is stored. The stored document keeps everything else, and above all its placement. The copy engine rewrites location and dataStoreName to a local location after the first write, and applying the entry's would send reads back to the source and leak a local copy that is never garbage-collected; a version still on the remote site takes the entry's. A restore this site performed is kept the same way. Several things follow that were unreachable before: - the stored document is always read, because it is what tells a first write from an update, and the entry cannot -- an insert is redelivered on replay and overlaps the bootstrap dump, so it is no promise that the object is absent here; - a delete always applies, where the ingestion guard skips one whose object has moved location, which for a replicated object it always has; - a scal version id names a version of the system an object was ingested from, and is ignored. The D/R source keys every entry by the document it comes from, so the sink acts on that very document: a version by its version id, null versions included, and a master as the master document itself, which only reaches the sink for an object with no version of its own. A master document that copies a version, or the latest version mongo returns in place of a missing master, is not what such an entry targets: the entry is written over it. The mock client returns the document it was given, so null versions are tested against mongodb. Issue: BB-811 --- .../PullReplicationMetadataPolicy.js | 105 +++++ .../mongoProcessor/metadataPolicy/index.js | 2 + .../ingestion/MongoQueueProcessor.js | 373 ++++++++++++++++++ tests/functional/ingestion/PullReplication.js | 150 +++++++ .../MongoProcessorConfigValidator.spec.js | 15 + 5 files changed, 645 insertions(+) create mode 100644 extensions/mongoProcessor/metadataPolicy/PullReplicationMetadataPolicy.js create mode 100644 tests/functional/ingestion/PullReplication.js diff --git a/extensions/mongoProcessor/metadataPolicy/PullReplicationMetadataPolicy.js b/extensions/mongoProcessor/metadataPolicy/PullReplicationMetadataPolicy.js new file mode 100644 index 000000000..d8fd2de09 --- /dev/null +++ b/extensions/mongoProcessor/metadataPolicy/PullReplicationMetadataPolicy.js @@ -0,0 +1,105 @@ +'use strict'; + +const { ObjectMD } = require('@scality/arsenal').models; + +const MetadataPolicy = require('./MetadataPolicy'); +const DeleteOpQueueEntry = require('../../../lib/models/DeleteOpQueueEntry'); +const { extractVersionId } = require('../../../lib/util/versioning'); +const getContentType = require('../utils/contentTypeHelper'); +const locationsConfig = require('../../../conf/locationConfig.json') || {}; + +class PullReplicationMetadataPolicy extends MetadataPolicy { + skipsMetadataFetch(entry, bucketInfo) { // eslint-disable-line no-unused-vars + // never: only the stored document tells a first write from an update, + // and a source insert is redelivered on replay and overlaps the + // bootstrap dump + return false; + } + + targetVersionId(entry) { + // the document the entry comes from: a version by its key, a master + // as the master document itself + const documentVersionId = extractVersionId(entry.getObjectVersionedKey()); + if (documentVersionId) { + return documentVersionId; + } + if (entry instanceof DeleteOpQueueEntry || entry.getIsNull()) { + return undefined; + } + return entry.getVersionId(); + } + + apply(entry, objMD) { + const versionId = this.targetVersionId(entry); + const stored = this._storedObject(versionId, objMD); + + if (stored) { + const content = getContentType(entry, stored); + this._mergeStoredMetadata(entry, stored); + return { content: content.length !== 0 ? content : ['METADATA'], versionId }; + } + + // the source-side pipeline shaped everything else this object keeps; + // ACLs are not replicated, so they reset as for an ingested object + entry.setAcl(new ObjectMD().getAcl()); + return { content: getContentType(entry), versionId }; + } + + /** + * The stored document, when it is the object the entry targets. The + * master document is only that object when it has no version of its own: + * a master that is a copy of a version, or the latest version mongo + * returns in place of a missing master, is not. + * + * @param {string|undefined} versionId - the version the entry targets + * @param {Object|undefined} objMD - metadata fetched from mongo + * @return {Object|undefined} the stored object, if any + */ + _storedObject(versionId, objMD) { + if (versionId === undefined && objMD?.versionId && !objMD.isNull) { + return undefined; + } + return objMD; + } + + /** + * A version whose data still lives on the remote site. + * + * @param {Object} objMD - metadata fetched from mongo + * @return {boolean} true if the stored version is not localized + */ + _isNotLocalized(objMD) { + return Boolean(locationsConfig[objMD.dataStoreName]?.isCRR); + } + + /** + * Keep the stored document and apply what an update brings. + * + * @param {ObjectQueueEntry} entry - object queue entry object + * @param {Object} objMD - metadata fetched from mongo + * @return {undefined} + */ + _mergeStoredMetadata(entry, objMD) { + // as the entry carries them, so that a value it leaves out is left + // out of the document too + const { tags, retentionMode, retentionDate, legalHold } = entry.getValue(); + const dataStoreName = entry.getDataStoreName(); + const location = entry.getLocation(); + + // eslint-disable-next-line no-param-reassign + entry._data = { ...objMD, tags, retentionMode, retentionDate, legalHold }; + + // update data location if data has not been pulled yet + if (locationsConfig[objMD.dataStoreName]?.isCRR) { + entry.setDataStoreName(dataStoreName); + entry.setLocation(location); + } + } + + skipsDelete(entry, objMD, location) { // eslint-disable-line no-unused-vars + // must process every deletion, whatever the state of the object + return false; + } +} + +module.exports = PullReplicationMetadataPolicy; diff --git a/extensions/mongoProcessor/metadataPolicy/index.js b/extensions/mongoProcessor/metadataPolicy/index.js index b11643cdc..0e861a2bf 100644 --- a/extensions/mongoProcessor/metadataPolicy/index.js +++ b/extensions/mongoProcessor/metadataPolicy/index.js @@ -1,7 +1,9 @@ const IngestionMetadataPolicy = require('./IngestionMetadataPolicy'); +const PullReplicationMetadataPolicy = require('./PullReplicationMetadataPolicy'); const metadataPolicies = { ingestion: IngestionMetadataPolicy, + dr: PullReplicationMetadataPolicy, }; module.exports = { diff --git a/tests/functional/ingestion/MongoQueueProcessor.js b/tests/functional/ingestion/MongoQueueProcessor.js index 3cf247f2d..33afae4d9 100644 --- a/tests/functional/ingestion/MongoQueueProcessor.js +++ b/tests/functional/ingestion/MongoQueueProcessor.js @@ -29,6 +29,9 @@ const bootstrapList = config.extensions.replication.destination.bootstrapList; const BUCKET = 'mqp-test-bucket'; const KEY = 'testkey1'; const LOCATION = 'us-east-1'; +// a location holding data owned by the remote site, as the source one is +// until the copy engine pulls the version's data +const CRR_LOCATION = 'location-crr-source'; const VERSION_ID = '98445230573829999999RG001 15.144.0'; // new version id > existing version id const NEW_VERSION_ID = '98445235075994999999RG001 14.90.2'; @@ -1071,3 +1074,373 @@ describe('MongoQueueProcessor', function mqp() { }); }); }); + +describe('MongoQueueProcessor in dr mode', function drMode() { + this.timeout(5000); + + let mqp; + let mongoClient; + + before(() => { + mqp = new MongoQueueProcessorMock(kafkaConfig, + { ...mongoProcessorConfig, mode: 'dr' }, mongoClientConfig, mConfig); + mqp.start(); + + mongoClient = mqp._mongoClient; + }); + + afterEach(() => { + mqp.reset(); + sinon.restore(); + }); + + function processEntry(entry, next) { + return async.waterfall([ + cb => mongoClient.getBucketAttributes(BUCKET, fakeLogger, cb), + (bucketInfo, cb) => mqp._processObjectQueueEntry(fakeLogger, entry, LOCATION, bucketInfo, cb), + ], next); + } + + it('should apply a new object as the source describes it', done => { + const key = 'dr-new-key'; + const versionKey = `${key}${VID_SEP}${VERSION_ID}`; + const objmd = new ObjectMD() + .setKey(key) + .setVersionId(VERSION_ID) + .setOwnerId('source-owner-id') + .setOwnerDisplayName('source-owner') + .setDataStoreName('cold-location') + // ACLs are not replicated: the matrix in the design resets them + .setAcl({ Canned: '', FULL_CONTROL: ['source-grantee'], WRITE_ACP: [], READ: [], READ_ACP: [] }); + const entry = new ObjectQueueEntry(BUCKET, versionKey, objmd); + + processEntry(entry, err => { + assert.ifError(err); + + const added = mqp.getAdded(); + assert.strictEqual(added.length, 1); + const { objVal } = added[0]; + assert.strictEqual(objVal['owner-id'], 'source-owner-id'); + assert.strictEqual(objVal['owner-display-name'], 'source-owner'); + assert.strictEqual(objVal.dataStoreName, 'cold-location'); + assert.deepStrictEqual(objVal.acl, new ObjectMD().getAcl()); + done(); + }); + }); + + it('should keep the stored placement when updating an object whose data was pulled', done => { + const versionKey = `${KEY}${VID_SEP}${VERSION_ID}`; + // the copy engine has since moved the object to a local location, so + // the entry's own placement must not be written back + const objmd = new ObjectMD() + .setKey(KEY) + .setVersionId(VERSION_ID) + .setTags({ mytag: 'mytags-value' }) + .setDataStoreName('source-site') + .setLocation([{ key: KEY, dataStoreName: 'source-site' }]) + .setLegalHold(true); + const entry = new ObjectQueueEntry(BUCKET, versionKey, objmd); + + processEntry(entry, err => { + assert.ifError(err); + + const added = mqp.getAdded(); + assert.strictEqual(added.length, 1); + const { objVal } = added[0]; + // placement from the stored document + assert.strictEqual(objVal.dataStoreName, LOCATION); + assert.strictEqual(objVal.location[0].dataStoreName, LOCATION); + // mutable metadata from the entry + assert.strictEqual(objVal.legalHold, true); + done(); + }); + }); + + it('should rewrite a replayed entry as it stands', done => { + // at-least-once delivery means the same entry arrives twice; the second + // writes the document back unchanged + const versionKey = `${KEY}${VID_SEP}${VERSION_ID}`; + const objmd = new ObjectMD() + .setKey(KEY) + .setVersionId(VERSION_ID) + .setTags({ mytag: 'mytags-value' }) + .setDataStoreName(LOCATION); + const entry = new ObjectQueueEntry(BUCKET, versionKey, objmd); + + processEntry(entry, err => { + assert.ifError(err); + const added = mqp.getAdded(); + assert.strictEqual(added.length, 1); + assert.deepStrictEqual(added[0].objVal.tags, { mytag: 'mytags-value' }); + done(); + }); + }); + + it('should take the entry placement for a version still on the source', done => { + // data not pulled yet, so the source still describes where it is: + // this is the production-side archive of an already-replicated version + sinon.stub(mongoClient, 'getObject').callsFake((b, k, p, l, cb) => cb(null, new ObjectMD() + .setKey(KEY) + .setVersionId(VERSION_ID) + .setTags({ mytag: 'mytags-value' }) + .setDataStoreName(CRR_LOCATION) + ._data)); + const versionKey = `${KEY}${VID_SEP}${VERSION_ID}`; + const objmd = new ObjectMD() + .setKey(KEY) + .setVersionId(VERSION_ID) + .setTags({ mytag: 'mytags-value' }) + .setDataStoreName('cold-location'); + const entry = new ObjectQueueEntry(BUCKET, versionKey, objmd); + + processEntry(entry, err => { + assert.ifError(err); + + const added = mqp.getAdded(); + assert.strictEqual(added.length, 1); + assert.strictEqual(added[0].objVal.dataStoreName, 'cold-location'); + done(); + }); + }); + + it('should not skip an update that only changes the object lock', done => { + const versionKey = `${KEY}${VID_SEP}${VERSION_ID}`; + // same tags as the stored object: the ingestion diff would call this a + // duplicate and drop the retention with it + const objmd = new ObjectMD() + .setKey(KEY) + .setVersionId(VERSION_ID) + .setTags({ mytag: 'mytags-value' }) + .setRetentionMode('GOVERNANCE') + .setRetentionDate('2099-01-01T00:00:00.000Z'); + const entry = new ObjectQueueEntry(BUCKET, versionKey, objmd); + + processEntry(entry, err => { + assert.ifError(err); + + const added = mqp.getAdded(); + assert.strictEqual(added.length, 1); + assert.strictEqual(added[0].objVal.retentionMode, 'GOVERNANCE'); + assert.strictEqual(added[0].objVal.retentionDate, '2099-01-01T00:00:00.000Z'); + done(); + }); + }); + + it('should read the stored object even when the bucket has no replication configuration', done => { + const getObject = sinon.spy(mongoClient, 'getObject'); + const versionKey = `${KEY}${VID_SEP}${VERSION_ID}`; + const objmd = new ObjectMD() + .setKey(KEY) + .setVersionId(VERSION_ID) + .setTags({ mytag: 'mytags-value' }) + .setLegalHold(true); + const entry = new ObjectQueueEntry(BUCKET, versionKey, objmd); + + async.waterfall([ + cb => mongoClient.getBucketAttributes(BUCKET, fakeLogger, cb), + (bucketInfo, cb) => { + sinon.stub(bucketInfo, 'getReplicationConfiguration').returns(null); + return mqp._processObjectQueueEntry(fakeLogger, entry, LOCATION, bucketInfo, cb); + }, + ], err => { + assert.ifError(err); + + sinon.assert.called(getObject); + // the merge happened, so the entry was recognised as an update + assert.strictEqual(mqp.getAdded()[0].objVal.dataStoreName, LOCATION); + done(); + }); + }); + + it('should ignore a scal version id the source object carries', done => { + // production may itself have ingested or cold-restored this object, so + // it can carry a scal version id of its own; it names a version of the + // system production ingested from, not anything here + const versionKey = `${KEY}${VID_SEP}${VERSION_ID}`; + const objmd = new ObjectMD() + .setKey(KEY) + .setVersionId(VERSION_ID) + .setTags({ mytag: 'mytags-value' }) + .setLegalHold(true) + .setUserMetadata({ + 'x-amz-meta-scal-version-id': encode(NEW_VERSION_ID), + }); + const entry = new ObjectQueueEntry(BUCKET, versionKey, objmd); + const getObject = sinon.spy(mongoClient, 'getObject'); + + processEntry(entry, err => { + assert.ifError(err); + + // read, and written back, on the version the entry names + assert.strictEqual(getObject.getCall(0).args[2].versionId, VERSION_ID); + assert.strictEqual(mqp.getAdded()[0].key, `${KEY}${VID_SEP}${VERSION_ID}`); + done(); + }); + }); + + it('should delete an object whose location differs from the bucket\'s', done => { + const versionKey = `${KEY}${VID_SEP}${VERSION_ID}`; + const entry = new DeleteOpQueueEntry(BUCKET, versionKey); + sinon.stub(mongoClient, 'getObject').callsFake((b, k, p, l, cb) => + cb(null, new ObjectMD() + .setKey(KEY) + .setVersionId(VERSION_ID) + // cold, or still pointing at the source: never the bucket's + // own location constraint + .setDataStoreName('cold-location') + .setLocation(null) + ._data)); + + async.waterfall([ + cb => mongoClient.getBucketAttributes(BUCKET, fakeLogger, cb), + (bucketInfo, cb) => mqp._processDeleteOpQueueEntry(fakeLogger, entry, LOCATION, bucketInfo, cb), + ], err => { + assert.ifError(err); + + const deleted = mqp.getDeleted(); + assert.strictEqual(deleted.length, 1); + assert.strictEqual(deleted[0].versionId, VERSION_ID); + done(); + }); + }); + + // a cold object with no version of its own, restored on this site + const STORED_ARCHIVE = { archiveInfo: { archiveId: 'stored-archive', archiveVersion: 1 } }; + const STORED_LAST_MODIFIED = '2026-09-01T10:00:00.000Z'; + + it('should write over a master that is a copy of a version', done => { + // versioning suspended on the source: a new null object is put over + // the master, which still copies a real version here + sinon.stub(mongoClient, 'getObject').callsFake((b, k, p, l, cb) => + cb(null, new ObjectMD() + .setKey(KEY) + .setVersionId(VERSION_ID) + .setTags({ stored: 'tag' }) + .setDataStoreName(LOCATION) + .setLastModified(STORED_LAST_MODIFIED) + ._data)); + // the same date, so that only the version tells the two apart + const objmd = new ObjectMD() + .setKey(KEY) + .setVersionId(NEW_VERSION_ID) + .setIsNull(true) + .setTags({ entry: 'tag' }) + .setDataStoreName(LOCATION) + .setLastModified(STORED_LAST_MODIFIED); + const entry = new ObjectQueueEntry(BUCKET, KEY, objmd); + + processEntry(entry, err => { + assert.ifError(err); + + const added = mqp.getAdded(); + assert.strictEqual(added.length, 1); + // written in place, under the object key alone + assert.strictEqual(added[0].key, KEY); + assert.strictEqual(added[0].objVal.versionId, NEW_VERSION_ID); + assert.deepStrictEqual(added[0].objVal.tags, { entry: 'tag' }); + done(); + }); + }); + + function storeUnversioned() { + sinon.stub(mongoClient, 'getObject').callsFake((b, k, p, l, cb) => + cb(null, new ObjectMD() + .setKey(KEY) + .setContentMd5('7d793037a0760186574b0282f2f435e7') + .setContentType('text/plain') + .setUserMetadata({ 'x-amz-meta-colour': 'red' }) + .setDataStoreName(LOCATION) + .setArchive(STORED_ARCHIVE) + .setLastModified(STORED_LAST_MODIFIED) + .setAmzRestore({ 'ongoing-request': false }) + ._data)); + } + + it('should merge an update to an object with no version of its own', done => { + storeUnversioned(); + // the pipeline strips the restore state, which is this site's own + const objmd = new ObjectMD() + .setKey(KEY) + .setContentMd5('7d793037a0760186574b0282f2f435e7') + .setContentType('application/octet-stream') + .setUserMetadata({ 'x-amz-meta-colour': 'red' }) + .setDataStoreName(LOCATION) + .setArchive(STORED_ARCHIVE) + .setLastModified(STORED_LAST_MODIFIED) + .setTags({ mytag: 'mytags-value' }); + const entry = new ObjectQueueEntry(BUCKET, KEY, objmd); + + processEntry(entry, err => { + assert.ifError(err); + + const added = mqp.getAdded(); + assert.strictEqual(added.length, 1); + const { objVal } = added[0]; + assert.deepStrictEqual(objVal.tags, { mytag: 'mytags-value' }); + assert.deepStrictEqual(objVal['x-amz-restore'], { 'ongoing-request': false }); + // not a field an update brings + assert.strictEqual(objVal['content-type'], 'text/plain'); + done(); + }); + }); + + it('should rewrite a replayed entry for an object with no version of its own as it stands', done => { + const objmd = new ObjectMD() + .setKey(KEY) + .setContentMd5('9e107d9d372bb6826bd81d3542a419d6') + .setUserMetadata({ 'x-amz-meta-colour': 'blue' }) + .setDataStoreName(LOCATION) + .setArchive(STORED_ARCHIVE) + .setLastModified(STORED_LAST_MODIFIED); + const entry = new ObjectQueueEntry(BUCKET, KEY, objmd); + // what the first delivery left behind, as mongo gives it back + const stored = JSON.parse(JSON.stringify({ + ...entry.getValue(), + acl: new ObjectMD().getAcl(), + })); + sinon.stub(mongoClient, 'getObject').callsFake((b, k, p, l, cb) => cb(null, stored)); + + processEntry(entry, err => { + assert.ifError(err); + const added = mqp.getAdded(); + assert.strictEqual(added.length, 1); + const { objVal } = added[0]; + for (const field of ['key', 'content-md5', 'x-amz-meta-colour', 'archive', 'last-modified', 'acl']) { + assert.deepStrictEqual(objVal[field], stored[field], field); + } + done(); + }); + }); + + it('should keep merging an update to a version of its own', done => { + // a version is immutable, so an entry for one carries an update and not + // a rewrite: the stored document still decides what it keeps + const versionKey = `${KEY}${VID_SEP}${VERSION_ID}`; + sinon.stub(mongoClient, 'getObject').callsFake((b, k, p, l, cb) => + cb(null, new ObjectMD() + .setKey(KEY) + .setVersionId(VERSION_ID) + .setUserMetadata({ 'x-amz-meta-colour': 'red' }) + .setDataStoreName(LOCATION) + ._data)); + const objmd = new ObjectMD() + .setKey(KEY) + .setVersionId(VERSION_ID) + .setTags({ mytag: 'mytags-value' }) + .setUserMetadata({ 'x-amz-meta-colour': 'blue' }) + .setDataStoreName(LOCATION); + const entry = new ObjectQueueEntry(BUCKET, versionKey, objmd); + + processEntry(entry, err => { + assert.ifError(err); + + const added = mqp.getAdded(); + assert.strictEqual(added.length, 1); + const { objVal } = added[0]; + assert.deepStrictEqual(objVal.tags, { mytag: 'mytags-value' }); + // not a field an update brings + assert.strictEqual(objVal['x-amz-meta-colour'], 'red'); + done(); + }); + }); +}); diff --git a/tests/functional/ingestion/PullReplication.js b/tests/functional/ingestion/PullReplication.js new file mode 100644 index 000000000..3b22fa11b --- /dev/null +++ b/tests/functional/ingestion/PullReplication.js @@ -0,0 +1,150 @@ +'use strict'; + +const assert = require('assert'); +const { promisify } = require('util'); +const sinon = require('sinon'); +const { ObjectMD, BucketInfo } = require('@scality/arsenal').models; +const { VersionID } = require('@scality/arsenal').versioning; +const VID_SEP = require('@scality/arsenal').versioning.VersioningConstants + .VersionId.Separator; + +const config = require('../../config.json'); +const MongoQueueProcessor = + require('../../../extensions/mongoProcessor/MongoQueueProcessor'); +const fakeLogger = require('../../utils/fakeLogger'); + +const REPLICATION_GROUP_ID = 'RG001'; +const LOCATION = 'us-east-1'; +const COLD_LOCATION = 'cold-location'; +const KEY = 'pull-replication-key'; + +// the mock client hands back the very document it was given, so what it +// cannot show is how mongo lays versions out and what it gives back +describe('MongoQueueProcessor in dr mode against mongodb', function drMongo() { + this.timeout(30000); + + let mqp; + let mongoClient; + let bucket; + let bucketId = 0; + + before(async () => { + mqp = new MongoQueueProcessor(config.kafka, + { ...config.extensions.mongoProcessor, mode: 'dr' }, + { + ...config.queuePopulator.mongo, + replicaSet: 'rs0', + writeConcern: 'majority', + database: `pullreplication${Date.now()}`, + replicationGroupId: REPLICATION_GROUP_ID, + }, + {}); + mqp._mProducer = { publishMetrics: () => {}, close: () => {} }; + mqp._bootstrapList = []; + mongoClient = mqp._mongoClient; + await promisify(mongoClient.setup.bind(mongoClient))(); + }); + + after(async () => { + await mongoClient.db.dropDatabase(); + await promisify(mongoClient.close.bind(mongoClient))(); + }); + + beforeEach(async () => { + bucketId += 1; + bucket = `pull-replication-bucket-${bucketId}`; + const bucketInfo = new BucketInfo(bucket, 'owner-id', 'owner', + new Date().toJSON(), BucketInfo.currentModelVersion()); + bucketInfo.setLocationConstraint(LOCATION); + bucketInfo.setVersioningConfiguration({ Status: 'Enabled' }); + await promisify(mongoClient.createBucket.bind(mongoClient))( + bucket, bucketInfo, fakeLogger); + }); + + afterEach(() => sinon.restore()); + + function coldObject(versionId, fields = {}) { + const objmd = new ObjectMD() + .setKey(KEY) + .setContentLength(1024) + .setContentMd5('9e107d9d372bb6826bd81d3542a419d6') + .setLastModified('2026-09-01T10:00:00.000Z') + .setAmzStorageClass(COLD_LOCATION) + .setDataStoreName(COLD_LOCATION) + .setArchive({ archiveInfo: { archiveId: `archive-${versionId || 'null'}`, archiveVersion: 1 } }); + if (versionId) { + objmd.setVersionId(versionId); + } + const value = { ...objmd.getValue(), ...fields }; + // the source pipeline strips the placement + delete value.location; + return value; + } + + // the entry the source pipeline produces for the document stored under + // `documentKey` + function replicate(documentKey, value) { + return promisify(mqp.processKafkaEntry.bind(mqp))({ + value: JSON.stringify({ type: 'put', bucket, key: documentKey, value }), + }); + } + + function getObject(versionId) { + return promisify(mongoClient.getObject.bind(mongoClient))( + bucket, KEY, versionId ? { versionId } : {}, fakeLogger); + } + + // a null master is moved into a version document of its own once a new + // version is put over it, and the master then follows the new version + [ + { + title: 'written while versioning was suspended', + nullVersionId: VersionID.generateVersionId('', REPLICATION_GROUP_ID), + master: versionId => coldObject(versionId, { isNull: true }), + }, + { + title: 'written before versioning was enabled', + nullVersionId: VersionID.getInfVid(REPLICATION_GROUP_ID), + master: () => coldObject(undefined), + }, + ].forEach(({ title, nullVersionId, master }) => + it(`should keep a null version ${title} once a newer one replicates`, async () => { + const newVersionId = VersionID.generateVersionId('', REPLICATION_GROUP_ID); + + const nullMaster = master(nullVersionId); + await replicate(KEY, nullMaster); + await replicate(`${KEY}${VID_SEP}${nullVersionId}`, { + ...nullMaster, + versionId: nullVersionId, + isNull: true, + originOp: 's3:StoreNullVersion', + }); + const newVersion = coldObject(newVersionId, { nullVersionId }); + await replicate(`${KEY}${VID_SEP}${newVersionId}`, newVersion); + await replicate(KEY, newVersion); + + const nullVersion = await getObject(nullVersionId); + assert.strictEqual(nullVersion.isNull, true); + assert.deepStrictEqual(nullVersion.archive, nullMaster.archive); + assert.strictEqual((await getObject(newVersionId)).versionId, newVersionId); + assert.strictEqual((await getObject()).versionId, newVersionId); + })); + + [ + { title: 'a version', versioned: true }, + { title: 'an object with no version of its own', versioned: false }, + ].forEach(({ title, versioned }) => + it(`should leave the document unchanged on a replayed entry for ${title}`, async () => { + const versionId = versioned ? + VersionID.generateVersionId('', REPLICATION_GROUP_ID) : undefined; + const documentKey = versionId ? `${KEY}${VID_SEP}${versionId}` : KEY; + // a dotted key is escaped when stored, and unescaped when read back + const value = coldObject(versionId, { tags: { 'dotted.key': 'value' } }); + + await replicate(documentKey, value); + const written = await getObject(versionId); + await replicate(documentKey, value); + + assert.deepStrictEqual(await getObject(versionId), written); + })); +}); diff --git a/tests/unit/mongoProcessor/MongoProcessorConfigValidator.spec.js b/tests/unit/mongoProcessor/MongoProcessorConfigValidator.spec.js index e4de5fe43..ea7d60f11 100644 --- a/tests/unit/mongoProcessor/MongoProcessorConfigValidator.spec.js +++ b/tests/unit/mongoProcessor/MongoProcessorConfigValidator.spec.js @@ -48,6 +48,21 @@ describe('MongoProcessorConfigValidator mode', () => { assert.strictEqual(validated.mode, 'ingestion'); }); + it('should accept the dr mode', () => { + const validated = configValidator(globalConfig, { ...baseExtConfig, mode: 'dr' }); + assert.strictEqual(validated.mode, 'dr'); + }); + + it('should accept the mode from the environment', () => { + process.env.EXTENSIONS_MONGO_PROCESSOR_MODE = 'dr'; + try { + const validated = configValidator(globalConfig, { ...baseExtConfig }); + assert.strictEqual(validated.mode, 'dr'); + } finally { + delete process.env.EXTENSIONS_MONGO_PROCESSOR_MODE; + } + }); + it('should reject a mode with no implementation', () => { let err; try { From 39c9f403a71e851dc154b39d84d587e173126c42 Mon Sep 17 00:00:00 2001 From: Thomas Flament Date: Tue, 29 Sep 2026 12:30:00 +0200 Subject: [PATCH 6/8] Replace an object overwritten in place An object with no version of its own is rewritten in place, and so is the master a versioning suspended bucket marks null. Merging such an entry kept the previous object's content, headers, user metadata and archive on the sink, which then described an object that no longer existed. The Kafka Connect sink this replaces wrote the whole document, so the merge was a regression against it. An overwrite moves the modification date, which a tag, retention, legal hold or restore update keeps: an entry whose date differs from the stored one replaces the document, and any other entry still merges, so a tag change keeps a restore this site performed. The content digest cannot tell an overwrite either, a PUT replacing the whole metadata whatever the bytes do. A version is immutable and always merges. Nothing of a replaced document is kept. A restore this site performed describes bytes that are gone; reclaiming them is left to the D/R garbage collection still to come. Issue: BB-811 --- .../PullReplicationMetadataPolicy.js | 32 ++++-- .../ingestion/MongoQueueProcessor.js | 105 ++++++++++++++++++ 2 files changed, 130 insertions(+), 7 deletions(-) diff --git a/extensions/mongoProcessor/metadataPolicy/PullReplicationMetadataPolicy.js b/extensions/mongoProcessor/metadataPolicy/PullReplicationMetadataPolicy.js index d8fd2de09..0b0302091 100644 --- a/extensions/mongoProcessor/metadataPolicy/PullReplicationMetadataPolicy.js +++ b/extensions/mongoProcessor/metadataPolicy/PullReplicationMetadataPolicy.js @@ -33,14 +33,15 @@ class PullReplicationMetadataPolicy extends MetadataPolicy { const versionId = this.targetVersionId(entry); const stored = this._storedObject(versionId, objMD); - if (stored) { + if (stored && this._isMetadataUpdate(entry, stored)) { const content = getContentType(entry, stored); this._mergeStoredMetadata(entry, stored); return { content: content.length !== 0 ? content : ['METADATA'], versionId }; } - // the source-side pipeline shaped everything else this object keeps; - // ACLs are not replicated, so they reset as for an ingested object + // a first write, or an object overwritten in place: the source-side + // pipeline shaped everything else this object keeps; ACLs are not + // replicated, so they reset as for an ingested object entry.setAcl(new ObjectMD().getAcl()); return { content: getContentType(entry), versionId }; } @@ -63,13 +64,30 @@ class PullReplicationMetadataPolicy extends MetadataPolicy { } /** - * A version whose data still lives on the remote site. + * Whether the entry updates the metadata of the object already stored, + * rather than describing one that replaced it. * + * A version is immutable, so an entry for one always updates it, a null + * version included. An object with no version of its own is rewritten in + * place, and so is the master a versioning suspended bucket marks null: + * an overwrite moves the modification date, which a metadata update + * keeps. + * + * Nothing of a replaced document is kept. Placement is the one field this + * site could own, and cannot here: a cold object holds what the source + * pipeline derives from its storage class, and clean room -which pulls + * the data to this site- replicates versioned buckets only. + * + * @param {ObjectQueueEntry} entry - object queue entry object * @param {Object} objMD - metadata fetched from mongo - * @return {boolean} true if the stored version is not localized + * @return {boolean} true if the entry updates the stored object */ - _isNotLocalized(objMD) { - return Boolean(locationsConfig[objMD.dataStoreName]?.isCRR); + _isMetadataUpdate(entry, objMD) { + if (this.targetVersionId(entry)) { + return true; + } + + return entry.getLastModified() === objMD['last-modified']; } /** diff --git a/tests/functional/ingestion/MongoQueueProcessor.js b/tests/functional/ingestion/MongoQueueProcessor.js index 33afae4d9..84b045187 100644 --- a/tests/functional/ingestion/MongoQueueProcessor.js +++ b/tests/functional/ingestion/MongoQueueProcessor.js @@ -1307,6 +1307,8 @@ describe('MongoQueueProcessor in dr mode', function drMode() { // a cold object with no version of its own, restored on this site const STORED_ARCHIVE = { archiveInfo: { archiveId: 'stored-archive', archiveVersion: 1 } }; const STORED_LAST_MODIFIED = '2026-09-01T10:00:00.000Z'; + const REWRITE_ARCHIVE = { archiveInfo: { archiveId: 'rewrite-archive', archiveVersion: 2 } }; + const REWRITE_LAST_MODIFIED = '2026-09-02T10:00:00.000Z'; it('should write over a master that is a copy of a version', done => { // versioning suspended on the source: a new null object is put over @@ -1356,6 +1358,109 @@ describe('MongoQueueProcessor in dr mode', function drMode() { ._data)); } + it('should take every field of an object rewritten in place', done => { + storeUnversioned(); + const objmd = new ObjectMD() + .setKey(KEY) + .setContentMd5('9e107d9d372bb6826bd81d3542a419d6') + .setContentType('application/json') + .setUserMetadata({ 'x-amz-meta-colour': 'blue' }) + .setDataStoreName(LOCATION) + .setArchive(REWRITE_ARCHIVE) + .setLastModified(REWRITE_LAST_MODIFIED); + const entry = new ObjectQueueEntry(BUCKET, KEY, objmd); + + processEntry(entry, err => { + assert.ifError(err); + + const added = mqp.getAdded(); + assert.strictEqual(added.length, 1); + // written in place, under the object key alone + assert.strictEqual(added[0].key, KEY); + const { objVal } = added[0]; + assert.strictEqual(objVal['content-md5'], '9e107d9d372bb6826bd81d3542a419d6'); + assert.strictEqual(objVal['content-type'], 'application/json'); + assert.strictEqual(objVal['x-amz-meta-colour'], 'blue'); + assert.deepStrictEqual(objVal.archive, REWRITE_ARCHIVE); + // a restore this site holds describes bytes that are gone + assert.strictEqual(objVal['x-amz-restore'], undefined); + assert.deepStrictEqual(objVal.acl, new ObjectMD().getAcl()); + done(); + }); + }); + + it('should take a rewrite that kept the same bytes', done => { + storeUnversioned(); + // a PUT replaces the whole metadata whatever the content does, so the + // content md5 cannot tell a rewrite from an untouched object + const objmd = new ObjectMD() + .setKey(KEY) + .setContentMd5('7d793037a0760186574b0282f2f435e7') + .setContentType('application/json') + .setUserMetadata({ 'x-amz-meta-colour': 'red' }) + .setDataStoreName(LOCATION) + .setArchive(REWRITE_ARCHIVE) + .setLastModified(REWRITE_LAST_MODIFIED); + const entry = new ObjectQueueEntry(BUCKET, KEY, objmd); + + processEntry(entry, err => { + assert.ifError(err); + + const added = mqp.getAdded(); + assert.strictEqual(added.length, 1); + assert.strictEqual(added[0].objVal['content-type'], 'application/json'); + done(); + }); + }); + + it('should take a rewrite stamped with the archive of the object it replaced', done => { + storeUnversioned(); + // the archival of the overwritten object completed after the + // overwrite, and stamped its archive onto the new object + const objmd = new ObjectMD() + .setKey(KEY) + .setContentMd5('9e107d9d372bb6826bd81d3542a419d6') + .setContentType('application/json') + .setDataStoreName(LOCATION) + .setArchive(STORED_ARCHIVE) + .setLastModified(REWRITE_LAST_MODIFIED); + const entry = new ObjectQueueEntry(BUCKET, KEY, objmd); + + processEntry(entry, err => { + assert.ifError(err); + + const added = mqp.getAdded(); + assert.strictEqual(added.length, 1); + assert.strictEqual(added[0].objVal['content-type'], 'application/json'); + done(); + }); + }); + + it('should take the null version a suspended bucket rewrites', done => { + storeUnversioned(); + // versioning suspended: the master carries a version id and is marked + // null, and it is still the document that gets rewritten in place + const objmd = new ObjectMD() + .setKey(KEY) + .setVersionId(VERSION_ID) + .setIsNull(true) + .setContentMd5('9e107d9d372bb6826bd81d3542a419d6') + .setDataStoreName(LOCATION) + .setArchive(REWRITE_ARCHIVE) + .setLastModified(REWRITE_LAST_MODIFIED); + const entry = new ObjectQueueEntry(BUCKET, KEY, objmd); + + processEntry(entry, err => { + assert.ifError(err); + + const added = mqp.getAdded(); + assert.strictEqual(added.length, 1); + assert.strictEqual(added[0].key, KEY); + assert.strictEqual(added[0].objVal['content-md5'], '9e107d9d372bb6826bd81d3542a419d6'); + done(); + }); + }); + it('should merge an update to an object with no version of its own', done => { storeUnversioned(); // the pipeline strips the restore state, which is this site's own From 33616309f8279e3ebc1051ee0dab56abe378eb41 Mon Sep 17 00:00:00 2001 From: Thomas Flament Date: Tue, 29 Sep 2026 23:38:39 +0200 Subject: [PATCH 7/8] Accept a configuration that runs no queue populator Every backbeat process reads its mongodb client from the queuePopulator section, so a process that only needs that client, like the mongo-processor of a D/R sink, had to fill in a cron rule, a zookeeper path, a probe server and a log source for a populator that never runs. The 'none' log source says so: it needs none of those settings, and a queue populator started with it refuses to run. Issue: BB-811 --- lib/config.joi.js | 13 ++++++++----- lib/queuePopulator/QueuePopulator.js | 6 +++++- tests/unit/QueuePopulator.spec.js | 5 +++++ tests/unit/lib/config/config.joi.spec.js | 21 ++++++++++++++++++++- 4 files changed, 38 insertions(+), 7 deletions(-) diff --git a/lib/config.joi.js b/lib/config.joi.js index b67937c2e..132084954 100644 --- a/lib/config.joi.js +++ b/lib/config.joi.js @@ -39,7 +39,10 @@ const KAFKA_CONSUMER_PARAMS_SCHEMA = joi.object({ 'auto.offset.reset': joi.forbidden(), 'enable.auto.commit': joi.forbidden(), }).unknown(true).default({}); -const logSourcesJoi = joi.string().valid('bucketd', 'ingestion', 'dmd', 'kafka'); +// a process that only needs the mongodb connection sets 'none': it runs no +// queue populator +const logSourcesJoi = joi.string().valid('bucketd', 'ingestion', 'dmd', 'kafka', 'none'); +const optionalWhenNoLogSource = schema => schema.when('logSource', { is: 'none', then: joi.optional() }); const joiSchema = joi.object({ replicationGroupId: joi.string().length(7).default('RG00001'), @@ -68,10 +71,10 @@ const joiSchema = joi.object({ vaultAdmin: hostPortJoi, queuePopulator: { auth: authJoi, - cronRule: joi.string().required(), + cronRule: optionalWhenNoLogSource(joi.string().required()), batchMaxRead: joi.number().default(10000), batchTimeoutMs: joi.number().default(9000), - zookeeperPath: joi.string().required(), + zookeeperPath: optionalWhenNoLogSource(joi.string().required()), logSource: joi.alternatives().try(logSourcesJoi).required(), exhaustLogSource: joi.bool().default(false), @@ -85,11 +88,11 @@ const joiSchema = joi.object({ kafka: qpKafkaJoi.when('logSource', { is: 'kafka', then: joi.required() }), // TODO: BB-625 reset to being required after supporting probeserver in S3C // for bucket notification proceses - probeServer: probeServerJoi.when('...extensions', { + probeServer: optionalWhenNoLogSource(probeServerJoi.when('...extensions', { is: joi.object().keys({ notification: joi.exist() }), then: joi.optional(), otherwise: joi.required(), - }), + })), circuitBreaker: joi.object().optional(), }, log: logJoi, diff --git a/lib/queuePopulator/QueuePopulator.js b/lib/queuePopulator/QueuePopulator.js index 7a5f29365..2910f6190 100644 --- a/lib/queuePopulator/QueuePopulator.js +++ b/lib/queuePopulator/QueuePopulator.js @@ -1,3 +1,4 @@ +const assert = require('assert'); const async = require('async'); const Logger = require('werelogs').Logger; const { BreakerState, CircuitBreaker } = require('@scality/breakbeat').CircuitBreaker; @@ -136,7 +137,8 @@ class QueuePopulator { * @param {String} qpConfig.zookeeperPath - sub-path to use for * storing populator state in zookeeper * @param {String} qpConfig.logSource - type of source - * log: "bucketd" (raft log) or "dmd" (bucketfile) + * log: "bucketd" (raft log) or "dmd" (bucketfile); "none" is + * refused * @param {Object} [qpConfig.bucketd] - bucketd source * configuration (mandatory if logSource is "bucket") * @param {Object} [qpConfig.dmd] - dmd source @@ -158,6 +160,8 @@ class QueuePopulator { */ constructor(zkConfig, kafkaConfig, qpConfig, httpsConfig, mConfig, rConfig, vConfig, extConfigs) { + assert.notStrictEqual(qpConfig.logSource, 'none', + "no queue populator runs with the 'none' log source"); this.zkConfig = zkConfig; this.kafkaConfig = kafkaConfig; this.qpConfig = qpConfig; diff --git a/tests/unit/QueuePopulator.spec.js b/tests/unit/QueuePopulator.spec.js index 0c77a483d..85891591b 100644 --- a/tests/unit/QueuePopulator.spec.js +++ b/tests/unit/QueuePopulator.spec.js @@ -378,6 +378,11 @@ describe('QueuePopulator', () => { }); }); + it('should refuse to run with no log source', () => { + assert.throws(() => new QueuePopulator({}, {}, { logSource: 'none' }, + null, null, null, null, {}), /'none' log source/); + }); + describe('configured extensions', () => { const extConfigs = { replication: {}, lifecycle: {}, notification: {} }; diff --git a/tests/unit/lib/config/config.joi.spec.js b/tests/unit/lib/config/config.joi.spec.js index 61bab3a2d..c993ce070 100644 --- a/tests/unit/lib/config/config.joi.spec.js +++ b/tests/unit/lib/config/config.joi.spec.js @@ -48,7 +48,26 @@ describe('backbeat config schema', () => { it('should reject a log source it cannot read', () => { assert.match(validate({ logSource: 'mongo' }).error.message, - /"queuePopulator.logSource" must be one of \[bucketd, ingestion, dmd, kafka\]/); + /"queuePopulator.logSource" must be one of \[bucketd, ingestion, dmd, kafka, none\]/); + }); + + it('should require the populator settings of any log source but none', () => { + const { error } = schema.validate( + { queuePopulator: { logSource: 'kafka', kafka: {} }, extensions: { gc: {} } }, + { abortEarly: false }); + const missing = error.details.map(detail => detail.path.join('.')); + assert.deepStrictEqual(missing.sort(), [ + 'queuePopulator.cronRule', + 'queuePopulator.probeServer', + 'queuePopulator.zookeeperPath', + ]); + }); + + // the mongo-processor of a D/R sink reads the mongodb connection only + it('should need nothing else with no log source', () => { + const mongo = { replicaSetHosts: 'mongo:27017', database: 'metadata' }; + const config = { queuePopulator: { logSource: 'none', mongo }, extensions: { gc: {} } }; + assert.strictEqual(schema.validate(config).error, undefined); }); }); From d558470826575aa4a6c2545503a5882b63811e77 Mon Sep 17 00:00:00 2001 From: Francois Ferrand Date: Wed, 30 Sep 2026 19:25:02 +0200 Subject: [PATCH 8/8] pull replication: target a null master by its versionId A null master is a version like any other: the source gives it a document of its own on its first metadata update, and streams every later change under that version id. Requesting and writing it as the master would leave those changes on another document, so the sink targets a null master by the internal version id it carries, and only an object with no version id at all is still written in place, told from an overwrite by its modification date. That also leaves one place to decide whether an entry updates the stored document: a targeted version always does, and the guard for a master copying a version has no case left to cover. Issue: BB-811 --- .../PullReplicationMetadataPolicy.js | 56 ++++++----------- .../ingestion/MongoQueueProcessor.js | 62 +++++++++++-------- 2 files changed, 55 insertions(+), 63 deletions(-) diff --git a/extensions/mongoProcessor/metadataPolicy/PullReplicationMetadataPolicy.js b/extensions/mongoProcessor/metadataPolicy/PullReplicationMetadataPolicy.js index 0b0302091..55ebf8161 100644 --- a/extensions/mongoProcessor/metadataPolicy/PullReplicationMetadataPolicy.js +++ b/extensions/mongoProcessor/metadataPolicy/PullReplicationMetadataPolicy.js @@ -17,25 +17,21 @@ class PullReplicationMetadataPolicy extends MetadataPolicy { } targetVersionId(entry) { - // the document the entry comes from: a version by its key, a master - // as the master document itself - const documentVersionId = extractVersionId(entry.getObjectVersionedKey()); - if (documentVersionId) { - return documentVersionId; + if (entry instanceof DeleteOpQueueEntry) { + return extractVersionId(entry.getObjectVersionedKey()); } - if (entry instanceof DeleteOpQueueEntry || entry.getIsNull()) { - return undefined; - } - return entry.getVersionId(); + + // the internal version id: getVersionId() names a null version + // 'null', as the S3 API does, while mongo keys it by its own id + return entry.getValue().versionId; } apply(entry, objMD) { const versionId = this.targetVersionId(entry); - const stored = this._storedObject(versionId, objMD); - if (stored && this._isMetadataUpdate(entry, stored)) { - const content = getContentType(entry, stored); - this._mergeStoredMetadata(entry, stored); + if (objMD && this._isMetadataUpdate(entry, objMD, versionId)) { + const content = getContentType(entry, objMD); + this._mergeStoredMetadata(entry, objMD); return { content: content.length !== 0 ? content : ['METADATA'], versionId }; } @@ -46,32 +42,17 @@ class PullReplicationMetadataPolicy extends MetadataPolicy { return { content: getContentType(entry), versionId }; } - /** - * The stored document, when it is the object the entry targets. The - * master document is only that object when it has no version of its own: - * a master that is a copy of a version, or the latest version mongo - * returns in place of a missing master, is not. - * - * @param {string|undefined} versionId - the version the entry targets - * @param {Object|undefined} objMD - metadata fetched from mongo - * @return {Object|undefined} the stored object, if any - */ - _storedObject(versionId, objMD) { - if (versionId === undefined && objMD?.versionId && !objMD.isNull) { - return undefined; - } - return objMD; - } - /** * Whether the entry updates the metadata of the object already stored, * rather than describing one that replaced it. * - * A version is immutable, so an entry for one always updates it, a null - * version included. An object with no version of its own is rewritten in - * place, and so is the master a versioning suspended bucket marks null: - * an overwrite moves the modification date, which a metadata update - * keeps. + * A version is immutable, so an entry for one always updates it. An + * object with no version of its own is rewritten in place under the same + * key, and so is the master a versioning suspended bucket marks null: an + * overwrite moves the modification date, which a metadata update keeps. + * A master that copies a real version, or the latest version mongo + * returns in place of a missing master, carries that version's date and + * is replaced the same way. * * Nothing of a replaced document is kept. Placement is the one field this * site could own, and cannot here: a cold object holds what the source @@ -80,10 +61,11 @@ class PullReplicationMetadataPolicy extends MetadataPolicy { * * @param {ObjectQueueEntry} entry - object queue entry object * @param {Object} objMD - metadata fetched from mongo + * @param {string|undefined} versionId - the version the entry targets * @return {boolean} true if the entry updates the stored object */ - _isMetadataUpdate(entry, objMD) { - if (this.targetVersionId(entry)) { + _isMetadataUpdate(entry, objMD, versionId) { + if (versionId) { return true; } diff --git a/tests/functional/ingestion/MongoQueueProcessor.js b/tests/functional/ingestion/MongoQueueProcessor.js index 84b045187..84f622336 100644 --- a/tests/functional/ingestion/MongoQueueProcessor.js +++ b/tests/functional/ingestion/MongoQueueProcessor.js @@ -1310,36 +1310,45 @@ describe('MongoQueueProcessor in dr mode', function drMode() { const REWRITE_ARCHIVE = { archiveInfo: { archiveId: 'rewrite-archive', archiveVersion: 2 } }; const REWRITE_LAST_MODIFIED = '2026-09-02T10:00:00.000Z'; - it('should write over a master that is a copy of a version', done => { - // versioning suspended on the source: a new null object is put over - // the master, which still copies a real version here - sinon.stub(mongoClient, 'getObject').callsFake((b, k, p, l, cb) => - cb(null, new ObjectMD() - .setKey(KEY) - .setVersionId(VERSION_ID) - .setTags({ stored: 'tag' }) - .setDataStoreName(LOCATION) - .setLastModified(STORED_LAST_MODIFIED) - ._data)); - // the same date, so that only the version tells the two apart - const objmd = new ObjectMD() + it('should apply a null master and then its own version as one document', done => { + // versioning suspended on the source: the null master is written as a + // version of its own, so when the source later gives it a document of + // its own -- the same version -- that entry updates what the master + // created, and the second write is the first over again + const versionKey = `${KEY}${VID_SEP}${VERSION_ID}`; + const nullVersion = () => new ObjectMD() .setKey(KEY) - .setVersionId(NEW_VERSION_ID) + .setVersionId(VERSION_ID) .setIsNull(true) - .setTags({ entry: 'tag' }) + .setTags({ mytag: 'mytags-value' }) .setDataStoreName(LOCATION) .setLastModified(STORED_LAST_MODIFIED); - const entry = new ObjectQueueEntry(BUCKET, KEY, objmd); + const master = new ObjectQueueEntry(BUCKET, KEY, nullVersion()); + const version = new ObjectQueueEntry(BUCKET, versionKey, nullVersion()); + // what the first write left behind, as mongo gives it back + const getObject = sinon.stub(mongoClient, 'getObject').callsFake((b, k, p, l, cb) => { + const stored = mqp.getAdded().find(a => a.key === `${KEY}${VID_SEP}${p.versionId}`); + return stored ? + cb(null, JSON.parse(JSON.stringify(stored.objVal))) : + cb(errors.NoSuchKey); + }); - processEntry(entry, err => { + async.series([ + next => processEntry(master, next), + next => processEntry(version, next), + ], err => { assert.ifError(err); + // the null version is requested both times, never the master + getObject.getCalls().forEach(call => + assert.strictEqual(call.args[2].versionId, VERSION_ID)); const added = mqp.getAdded(); - assert.strictEqual(added.length, 1); - // written in place, under the object key alone - assert.strictEqual(added[0].key, KEY); - assert.strictEqual(added[0].objVal.versionId, NEW_VERSION_ID); - assert.deepStrictEqual(added[0].objVal.tags, { entry: 'tag' }); + assert.strictEqual(added.length, 2); + assert.strictEqual(added[0].key, versionKey); + assert.strictEqual(added[1].key, versionKey); + assert.deepStrictEqual( + JSON.parse(JSON.stringify(added[1].objVal)), + JSON.parse(JSON.stringify(added[0].objVal))); done(); }); }); @@ -1437,9 +1446,10 @@ describe('MongoQueueProcessor in dr mode', function drMode() { }); it('should take the null version a suspended bucket rewrites', done => { - storeUnversioned(); - // versioning suspended: the master carries a version id and is marked - // null, and it is still the document that gets rewritten in place + // versioning suspended: the put renews the null master's version id, + // so nothing is stored under it yet, and the entry is written as a + // version of its own, the master following it + sinon.stub(mongoClient, 'getObject').yields(errors.NoSuchKey); const objmd = new ObjectMD() .setKey(KEY) .setVersionId(VERSION_ID) @@ -1455,7 +1465,7 @@ describe('MongoQueueProcessor in dr mode', function drMode() { const added = mqp.getAdded(); assert.strictEqual(added.length, 1); - assert.strictEqual(added[0].key, KEY); + assert.strictEqual(added[0].key, `${KEY}${VID_SEP}${VERSION_ID}`); assert.strictEqual(added[0].objVal['content-md5'], '9e107d9d372bb6826bd81d3542a419d6'); done(); });