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 5e797a5c9..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,18 +226,19 @@ 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 }); 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', @@ -251,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. @@ -380,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, @@ -477,29 +389,19 @@ 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); }; 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', @@ -509,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', @@ -521,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. @@ -545,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/PullReplicationMetadataPolicy.js b/extensions/mongoProcessor/metadataPolicy/PullReplicationMetadataPolicy.js new file mode 100644 index 000000000..55ebf8161 --- /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) { + if (entry instanceof DeleteOpQueueEntry) { + return extractVersionId(entry.getObjectVersionedKey()); + } + + // 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); + + if (objMD && this._isMetadataUpdate(entry, objMD, versionId)) { + const content = getContentType(entry, objMD); + this._mergeStoredMetadata(entry, objMD); + return { content: content.length !== 0 ? content : ['METADATA'], versionId }; + } + + // 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 }; + } + + /** + * 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. 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 + * 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 + * @param {string|undefined} versionId - the version the entry targets + * @return {boolean} true if the entry updates the stored object + */ + _isMetadataUpdate(entry, objMD, versionId) { + if (versionId) { + return true; + } + + return entry.getLastModified() === objMD['last-modified']; + } + + /** + * 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 new file mode 100644 index 000000000..0e861a2bf --- /dev/null +++ b/extensions/mongoProcessor/metadataPolicy/index.js @@ -0,0 +1,12 @@ +const IngestionMetadataPolicy = require('./IngestionMetadataPolicy'); +const PullReplicationMetadataPolicy = require('./PullReplicationMetadataPolicy'); + +const metadataPolicies = { + ingestion: IngestionMetadataPolicy, + dr: PullReplicationMetadataPolicy, +}; + +module.exports = { + metadataPolicies, + defaultMode: 'ingestion', +}; 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/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/functional/ingestion/MongoQueueProcessor.js b/tests/functional/ingestion/MongoQueueProcessor.js index ffab98101..84f622336 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'; @@ -300,7 +303,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); @@ -1070,3 +1074,488 @@ 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'; + const REWRITE_ARCHIVE = { archiveInfo: { archiveId: 'rewrite-archive', archiveVersion: 2 } }; + const REWRITE_LAST_MODIFIED = '2026-09-02T10:00:00.000Z'; + + 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(VERSION_ID) + .setIsNull(true) + .setTags({ mytag: 'mytags-value' }) + .setDataStoreName(LOCATION) + .setLastModified(STORED_LAST_MODIFIED); + 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); + }); + + 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, 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(); + }); + }); + + 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 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 => { + // 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) + .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}${VID_SEP}${VERSION_ID}`); + 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 + 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/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.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)); 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); }); }); 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: [], }, 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..ea7d60f11 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,36 @@ 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 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 { + configValidator(globalConfig, { ...baseExtConfig, mode: 'sideways' }); + } catch (e) { + err = e; + } + assert(err, 'expected configValidator to throw on an unknown mode'); + assert.match(err.message, /mode/); + }); +});