Skip to content
Open
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
187 changes: 93 additions & 94 deletions extensions/replication/tasks/MultipleBackendTask.js

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Migrations here aren't the main point of this pr, but were kinda forced as MultipleBackend extends ReplicateObject.

I did the minimum required migration on this file

Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
const { promisify } = require('util');
const async = require('async');
const { v4: uuid } = require('uuid');
const { GetBucketReplicationCommand } = require('@aws-sdk/client-s3');
Expand Down Expand Up @@ -54,7 +55,7 @@ class MultipleBackendTask extends ReplicateObject {
return this.destConfig.replicationEndpoint.type;
}

_setupRolesOnce(entry, log, cb) {
async _setupRolesOnce(entry, log) {
log.debug('getting bucket replication', { entry: entry.getLogInfo() });
const entryRolesString = entry.getReplicationRoles(entry.getReplicationBackend());
let errMessage;
Expand All @@ -71,7 +72,7 @@ class MultipleBackendTask extends ReplicateObject {
entry: entry.getLogInfo(),
roles: entryRolesString,
});
return cb(errors.BadRole.customizeDescription(errMessage));
throw errors.BadRole.customizeDescription(errMessage);
}
this.sourceRole = entryRoles[0];

Expand All @@ -81,89 +82,89 @@ class MultipleBackendTask extends ReplicateObject {
Bucket: entry.getBucket(),
});
attachReqUids(command, log.getSerializedUids());
return this.S3source.send(command)
.then(data => {
const replicationEnabled = data.ReplicationConfiguration.Rules
.some(rule => rule.Status === 'Enabled' &&
entry.getObjectKey().startsWith(
rule.Filter?.Prefix ?? rule.Prefix ?? ''));
if (!replicationEnabled) {
errMessage = 'replication disabled for object';
log.debug(errMessage, {
method: 'MultipleBackendTask._setupRolesOnce',
entry: entry.getLogInfo(),
});
return cb(errors.PreconditionFailed.customizeDescription(
errMessage));
}
const roles = data.ReplicationConfiguration.Role.split(',');
if (roles.length > 2) {
errMessage = 'expecting no more than two roles in bucket ' +
'replication configuration when replicating to an ' +
'external location';
log.error(errMessage, {
method: 'MultipleBackendTask._setupRolesOnce',
entry: entry.getLogInfo(),
roles,
});
return cb(errors.BadRole.customizeDescription(errMessage));
}
if (roles[0] !== entryRoles[0]) {
log.error('role in replication entry for source does not ' +
'match role in bucket replication configuration', {
method: 'MultipleBackendTask._setupRolesOnce',
entry: entry.getLogInfo(),
entryRole: entryRoles[0],
bucketRole: roles[0],
});
return cb(errors.BadRole);
}
return cb();
})
.catch(err => {
log.error('error getting replication configuration from S3', {
method: 'MultipleBackendTask._setupRolesOnce',
entry: entry.getLogInfo(),
origin: 'source',
peer: this.sourceConfig.s3,
error: err.message,
httpStatus: err.$metadata?.httpStatusCode,
});
// eslint-disable-next-line no-param-reassign
err.origin = 'source';
return cb(err);

let data;
try {
data = await this.S3source.send(command);
} catch (err) {
log.error('error getting replication configuration from S3', {
method: 'MultipleBackendTask._setupRolesOnce',
entry: entry.getLogInfo(),
origin: 'source',
peer: this.sourceConfig.s3,
error: err.message,
httpStatus: err.$metadata?.httpStatusCode,
});
err.origin = 'source';
throw err;
}

const replicationEnabled = data.ReplicationConfiguration.Rules
.some(rule => rule.Status === 'Enabled' &&
entry.getObjectKey().startsWith(
rule.Filter?.Prefix ?? rule.Prefix ?? ''));
if (!replicationEnabled) {
errMessage = 'replication disabled for object';
log.debug(errMessage, {
method: 'MultipleBackendTask._setupRolesOnce',
entry: entry.getLogInfo(),
});
throw errors.PreconditionFailed.customizeDescription(errMessage);
}
const roles = data.ReplicationConfiguration.Role.split(',');
if (roles.length > 2) {
errMessage = 'expecting no more than two roles in bucket ' +
'replication configuration when replicating to an ' +
'external location';
log.error(errMessage, {
method: 'MultipleBackendTask._setupRolesOnce',
entry: entry.getLogInfo(),
roles,
});
throw errors.BadRole.customizeDescription(errMessage);
}
if (roles[0] !== entryRoles[0]) {
log.error('role in replication entry for source does not ' +
'match role in bucket replication configuration', {
method: 'MultipleBackendTask._setupRolesOnce',
entry: entry.getLogInfo(),
entryRole: entryRoles[0],
bucketRole: roles[0],
});
throw errors.BadRole;
}
return [this.sourceRole, entryRoles[1]];

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

entryRoles[1] is undefined for the single-role external-backend case , the value is discarded by _setupClients, but returning a knowingly-undefined slot is confusing; a comment or returning only what's meaningful would help.

}
Comment thread
SylvainSenechal marked this conversation as resolved.

_refreshSourceEntry(sourceEntry, log, cb) {
async _refreshSourceEntry(sourceEntry, log) {
const params = {
bucket: sourceEntry.getBucket(),
objectKey: sourceEntry.getObjectKey(),
versionId: sourceEntry.getEncodedVersionId() || 'null',
};
return this.backbeatSourceProxy.getMetadata(
params, log, (err, blob) => {
if (err) {
log.error('error getting metadata blob from S3', {
method: 'MultipleBackendTask._refreshSourceEntry',
error: err,
});
return cb(err);
}
const parsedEntry = ObjectQueueEntry.createFromBlob(blob.Body);
if (parsedEntry.error) {
log.error('error parsing metadata blob', {
error: parsedEntry.error,
method: 'MultipleBackendTask._refreshSourceEntry',
});
return cb(errors.InternalError.
customizeDescription('error parsing metadata blob'));
}
const refreshedEntry = new ObjectQueueEntry(sourceEntry.getBucket(),
sourceEntry.getObjectVersionedKey(), parsedEntry.result)
.setReplicationBackend(sourceEntry.getReplicationBackend());
return cb(null, refreshedEntry);
});
const getMetadata = promisify(
this.backbeatSourceProxy.getMetadata.bind(this.backbeatSourceProxy));
let blob;
try {
blob = await getMetadata(params, log);
} catch (err) {
log.error('error getting metadata blob from S3', {
method: 'MultipleBackendTask._refreshSourceEntry',
error: err,
});
throw err;
}
const parsedEntry = ObjectQueueEntry.createFromBlob(blob.Body);
if (parsedEntry.error) {
log.error('error parsing metadata blob', {
error: parsedEntry.error,
method: 'MultipleBackendTask._refreshSourceEntry',
});
throw errors.InternalError.customizeDescription('error parsing metadata blob');
}
return new ObjectQueueEntry(sourceEntry.getBucket(),
sourceEntry.getObjectVersionedKey(), parsedEntry.result)
.setReplicationBackend(sourceEntry.getReplicationBackend());
}

/**
Expand Down Expand Up @@ -1186,7 +1187,7 @@ class MultipleBackendTask extends ReplicateObject {
_setupClients(entry, log, cb) {
// Sets up source clients using the role from the replication
// configuration if the authentication type is as such.
return this._setupRoles(entry, log, cb);
return this._setupRoles(entry, log).then(() => cb(), cb);
}

processQueueEntry(sourceEntry, kafkaEntry, done) {
Expand All @@ -1196,18 +1197,16 @@ class MultipleBackendTask extends ReplicateObject {

return async.waterfall([
next => this._setupClients(sourceEntry, log, next),
next => this._refreshSourceEntry(sourceEntry, log, (err, res) => {
if (err && err.name === 'ObjNotFound' &&
sourceEntry.getReplicationIsNFS() && !sourceEntry.getIsDeleteMarker()) {
next => this._refreshSourceEntry(sourceEntry, log)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This will runs the next waterfall stage synchronously, and that stage is not async — metricsHandler.rpo() can all throw synchronously for example. A throw there lands in .catch, which calls next again; async.waterfall's onlyOnce then throws "Callback was already called" as an uncaught exception in the queue processor, masking the real error. Use the two-argument then so the success path can't feed the rejection handler

next => this._refreshSourceEntry(sourceEntry, log)
                .then(res => next(null, res), err => {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

aiya its more complicated than i thought to make these refactoring

.then(res => next(null, res), err => {
if (err.name === 'ObjNotFound' &&
sourceEntry.getReplicationIsNFS() && !sourceEntry.getIsDeleteMarker()) {
// The object was deleted before entry is processed, we
// can safely skip this entry.
return next(errors.InvalidObjectState);
}
if (err) {
}
return next(err);
}
return next(null, res);
}),
}),
(refreshedEntry, next) => {
const lastModified = new Date(refreshedEntry.getLastModified());
this.metricsHandler.rpo({
Expand Down Expand Up @@ -1267,18 +1266,18 @@ class MultipleBackendTask extends ReplicateObject {
}
return this._getAndPutObject(sourceEntry, log, next);
},
], err => this._handleReplicationOutcome(
err, sourceEntry, kafkaEntry, log, done));
], err => this._handleReplicationOutcome(err, sourceEntry, null, kafkaEntry, log)
.then(result => (result === null ? done() : done(null, result)), done));
}

_handleReplicationOutcome(err, sourceEntry, kafkaEntry, log, done) {
async _handleReplicationOutcome(err, sourceEntry, destEntry, kafkaEntry, log) {
if (!err) {
log.debug('replication succeeded for object, publishing ' +
'replication status as COMPLETED',
{ entry: sourceEntry.getLogInfo() });
this._publishReplicationStatus(
sourceEntry, 'COMPLETED', { kafkaEntry, log });
return done(null, { committable: false });
return { committable: false };
}
if (err.BadRole || err.name === 'BadRole' ||
(err.origin === 'source' &&
Expand All @@ -1291,18 +1290,18 @@ class MultipleBackendTask extends ReplicateObject {
entry: sourceEntry.getLogInfo(),
origin: err.origin,
error: err.description });
return done();
return null;
}
if (err.ObjNotFound || err.name === 'ObjNotFound') {
log.info('replication skipped: ' +
'source object version does not exist',
{ entry: sourceEntry.getLogInfo() });
return done();
return null;
}
if (err.InvalidObjectState || err.name === 'InvalidObjectState') {
log.info('replication skipped: invalid object state',
{ entry: sourceEntry.getLogInfo() });
return done();
return null;
}
log.debug('replication failed permanently for object, ' +
'publishing replication status as FAILED',
Expand All @@ -1315,7 +1314,7 @@ class MultipleBackendTask extends ReplicateObject {
reason: err.description,
kafkaEntry,
});
return done(null, { committable: false });
return { committable: false };
}
}

Expand Down
Loading
Loading