Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions extensions/mongoProcessor/MongoProcessorConfigValidator.js
Original file line number Diff line number Diff line change
@@ -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');

Expand All @@ -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),
Expand Down
165 changes: 26 additions & 139 deletions extensions/mongoProcessor/MongoQueueProcessor.js
Original file line number Diff line number Diff line change
Expand Up @@ -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');
Expand All @@ -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');

Expand Down Expand Up @@ -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);

Expand Down Expand Up @@ -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',
Expand All @@ -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.
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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',
Expand All @@ -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',
Expand All @@ -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.
Expand All @@ -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;
}

Expand Down
Loading
Loading