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
29 changes: 12 additions & 17 deletions lib/models/ObjectQueueEntry.js
Original file line number Diff line number Diff line change
Expand Up @@ -212,39 +212,34 @@ class ObjectQueueEntry extends ObjectMD {
};
}

_getGlobalReplicationStatus() {
const data = this.getValue();
// Check the global status relative to the other backends
if (Array.isArray(data.replicationInfo.backends)) {
const statuses = data.replicationInfo.backends.map(
backend => backend.status);
// If any site replication failed, set the global status
// to FAILED.
if (statuses.includes('FAILED')) {
return 'FAILED';
_setGlobalReplicationStatus() {
Comment thread
francoisferrand marked this conversation as resolved.
let status = 'COMPLETED';
for (const backend of this.getReplicationBackends() ?? []) {
if (backend.status === 'FAILED') {
return this.setReplicationStatus('FAILED');
}
if (statuses.includes('PENDING')) {
return 'PROCESSING';
if (backend.status === 'PENDING') {
status = 'PROCESSING';
}
}
return 'COMPLETED';
return this.setReplicationStatus(status);
}

toReplicaEntry(backend) {
Comment thread
maeldonn marked this conversation as resolved.
const newEntry = this.clone();
newEntry
return newEntry
.setAccountId(this.getAccountId())
.setBucket(this.getReplicationTargetBucket(backend))
.setReplicationBackends([newEntry._findBackend(backend)])
Comment thread
maeldonn marked this conversation as resolved.
.setReplicationSiteStatus(backend, 'REPLICA')
.setReplicationStatus('REPLICA');
return newEntry;
}

toCompletedEntry(backend) {
return this.clone()
.setAccountId(this.getAccountId())
.setReplicationSiteStatus(backend, 'COMPLETED')
.setReplicationStatus(this._getGlobalReplicationStatus())
._setGlobalReplicationStatus()
.setOriginOp('s3:Replication:OperationCompletedReplication');
}

Expand All @@ -260,7 +255,7 @@ class ObjectQueueEntry extends ObjectMD {
return this.clone()
.setAccountId(this.getAccountId())
.setReplicationSiteStatus(backend, 'PENDING')
.setReplicationStatus(this._getGlobalReplicationStatus())
._setGlobalReplicationStatus()
.setOriginOp('s3:Replication:OperationPendingReplication');
}

Expand Down
2 changes: 1 addition & 1 deletion package.json
Original file line number Diff line number Diff line change
Expand Up @@ -59,7 +59,7 @@
"@scality/cloudserverclient": "^1.0.12",
"@smithy/node-http-handler": "^3.3.3",
"JSONStream": "^1.3.5",
"arsenal": "git+https://github.com/scality/arsenal#8.5.6",
"arsenal": "git+https://github.com/scality/arsenal#8.5.18",
"async": "^2.3.0",
"backo": "^1.1.0",
"breakbeat": "scality/breakbeat#v1.0.3",
Expand Down
8 changes: 0 additions & 8 deletions tests/functional/replication/queueProcessor.js
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@
/* eslint-disable max-len */

function waitForSourceMdPut(s3mock) {
return new Promise(resolve => { s3mock.onPutSourceMd = resolve; });

Check warning on line 27 in tests/functional/replication/queueProcessor.js

View workflow job for this annotation

GitHub Actions / lint

Assignment to property of function parameter 's3mock'
}

function getMD5(body) {
Expand Down Expand Up @@ -599,14 +599,6 @@
site: 'sf',
status: 'REPLICA',
dataStoreVersionId: '',
}, {
site: 'replicationaws',
status: 'PENDING',
dataStoreVersionId: '',
}, {
site: 'toazure',
status: 'PENDING',
dataStoreVersionId: '',
}],
content: replicatedContent,
destination: this.getParam('source.md.replicationInfo.destination'),
Expand Down
15 changes: 15 additions & 0 deletions tests/unit/lib/models/ObjectQueueEntry.spec.js
Original file line number Diff line number Diff line change
Expand Up @@ -144,6 +144,21 @@ describe('ObjectQueueEntry', () => {
assert.strictEqual(replicaA.getBucket(), 'bucket-a');
assert.strictEqual(replicaB.getBucket(), 'bucket-b');
});

it('toReplicaEntry drops the other destinations statuses', () => {
const entry = _makeEntryWithBackends([
{ site: 'siteA', status: 'PENDING', dataStoreVersionId: '' },
{ site: 'siteB', status: 'PENDING', dataStoreVersionId: '' },
]);

const replica = entry.toReplicaEntry({ site: 'siteB' });
assert.deepStrictEqual(
replica.getReplicationBackends().map(b => b.site), ['siteB']);
assert.strictEqual(replica.getReplicationSiteStatus({ site: 'siteB' }), 'REPLICA');
assert.strictEqual(replica.getReplicationStatus(), 'REPLICA');
assert.strictEqual(entry.getReplicationBackends().length, 2);
assert.strictEqual(entry.getReplicationSiteStatus({ site: 'siteB' }), 'PENDING');
});
});

describe('same-site backend disambiguation', () => {
Expand Down
17 changes: 9 additions & 8 deletions tests/unit/replication/QueueEntry.spec.js
Original file line number Diff line number Diff line change
Expand Up @@ -39,15 +39,15 @@ describe('QueueEntry helper class', () => {
'REPLICA');
assert.strictEqual(
replica.getReplicationSiteStatus({ site: 'replicationaws' }),
'PENDING');
undefined);
assert.strictEqual(replica.getReplicationStatus(), 'REPLICA');

// If one site is FAILED, the global status should be FAILED
const failed = entry.toFailedEntry({ site: 'sf' });
assert.strictEqual(failed.getReplicationSiteStatus({ site: 'sf' }),
'FAILED');
assert.strictEqual(
replica.getReplicationSiteStatus({ site: 'replicationaws' }),
failed.getReplicationSiteStatus({ site: 'replicationaws' }),
'PENDING');
assert.strictEqual(failed.getReplicationStatus(), 'FAILED');

Expand All @@ -62,14 +62,15 @@ describe('QueueEntry helper class', () => {
assert.strictEqual(completed.getReplicationStatus(), 'PROCESSING');

// If all sites are COMPLETED, the global status should be COMPLETED
const completed1 = entry.toCompletedEntry({ site: 'sf' });
const completed2 = entry.toCompletedEntry({ site: 'replicationaws' });
assert.strictEqual(completed2
.getReplicationSiteStatus({ site: 'replicationaws' }),
const allCompleted = entry.toCompletedEntry({ site: 'sf' })
.toCompletedEntry({ site: 'replicationaws' });
assert.strictEqual(allCompleted.getReplicationSiteStatus({ site: 'sf' }),
'COMPLETED');
assert.strictEqual(completed1.getReplicationSiteStatus({ site: 'sf' }),
assert.strictEqual(allCompleted
.getReplicationSiteStatus({ site: 'replicationaws' }),
'COMPLETED');
assert.strictEqual(completed1.getReplicationStatus(), 'COMPLETED');
assert.strictEqual(allCompleted.getReplicationStatus(), 'COMPLETED');
assert.strictEqual(entry.getReplicationSiteStatus({ site: 'sf' }), 'PENDING');
});
});
});
14 changes: 7 additions & 7 deletions yarn.lock
Original file line number Diff line number Diff line change
Expand Up @@ -5500,9 +5500,9 @@ arraybuffer.prototype.slice@^1.0.4:
optionalDependencies:
ioctl "^2.0.2"

"arsenal@git+https://github.com/scality/arsenal#8.5.6":
version "8.5.6"
resolved "git+https://github.com/scality/arsenal#0db557930c7d13204167188a7503e6e00d154df5"
"arsenal@git+https://github.com/scality/arsenal#8.5.18":
version "8.5.18"
resolved "git+https://github.com/scality/arsenal#e0d5315a1bde8c94e24e1dad9129074333eddbcd"
dependencies:
"@aws-sdk/client-kms" "^3.975.0"
"@aws-sdk/client-s3" "^3.975.0"
Expand Down Expand Up @@ -5541,7 +5541,7 @@ arraybuffer.prototype.slice@^1.0.4:
simple-glob "^0.2.0"
socket.io "^4.8.0"
socket.io-client "^4.8.0"
sproxydclient "github:scality/sproxydclient#8.2.1"
sproxydclient "github:scality/sproxydclient#8.2.2"
utf8 "^3.0.0"
uuid "^10.0.0"
werelogs scality/werelogs#8.2.4
Expand Down Expand Up @@ -10554,9 +10554,9 @@ sprintf-js@~1.0.2:
httpagent "github:scality/httpagent#1.1.0"
werelogs scality/werelogs#8.2.0

"sproxydclient@github:scality/sproxydclient#8.2.1":
version "8.2.1"
resolved "https://codeload.github.com/scality/sproxydclient/tar.gz/501829f5521787e7e946de6792f5806ffa6ec437"
"sproxydclient@github:scality/sproxydclient#8.2.2":
version "8.2.2"
resolved "https://codeload.github.com/scality/sproxydclient/tar.gz/62d3aac1c844beda2286f33a1a7c7f2a31937579"
dependencies:
async "^3.2.6"
httpagent "github:scality/httpagent#1.1.0"
Expand Down
Loading