mirror of
https://github.com/cryptpad/cryptpad.git
synced 2026-09-12 19:49:59 +05:00
Merge pull request #1800 from cryptpad/blobmd
Blob metadata refactoring
This commit is contained in:
commit
29a84c5114
25
lib/api.js
25
lib/api.js
@ -27,6 +27,31 @@ nThen(function (w) {
|
||||
console.error(err);
|
||||
}
|
||||
}));
|
||||
}).nThen(function () {
|
||||
if (Env.proofsMigrated) { return; }
|
||||
const { Worker } = require('node:worker_threads');
|
||||
const Admin = require("./commands/admin-rpc");
|
||||
|
||||
const worker = new Worker('./scripts/migrations/migrate-blob-proofs.js');
|
||||
|
||||
worker.on('message', message => {
|
||||
if (message === 'READY') {
|
||||
log.info('BLOB_PROOFS_MIGRATION');
|
||||
return void worker.postMessage({
|
||||
start: 1,
|
||||
});
|
||||
}
|
||||
if (message === 'MIGRATED') {
|
||||
return void log.info('BLOB_PROOFS_DELETION');
|
||||
}
|
||||
if (message === 'CLEANED') {
|
||||
log.info('BLOB_PROOFS_MIGRATED');
|
||||
Admin.sendDecree(Env, null, function (err) {
|
||||
if (err) { return void log.error('BLOB_PROOF', err); }
|
||||
Env.flushCache();
|
||||
}, ['PROOFS_MIGRATED', ['PROOFS_MIGRATED', 1]], 'server');
|
||||
}
|
||||
});
|
||||
}).nThen(function (w) {
|
||||
let admins = Env.admins || [];
|
||||
|
||||
|
||||
@ -13,6 +13,7 @@ const Metadata = require("./commands/metadata");
|
||||
const Meta = require("./metadata");
|
||||
const Logger = require("./log");
|
||||
const plugins = require("./plugin-manager");
|
||||
const HK = require('./hk-util');
|
||||
|
||||
let SSOUtils = plugins.SSO && plugins.SSO.utils;
|
||||
|
||||
@ -53,7 +54,13 @@ const init = (cb) => {
|
||||
Env.computeMetadata = function (channel, cb) {
|
||||
const ref = {};
|
||||
const lineHandler = Meta.createLineHandler(ref, (err) => { console.log(err); });
|
||||
return void Env.store.readChannelMetadata(channel, lineHandler, function (err) {
|
||||
|
||||
let f = Env.store.readChannelMetadata;
|
||||
if (channel.length === HK.BLOB_ID_LENGTH) {
|
||||
f = Env.blobStore.readMetadata;
|
||||
}
|
||||
|
||||
return void f(channel, lineHandler, function (err) {
|
||||
if (err) {
|
||||
// stream errors?
|
||||
return void cb(err);
|
||||
@ -132,9 +139,12 @@ COMMANDS.start = (edPublic, blockId, reason) => {
|
||||
n = n((w) => {
|
||||
// Blobs
|
||||
if (Env.blobStore.isFileId(chanId)) {
|
||||
return void Env.blobStore.isOwnedBy(safeKey, chanId, w((err, owned) => {
|
||||
if (err || !owned) { return; }
|
||||
blobsToArchive.push(chanId);
|
||||
return Env.computeMetadata(chanId, w((e, md) => {
|
||||
if (e || !md) { return; }
|
||||
if (md && md.owners
|
||||
&& md.owners.includes(edPublic)) {
|
||||
blobsToArchive.push(chanId);
|
||||
}
|
||||
}));
|
||||
}
|
||||
// Pads
|
||||
|
||||
@ -409,7 +409,7 @@ var getChannelMetadata = function (Env, Server, cb, data) {
|
||||
};
|
||||
|
||||
// CryptPad_AsyncStore.rpc.send('ADMIN', [ 'ADMIN_DECREE', ['RESTRICT_REGISTRATION', [true]]], console.log)
|
||||
var adminDecree = function (Env, Server, cb, data, unsafeKey) {
|
||||
var adminDecree = Admin.sendDecree = function (Env, Server, cb, data, unsafeKey) {
|
||||
var value = data[1];
|
||||
if (!Array.isArray(value)) { return void cb('INVALID_DECREE'); }
|
||||
|
||||
|
||||
@ -13,7 +13,8 @@ Data.getMetadataRaw = function (Env, channel /* channelName */, _cb) {
|
||||
const cb = Util.once(Util.mkAsync(_cb));
|
||||
if (!Core.isValidId(channel)) { return void cb('INVALID_CHAN'); }
|
||||
if (channel.length !== HK.STANDARD_CHANNEL_LENGTH &&
|
||||
channel.length !== HK.ADMIN_CHANNEL_LENGTH) { return cb("INVALID_CHAN_LENGTH"); }
|
||||
channel.length !== HK.ADMIN_CHANNEL_LENGTH &&
|
||||
channel.length !== HK.BLOB_ID_LENGTH) { return cb("INVALID_CHAN_LENGTH"); }
|
||||
|
||||
// return synthetic metadata for admin broadcast channels as a safety net
|
||||
// in case anybody manages to write metadata
|
||||
@ -79,6 +80,7 @@ Data.setMetadata = function (Env, safeKey, data, cb, Server) {
|
||||
|
||||
var channel = data.channel;
|
||||
var command = data.command;
|
||||
// XXX BLOBMD allow blobs
|
||||
if (!channel || !Core.isValidId(channel)) { return void cb ('INVALID_CHAN'); }
|
||||
if (!command || typeof (command) !== 'string') { return void cb('INVALID_COMMAND'); }
|
||||
if (Meta.commands.indexOf(command) === -1) { return void cb('UNSUPPORTED_COMMAND'); }
|
||||
@ -140,6 +142,7 @@ Data.setMetadata = function (Env, safeKey, data, cb, Server) {
|
||||
cb(void 0, metadata);
|
||||
return void next();
|
||||
}
|
||||
// XXX BLOBMD use correct store for blobs
|
||||
Env.msgStore.writeMetadata(channel, JSON.stringify(line), function (e) {
|
||||
if (e) {
|
||||
cb(e);
|
||||
@ -155,6 +158,7 @@ Data.setMetadata = function (Env, safeKey, data, cb, Server) {
|
||||
|
||||
// update the cached metadata
|
||||
metadata_cache[channel] = metadata;
|
||||
Env.checkCache(channel); // XXX ???
|
||||
|
||||
// it's easy to check if the channel is restricted
|
||||
const isRestricted = metadata.restricted;
|
||||
|
||||
@ -395,6 +395,14 @@ commands.ADD_ADMIN_KEY = function (Env, args) {
|
||||
return true;
|
||||
};
|
||||
|
||||
commands.PROOFS_MIGRATED = function (Env, args) {
|
||||
if (args !== 1) {
|
||||
throw new Error("INVALID_ARGS");
|
||||
}
|
||||
Env.proofsMigrated = true;
|
||||
return true;
|
||||
};
|
||||
|
||||
commands.SET_BEARER_SECRET = function (Env, args) {
|
||||
if (!args_isString(args) || args.length !== 1 || !args[0]) {
|
||||
throw new Error("INVALID_ARGS");
|
||||
|
||||
@ -415,6 +415,7 @@ const BAD = [
|
||||
'limits',
|
||||
'customLimits',
|
||||
'scheduleDecree',
|
||||
'plugins',
|
||||
|
||||
'httpServer',
|
||||
|
||||
|
||||
144
lib/eviction.js
144
lib/eviction.js
@ -51,7 +51,6 @@ var evictArchived = function (Env, cb) {
|
||||
var report = {
|
||||
// archivedChannelsRemoved,
|
||||
// archivedAccountsRemoved,
|
||||
// archivedBlobProofsRemoved,
|
||||
// archivedBlobsRemoved,
|
||||
|
||||
// totalChannels,
|
||||
@ -237,37 +236,6 @@ var evictArchived = function (Env, cb) {
|
||||
store.listArchivedChannels(handler, w(done));
|
||||
};
|
||||
|
||||
var removeArchivedBlobProofs = function (w) {
|
||||
if (typeof(Env.archiveRetentionTime) !== "number") { return; }
|
||||
// Iterate over archive blob ownership proofs and remove them
|
||||
// if they are older than the specified retention time
|
||||
var removed = 0;
|
||||
blobs.list.archived.proofs(function (err, item, next) {
|
||||
next = Util.mkAsync(next, THROTTLE_FACTOR);
|
||||
if (err) {
|
||||
Log.error("EVICT_BLOB_LIST_ARCHIVED_PROOF_ERROR", err);
|
||||
return void next();
|
||||
}
|
||||
if (item && item.ctime > retentionTime) { return void next(); }
|
||||
if (Env.DRY_RUN) {
|
||||
removed++;
|
||||
return void Log.info("EVICT_ARCHIVED_BLOB_PROOF_DRY_RUN", item, next);
|
||||
}
|
||||
blobs.remove.archived.proof(item.safeKey, item.blobId, (function (err) {
|
||||
if (err) {
|
||||
Log.error("EVICT_ARCHIVED_BLOB_PROOF_ERROR", item);
|
||||
return void next();
|
||||
}
|
||||
Log.info("EVICT_ARCHIVED_BLOB_PROOF", item);
|
||||
removed++;
|
||||
next();
|
||||
}));
|
||||
}, w(function () {
|
||||
report.archivedBlobProofsRemoved = removed;
|
||||
Log.info('EVICT_ARCHIVED_BLOB_PROOFS_REMOVED', removed);
|
||||
}));
|
||||
};
|
||||
|
||||
var removeArchivedBlobs = function (w) {
|
||||
if (typeof(Env.archiveRetentionTime) !== "number") { return; }
|
||||
// Iterate over archived blobs and remove them
|
||||
@ -303,7 +271,6 @@ var evictArchived = function (Env, cb) {
|
||||
nThen(loadStorage)
|
||||
.nThen(migrateIncorrectBlobs)
|
||||
.nThen(removeArchivedChannels)
|
||||
.nThen(removeArchivedBlobProofs)
|
||||
.nThen(removeArchivedBlobs)
|
||||
.nThen(function () {
|
||||
cb(void 0, report);
|
||||
@ -315,7 +282,6 @@ module.exports = function (Env, cb) {
|
||||
var report = {
|
||||
// archivedChannelsRemoved,
|
||||
// archivedAccountsRemoved,
|
||||
// archivedBlobProofsRemoved,
|
||||
// archivedBlobsRemoved,
|
||||
|
||||
// totalChannels,
|
||||
@ -612,92 +578,45 @@ module.exports = function (Env, cb) {
|
||||
if (pinnedDocs.test(item.blobId)) { return void next(); }
|
||||
if (activeDocs.test(item.blobId)) { return void next(); }
|
||||
|
||||
// This seems redundant because we're already checking the bloom filter
|
||||
// but we can't implement a 'fast mode' for the iterator
|
||||
// unless we address this race condition with this last-minute double-check
|
||||
if (item.mtime > inactiveTime) { return void next(); }
|
||||
|
||||
if (Env.DRY_RUN) {
|
||||
removed++;
|
||||
return void Log.info("EVICT_ARCHIVE_BLOB_DRY_RUN", {
|
||||
item: item,
|
||||
}, next);
|
||||
}
|
||||
blobs.archive.blob(item.blobId, 'INACTIVE', function (err) {
|
||||
if (err) {
|
||||
return Log.error("EVICT_ARCHIVE_BLOB_ERROR", {
|
||||
error: err,
|
||||
// NOTE: fast mode allows us to skip getStats for
|
||||
// the pinned and active channels
|
||||
nThen(function (w) {
|
||||
// double check that the channel really is inactive before archiving it
|
||||
// because it might have been created after the initial activity scan
|
||||
blobs.getStats(item.blobId, w(function (err, newerItem) {
|
||||
if (err) { return; }
|
||||
if (newerItem && getNewestTime(newerItem) > retentionTime) {
|
||||
// it's actually active, so don't archive it.
|
||||
w.abort();
|
||||
cb();
|
||||
}
|
||||
// else fall through to the archival
|
||||
}));
|
||||
}).nThen(function () {
|
||||
if (Env.DRY_RUN) {
|
||||
removed++;
|
||||
return void Log.info("EVICT_ARCHIVE_BLOB_DRY_RUN", {
|
||||
item: item,
|
||||
}, next);
|
||||
}
|
||||
removed++;
|
||||
Log.info("EVICT_ARCHIVE_BLOB", {
|
||||
item: item,
|
||||
}, next);
|
||||
blobs.archive.blob(item.blobId, 'INACTIVE', function (err) {
|
||||
if (err) {
|
||||
return Log.error("EVICT_ARCHIVE_BLOB_ERROR", {
|
||||
error: err,
|
||||
item: item,
|
||||
}, next);
|
||||
}
|
||||
removed++;
|
||||
Log.info("EVICT_ARCHIVE_BLOB", {
|
||||
item: item,
|
||||
}, next);
|
||||
});
|
||||
});
|
||||
}, w(function () {
|
||||
report.totalBlobs = total;
|
||||
report.activeBlobs = total - removed;
|
||||
Log.info('EVICT_BLOBS_REMOVED', removed, w());
|
||||
}));
|
||||
};
|
||||
|
||||
var archiveInactiveBlobProofs = function (w) {
|
||||
// iterate over blob proofs and remove them
|
||||
// if they don't correspond to a pinned or active file
|
||||
var removed = 0;
|
||||
var total = 0;
|
||||
|
||||
Log.info("EVICT_ARCHIVE_INACTIVE_BLOB_PROOFS_START", {});
|
||||
blobs.list.proofs(function (err, item, next) {
|
||||
next = Util.mkAsync(next, THROTTLE_FACTOR);
|
||||
if (err) {
|
||||
return void Log.error("EVICT_BLOB_LIST_PROOFS_ERROR", err, next);
|
||||
}
|
||||
if (!item) {
|
||||
return void Log.error('EVICT_BLOB_LIST_PROOFS_NO_ITEM', item, next);
|
||||
}
|
||||
total++;
|
||||
|
||||
if (total % PROGRESS_FACTOR === 0) {
|
||||
Log.info('EVICT_BLOB_PROOF_PROGRESS', {
|
||||
proofs: total,
|
||||
});
|
||||
}
|
||||
|
||||
if (pinnedDocs.test(item.blobId)) { return void next(); }
|
||||
if (item.mtime > inactiveTime) { return void next(); }
|
||||
nThen(function (w) {
|
||||
blobs.size(item.blobId, w(function (err, size) {
|
||||
if (err && err === 'ENOENT') { return; } // XXX delete the proof
|
||||
if (err) {
|
||||
w.abort();
|
||||
return void Log.error("EVICT_BLOB_LIST_PROOFS_ERROR", err, next);
|
||||
}
|
||||
if (size !== 0) {
|
||||
w.abort();
|
||||
next();
|
||||
}
|
||||
}));
|
||||
}).nThen(function () {
|
||||
if (Env.DRY_RUN) {
|
||||
removed++;
|
||||
return void Log.info("EVICT_BLOB_PROOF_LONELY_DRY_RUN", item, next);
|
||||
}
|
||||
blobs.remove.proof(item.safeKey, item.blobId, function (err) {
|
||||
if (err) {
|
||||
return Log.error("EVICT_BLOB_PROOF_LONELY_ERROR", item, next);
|
||||
}
|
||||
removed++;
|
||||
return Log.info("EVICT_BLOB_PROOF_LONELY", item, next);
|
||||
});
|
||||
});
|
||||
}, w(function () {
|
||||
Log.info("EVICT_BLOB_PROOFS_REMOVED", {
|
||||
removed,
|
||||
total,
|
||||
}, w());
|
||||
}));
|
||||
}), true);
|
||||
};
|
||||
|
||||
var archiveInactiveChannels = function (w) {
|
||||
@ -802,7 +721,6 @@ module.exports = function (Env, cb) {
|
||||
// (documents which are not in either bloom filter)
|
||||
|
||||
.nThen(archiveInactiveBlobs)
|
||||
.nThen(archiveInactiveBlobProofs)
|
||||
.nThen(archiveInactiveChannels)
|
||||
.nThen(function () {
|
||||
var runningTime = report.runningTime = msSinceStart();
|
||||
|
||||
@ -42,6 +42,8 @@ const ADMIN_CHANNEL_LENGTH = HK.ADMIN_CHANNEL_LENGTH = 33;
|
||||
// with a 34 character id
|
||||
const EPHEMERAL_CHANNEL_LENGTH = HK.EPHEMERAL_CHANNEL_LENGTH = 34;
|
||||
|
||||
HK.BLOB_ID_LENGTH = 48;
|
||||
|
||||
// Temporary channels are archived X ms after everyone has left them
|
||||
const TEMPORARY_CHANNEL_LIFETIME = 30 * 1000;
|
||||
|
||||
|
||||
@ -33,6 +33,7 @@ var createLineHandler = Pins.createLineHandler = function (ref, errorHandler) {
|
||||
// it's a weird API but it's faster than unpinning manually
|
||||
var pins = ref.pins = {};
|
||||
ref.index = 0;
|
||||
ref.first = 0;
|
||||
ref.latest = 0; // the latest message (timestamp in ms)
|
||||
ref.surplus = 0; // how many lines exist behind a reset
|
||||
|
||||
@ -58,7 +59,7 @@ var createLineHandler = Pins.createLineHandler = function (ref, errorHandler) {
|
||||
return sanitized;
|
||||
};
|
||||
|
||||
return function (line) {
|
||||
return function (line, i) {
|
||||
ref.index++;
|
||||
if (!Boolean(line)) { return; }
|
||||
|
||||
@ -74,6 +75,7 @@ var createLineHandler = Pins.createLineHandler = function (ref, errorHandler) {
|
||||
}
|
||||
|
||||
if (typeof(l[2]) === 'number') {
|
||||
if (!ref.first) { ref.first = l[2]; }
|
||||
ref.latest = l[2]; // date
|
||||
}
|
||||
|
||||
@ -109,6 +111,11 @@ var createLineHandler = Pins.createLineHandler = function (ref, errorHandler) {
|
||||
default:
|
||||
errorHandler("PIN_LINE_UNSUPPORTED_COMMAND", l);
|
||||
}
|
||||
|
||||
if (i === 0) { // First line when using Pins.load
|
||||
if (l[0] === 'PIN' || ref.block) { ref.user = true; } // teams always start with RESET
|
||||
}
|
||||
|
||||
};
|
||||
};
|
||||
|
||||
|
||||
@ -10,6 +10,9 @@ var BlobStore = module.exports;
|
||||
var nThen = require("nthen");
|
||||
var Semaphore = require("saferphore");
|
||||
var Util = require("../common-util");
|
||||
const PERMISSIVE = 511;
|
||||
|
||||
const readFileBin = require("../stream-file").readFileBin;
|
||||
|
||||
const BLOB_LENGTH = 48;
|
||||
|
||||
@ -31,37 +34,30 @@ var prependArchive = function (Env, path) {
|
||||
return Path.join(Env.archivePath, 'blob', relativePathToBlob);
|
||||
};
|
||||
|
||||
// /blob/<safeKeyPrefix>/<safeKey>/<blobPrefix>/<blobId>
|
||||
// /blob/<blobPrefix>/<blobId>
|
||||
var makeBlobPath = function (Env, blobId) {
|
||||
return Path.join(Env.blobPath, blobId.slice(0, 2), blobId);
|
||||
};
|
||||
|
||||
|
||||
var makeActivityPath = function (Env, blobId) {
|
||||
return makeBlobPath(Env, blobId) + '.activity';
|
||||
};
|
||||
|
||||
// /blob/<blobPrefix>/<blobId>.metadata.ndjson
|
||||
var mkMetadataPath = function (Env, blobId) {
|
||||
return Path.join(Env.blobPath, blobId.slice(0, 2), blobId) + '.metadata.ndjson';
|
||||
};
|
||||
|
||||
// /blobstate/<safeKeyPrefix>/<safeKey>
|
||||
var makeStagePath = function (Env, safeKey) {
|
||||
return Path.join(Env.blobStagingPath, safeKey.slice(0, 2), safeKey);
|
||||
};
|
||||
|
||||
// /blob/<safeKeyPrefix>/<safeKey>/<blobPrefix>/<blobId>
|
||||
var makeProofPath = function (Env, safeKey, blobId) {
|
||||
return Path.join(Env.blobPath, safeKey.slice(0, 3), safeKey, blobId.slice(0, 2), blobId);
|
||||
};
|
||||
|
||||
var mkPlaceholderPath = function (Env, blobId) {
|
||||
return makeBlobPath(Env, blobId) + '.placeholder';
|
||||
};
|
||||
|
||||
var parseProofPath = function (path) {
|
||||
var parts = path.split('/');
|
||||
return {
|
||||
blobId: parts[parts.length -1],
|
||||
safeKey: parts[parts.length - 3],
|
||||
};
|
||||
};
|
||||
|
||||
// Placeholder for deleted files
|
||||
var addPlaceholder = function (Env, blobId, reason, cb) {
|
||||
if (!reason) { return cb(); }
|
||||
@ -120,6 +116,18 @@ var isFile = function (filePath, cb) {
|
||||
});
|
||||
};
|
||||
|
||||
// PROOFS
|
||||
// DEPRECATED, keep for compatibility
|
||||
// /blob/<safeKeyPrefix>/<safeKey>/<blobPrefix>/<blobId>
|
||||
var makeProofPath = function (Env, safeKey, blobId) {
|
||||
return Path.join(Env.blobPath, safeKey.slice(0, 3), safeKey, blobId.slice(0, 2), blobId);
|
||||
};
|
||||
// isOwnedBy(id, safeKey)
|
||||
var isOwnedBy = function (Env, safeKey, blobId, cb) {
|
||||
var proofPath = makeProofPath(Env, safeKey, blobId);
|
||||
isFile(proofPath, cb);
|
||||
};
|
||||
|
||||
var makeFileStream = function (full, _cb) {
|
||||
var cb = Util.once(Util.mkAsync(_cb));
|
||||
Fse.mkdirp(Path.dirname(full), function (e) {
|
||||
@ -190,6 +198,76 @@ var getActivity = function (Env, blobId, cb) {
|
||||
});
|
||||
};
|
||||
|
||||
// destroyStream && createIdleStreamCollector
|
||||
// copied from lib/storage/file.js
|
||||
// see comments there
|
||||
const STREAM_CLOSE_TIMEOUT = 120000;
|
||||
const STREAM_DESTROY_TIMEOUT = 30000;
|
||||
const destroyStream = function (stream) {
|
||||
if (!stream) { return; }
|
||||
try {
|
||||
stream.close();
|
||||
if (stream.closed && stream.fd === null) { return; }
|
||||
} catch (err) {
|
||||
console.error(err);
|
||||
}
|
||||
setTimeout(function () {
|
||||
try { stream.destroy(); } catch (err) { console.error(err); }
|
||||
}, STREAM_DESTROY_TIMEOUT);
|
||||
};
|
||||
const createIdleStreamCollector = function (stream) {
|
||||
var collector = Util.once(Util.mkAsync(Util.bake(destroyStream, [stream])));
|
||||
collector.keepAlive = Util.throttle(collector, STREAM_CLOSE_TIMEOUT);
|
||||
collector.keepAlive();
|
||||
return collector;
|
||||
};
|
||||
|
||||
// writeMetadata appends to the dedicated log of metadata amendments
|
||||
var writeMetadata = function (env, channelId, data, cb) {
|
||||
var path = mkMetadataPath(env, channelId);
|
||||
|
||||
Fse.mkdirp(Path.dirname(path), PERMISSIVE, function (err) {
|
||||
if (err && err.code !== 'EEXIST') { return void cb(err); }
|
||||
Fs.appendFile(path, data + '\n', cb);
|
||||
});
|
||||
};
|
||||
var archiveMetadata = (Env, blobId, cb) => {
|
||||
var path = mkMetadataPath(Env, blobId);
|
||||
var archivePath = prependArchive(Env, path);
|
||||
// XXX eviction clean lone md files
|
||||
// if we fail to delete the metadata file, it can still be removed later by the eviction script
|
||||
Fse.move(path, archivePath, { overwrite: true }, cb);
|
||||
};
|
||||
var restoreMetadata = function (Env, blobId, cb) {
|
||||
var path = mkMetadataPath(Env, blobId);
|
||||
var archivePath = prependArchive(Env, path);
|
||||
Fse.move(archivePath, path, cb);
|
||||
};
|
||||
var readBlobMetadata = function (env, blobId, handler, _cb) {
|
||||
var metadataPath = mkMetadataPath(env, blobId);
|
||||
var stream = Fs.createReadStream(metadataPath, {start: 0});
|
||||
|
||||
const collector = createIdleStreamCollector(stream);
|
||||
var cb = Util.both(_cb, collector);
|
||||
|
||||
readFileBin(stream, function (msgObj, readMore) {
|
||||
collector.keepAlive();
|
||||
var line = msgObj.buff.toString('utf8');
|
||||
try {
|
||||
var parsed = JSON.parse(line);
|
||||
handler(null, parsed);
|
||||
} catch (err) {
|
||||
handler(err, line);
|
||||
}
|
||||
readMore();
|
||||
}, function (err) {
|
||||
// ENOENT => there is no metadata log
|
||||
if (!err || err.code === 'ENOENT') { return void cb(); }
|
||||
// otherwise stream errors?
|
||||
cb(err);
|
||||
});
|
||||
};
|
||||
|
||||
/********** METHODS **************/
|
||||
|
||||
var upload = function (Env, safeKey, content, cb) {
|
||||
@ -313,6 +391,9 @@ var tryId = function (path, cb) {
|
||||
};
|
||||
|
||||
// owned_upload_complete
|
||||
let unescapeKeyCharacters = function (key) {
|
||||
return key.replace(/\-/g, '/');
|
||||
};
|
||||
var owned_upload_complete = function (Env, safeKey, id, cb) {
|
||||
closeBlobstage(Env, safeKey);
|
||||
if (!isValidId(id)) {
|
||||
@ -325,12 +406,9 @@ var owned_upload_complete = function (Env, safeKey, id, cb) {
|
||||
}
|
||||
|
||||
var finalPath = makeBlobPath(Env, id);
|
||||
let unsafeKey = unescapeKeyCharacters(safeKey);
|
||||
|
||||
var finalOwnPath = makeProofPath(Env, safeKey, id);
|
||||
|
||||
// the user wants to move it into blob and create a empty file with the same id
|
||||
// in their own space:
|
||||
// /blob/safeKeyPrefix/safeKey/blobPrefix/blobID
|
||||
// the user wants to move it into blob and create a metadata log with an owner
|
||||
|
||||
nThen(function (w) {
|
||||
// make the requisite directory structure using Mkdirp
|
||||
@ -340,12 +418,6 @@ var owned_upload_complete = function (Env, safeKey, id, cb) {
|
||||
return void cb(e.code);
|
||||
}
|
||||
}));
|
||||
Fse.mkdirp(Path.dirname(finalOwnPath), w(function (e /*, path */) {
|
||||
if (e) { // does not throw error if the directory already existed
|
||||
w.abort();
|
||||
return void cb(e.code);
|
||||
}
|
||||
}));
|
||||
}).nThen(function (w) {
|
||||
// make sure the id does not collide with another
|
||||
tryId(finalPath, w(function (e) {
|
||||
@ -355,8 +427,11 @@ var owned_upload_complete = function (Env, safeKey, id, cb) {
|
||||
}
|
||||
}));
|
||||
}).nThen(function (w) {
|
||||
// Create the empty file proving ownership
|
||||
Fs.writeFile(finalOwnPath, '', w(function (e) {
|
||||
// Write the metadata
|
||||
let md = JSON.stringify({
|
||||
owners: [unsafeKey]
|
||||
});
|
||||
writeMetadata(Env, id, md, w((e) => {
|
||||
if (e) {
|
||||
w.abort();
|
||||
return void cb(e.code);
|
||||
@ -367,12 +442,6 @@ var owned_upload_complete = function (Env, safeKey, id, cb) {
|
||||
// move the existing file to its new path
|
||||
Fse.move(oldPath, finalPath, w(function (e) {
|
||||
if (e) {
|
||||
// if there's an error putting the file into its final location...
|
||||
// ... you should remove the ownership file
|
||||
Fs.unlink(finalOwnPath, function () {
|
||||
// but if you can't, it's not catestrophic
|
||||
// we can clean it up later
|
||||
});
|
||||
w.abort();
|
||||
return void cb(e.code);
|
||||
}
|
||||
@ -393,24 +462,12 @@ var remove = function (Env, blobId, cb) {
|
||||
clearActivity(Env, blobId, () => {});
|
||||
};
|
||||
|
||||
// removeProof
|
||||
var removeProof = function (Env, safeKey, blobId, cb) {
|
||||
var proofPath = makeProofPath(Env, safeKey, blobId);
|
||||
Fs.unlink(proofPath, cb);
|
||||
};
|
||||
|
||||
// isOwnedBy(id, safeKey)
|
||||
var isOwnedBy = function (Env, safeKey, blobId, cb) {
|
||||
var proofPath = makeProofPath(Env, safeKey, blobId);
|
||||
isFile(proofPath, cb);
|
||||
};
|
||||
|
||||
|
||||
// archiveBlob
|
||||
var archiveBlob = function (Env, blobId, reason, cb) {
|
||||
var blobPath = makeBlobPath(Env, blobId);
|
||||
var archivePath = prependArchive(Env, blobPath);
|
||||
Fse.move(blobPath, archivePath, { overwrite: true }, cb);
|
||||
archiveMetadata(Env, blobId, () => {});
|
||||
archiveActivity(Env, blobId, () => {});
|
||||
addPlaceholder(Env, blobId, reason, () => {});
|
||||
};
|
||||
@ -426,29 +483,11 @@ var restoreBlob = function (Env, blobId, cb) {
|
||||
var blobPath = makeBlobPath(Env, blobId);
|
||||
var archivePath = prependArchive(Env, blobPath);
|
||||
Fse.move(archivePath, blobPath, cb);
|
||||
restoreMetadata(Env, blobId, () => {});
|
||||
restoreActivity(Env, blobId, () => {});
|
||||
clearPlaceholder(Env, blobId, () => {});
|
||||
};
|
||||
|
||||
// archiveProof
|
||||
var archiveProof = function (Env, safeKey, blobId, cb) {
|
||||
var proofPath = makeProofPath(Env, safeKey, blobId);
|
||||
var archivePath = prependArchive(Env, proofPath);
|
||||
Fse.move(proofPath, archivePath, { overwrite: true }, cb);
|
||||
};
|
||||
|
||||
var removeArchivedProof = function (Env, safeKey, blobId, cb) {
|
||||
var archivedPath = prependArchive(Env, makeProofPath(Env, safeKey, blobId));
|
||||
Fs.unlink(archivedPath, cb);
|
||||
};
|
||||
|
||||
// restoreProof
|
||||
var restoreProof = function (Env, safeKey, blobId, cb) {
|
||||
var proofPath = makeProofPath(Env, safeKey, blobId);
|
||||
var archivePath = prependArchive(Env, proofPath);
|
||||
Fse.move(archivePath, proofPath, cb);
|
||||
};
|
||||
|
||||
var makeWalker = function (n, handleChild, done) {
|
||||
if (!n || typeof(n) !== 'number' || n < 2) { n = 2; }
|
||||
|
||||
@ -480,7 +519,7 @@ var makeWalker = function (n, handleChild, done) {
|
||||
}
|
||||
if (!stats.isDirectory()) {
|
||||
w.abort();
|
||||
if (/\.activity$/.test(path)) {
|
||||
if (/\.activity$/.test(path)) {
|
||||
// NOTE: some activity files were created for deleted blobs due to
|
||||
// a bug. We're going to detect them here in order to be able to clean
|
||||
// them.
|
||||
@ -513,46 +552,6 @@ var makeWalker = function (n, handleChild, done) {
|
||||
return recurse;
|
||||
};
|
||||
|
||||
var listProofs = function (root, handler, cb) {
|
||||
Fs.readdir(root, function (err, dir) {
|
||||
if (err) { return void cb(err); }
|
||||
|
||||
var walk = makeWalker(20, function (err, path, next, loneActivity) {
|
||||
if (loneActivity) { return void next(); }
|
||||
// path is the path to a child node on the filesystem
|
||||
|
||||
// next handles the next job in a queue
|
||||
|
||||
// iterate over proofs
|
||||
// check for presence of corresponding files
|
||||
Fs.stat(path, function (err, stats) {
|
||||
if (err) {
|
||||
return void handler(err, void 0, next);
|
||||
}
|
||||
|
||||
var parsed = parseProofPath(path);
|
||||
handler(void 0, {
|
||||
path: path,
|
||||
blobId: parsed.blobId,
|
||||
safeKey: parsed.safeKey,
|
||||
atime: stats.atime,
|
||||
ctime: stats.ctime,
|
||||
mtime: stats.mtime,
|
||||
}, next);
|
||||
});
|
||||
}, function () {
|
||||
// called when there are no more directories or children to process
|
||||
cb();
|
||||
});
|
||||
|
||||
dir.forEach(function (d) {
|
||||
// ignore directories that aren't 3 characters long...
|
||||
if (d.length !== 3) { return; }
|
||||
walk(Path.join(root, d));
|
||||
});
|
||||
});
|
||||
};
|
||||
|
||||
var getActivityStat = function (path, base, cb) {
|
||||
var suffix = base ? '' : '.activity';
|
||||
Fs.stat(path+suffix, function (err, stats) {
|
||||
@ -560,32 +559,92 @@ var getActivityStat = function (path, base, cb) {
|
||||
cb(err, stats);
|
||||
});
|
||||
};
|
||||
var listBlobs = function (root, handler, cb) {
|
||||
// iterate over files
|
||||
Fs.readdir(root, function (err, dir) {
|
||||
if (err) { return void cb(err); }
|
||||
var walk = makeWalker(20, function (err, path, next, loneActivity) {
|
||||
if (loneActivity) { return void next(); }
|
||||
getActivityStat(path, false, function (err, stats) {
|
||||
if (err) {
|
||||
return void handler(err, void 0, next);
|
||||
}
|
||||
var getStats = function (Env, blobId, cb) {
|
||||
var path = makeBlobPath(Env, blobId);
|
||||
getActivityStat(path, false, cb);
|
||||
};
|
||||
|
||||
handler(void 0, {
|
||||
blobId: Path.basename(path),
|
||||
atime: stats.atime,
|
||||
ctime: stats.ctime,
|
||||
mtime: stats.mtime,
|
||||
}, next);
|
||||
});
|
||||
}, function () {
|
||||
cb();
|
||||
});
|
||||
let blobRegex = /^[0-9a-fA-F]{48}(\.metadata)*(\.ndjson)*$/;
|
||||
var listBlobs = function (root, handler, fast, cb) {
|
||||
var dirList = [];
|
||||
|
||||
dir.forEach(function (d) {
|
||||
if (d.length !== 2) { return; }
|
||||
walk(Path.join(root, d));
|
||||
nThen(function (w) {
|
||||
// the root of your datastore contains nested directories...
|
||||
Fs.readdir(root, w(function (err, list) {
|
||||
if (err) {
|
||||
w.abort();
|
||||
// TODO check if we normally return strings or errors
|
||||
return void cb(err);
|
||||
}
|
||||
dirList = list;
|
||||
}));
|
||||
}).nThen(function (waitFor) {
|
||||
// search inside the nested directories
|
||||
// stream it so you don't put unnecessary data in memory
|
||||
var n = nThen;
|
||||
dirList.forEach(function (dir) {
|
||||
if (dir.length !== 2) { return; }
|
||||
// Handle one directory at a time to save some memory
|
||||
n = n(function (w) {
|
||||
// do twenty things at a time
|
||||
var sema = Semaphore.create(20);
|
||||
var nestedDirPath = Path.join(root, dir);
|
||||
Fs.readdir(nestedDirPath, w(function (err, list) {
|
||||
if (err) { return void handler(err); } // Is this correct?
|
||||
list.forEach(function (item) {
|
||||
// ignore hidden files
|
||||
if (/^\./.test(item)) { return; }
|
||||
// ignore anything that isn't channel or metadata
|
||||
if (!blobRegex.test(item)) { return; }
|
||||
|
||||
var isLonelyMetadata = false;
|
||||
var blobName;
|
||||
|
||||
// if the current file is not the channel data, then it must be metadata
|
||||
if (!/^[0-9a-fA-F]{48}$/.test(item)) {
|
||||
blobName = item.replace(/\.metadata\.ndjson/, '');
|
||||
// check if blob already exists
|
||||
if (list.indexOf(blobName) !== -1) { return; }
|
||||
// otherwise set a flag indicating that we should
|
||||
// handle the metadata on its own
|
||||
isLonelyMetadata = true;
|
||||
} else {
|
||||
blobName = item;
|
||||
}
|
||||
if (blobName.length !== 48) { return; }
|
||||
|
||||
sema.take(function (give) {
|
||||
var next = w(give());
|
||||
|
||||
if (fast) {
|
||||
return void handler(void 0, {
|
||||
blobId: blobName
|
||||
}, next);
|
||||
}
|
||||
|
||||
var filePath = Path.join(nestedDirPath, blobName);
|
||||
if (isLonelyMetadata) {
|
||||
// Set time to 0 to delete this
|
||||
// lonely metadata file
|
||||
return void handler(void 0, {
|
||||
blobId: blobName,
|
||||
mtime: 0,
|
||||
atime: 0,
|
||||
ctime: 0
|
||||
}, next);
|
||||
}
|
||||
return void getActivityStat(filePath, false, (err, data) => {
|
||||
data.blobId = blobName;
|
||||
handler(err, data, next);
|
||||
});
|
||||
});
|
||||
});
|
||||
}));
|
||||
}).nThen;
|
||||
});
|
||||
n(waitFor());
|
||||
}).nThen(function () {
|
||||
cb();
|
||||
});
|
||||
};
|
||||
|
||||
@ -673,6 +732,20 @@ BlobStore.create = function (config, _cb) {
|
||||
if (!isValidSafeKey(safeKey)) { return void cb('INVALID_SAFEKEY'); }
|
||||
isOwnedBy(Env, safeKey, blobId, cb);
|
||||
},
|
||||
readMetadata: (blobId, handler, cb) => {
|
||||
if (!isValidId(blobId)) { return void cb("INVALID_ID"); }
|
||||
readBlobMetadata(Env, blobId, handler, cb);
|
||||
},
|
||||
writeMetadata: (blobId, data, cb) => {
|
||||
if (!isValidId(blobId)) { return void cb("INVALID_ID"); }
|
||||
writeMetadata(Env, blobId, data, cb);
|
||||
},
|
||||
hasMetadata: (blobId, _cb) => {
|
||||
var cb = Util.once(Util.mkAsync(_cb));
|
||||
if (!isValidId(blobId)) { return void cb("INVALID_ID"); }
|
||||
var path = mkMetadataPath(Env, blobId);
|
||||
isFile(path, cb);
|
||||
},
|
||||
|
||||
remove: {
|
||||
blob: function (blobId, _cb) {
|
||||
@ -680,24 +753,12 @@ BlobStore.create = function (config, _cb) {
|
||||
if (!isValidId(blobId)) { return void cb("INVALID_ID"); }
|
||||
remove(Env, blobId, cb);
|
||||
},
|
||||
proof: function (safeKey, blobId, _cb) {
|
||||
var cb = Util.once(Util.mkAsync(_cb));
|
||||
if (!isValidSafeKey(safeKey)) { return void cb('INVALID_SAFEKEY'); }
|
||||
if (!isValidId(blobId)) { return void cb("INVALID_ID"); }
|
||||
removeProof(Env, safeKey, blobId, cb);
|
||||
},
|
||||
archived: {
|
||||
blob: function (blobId, _cb) {
|
||||
var cb = Util.once(Util.mkAsync(_cb));
|
||||
if (!isValidId(blobId)) { return void cb("INVALID_ID"); }
|
||||
removeArchivedBlob(Env, blobId, cb);
|
||||
},
|
||||
proof: function (safeKey, blobId, _cb) {
|
||||
var cb = Util.once(Util.mkAsync(_cb));
|
||||
if (!isValidSafeKey(safeKey)) { return void cb('INVALID_SAFEKEY'); }
|
||||
if (!isValidId(blobId)) { return void cb("INVALID_ID"); }
|
||||
removeArchivedProof(Env, safeKey, blobId, cb);
|
||||
},
|
||||
},
|
||||
loneActivity: function (_cb) {
|
||||
var cb = Util.once(Util.mkAsync(_cb));
|
||||
@ -711,12 +772,6 @@ BlobStore.create = function (config, _cb) {
|
||||
if (!isValidId(blobId)) { return void cb("INVALID_ID"); }
|
||||
archiveBlob(Env, blobId, reason, cb);
|
||||
},
|
||||
proof: function (safeKey, blobId, _cb) {
|
||||
var cb = Util.once(Util.mkAsync(_cb));
|
||||
if (!isValidSafeKey(safeKey)) { return void cb('INVALID_SAFEKEY'); }
|
||||
if (!isValidId(blobId)) { return void cb("INVALID_ID"); }
|
||||
archiveProof(Env, safeKey, blobId, cb);
|
||||
},
|
||||
},
|
||||
|
||||
restore: {
|
||||
@ -725,12 +780,6 @@ BlobStore.create = function (config, _cb) {
|
||||
if (!isValidId(blobId)) { return void cb("INVALID_ID"); }
|
||||
restoreBlob(Env, blobId, cb);
|
||||
},
|
||||
proof: function (safeKey, blobId, _cb) {
|
||||
var cb = Util.once(Util.mkAsync(_cb));
|
||||
if (!isValidSafeKey(safeKey)) { return void cb('INVALID_SAFEKEY'); }
|
||||
if (!isValidId(blobId)) { return void cb("INVALID_ID"); }
|
||||
restoreProof(Env, safeKey, blobId, cb);
|
||||
},
|
||||
},
|
||||
|
||||
isBlobAvailable: function (blobId, _cb) {
|
||||
@ -781,24 +830,21 @@ BlobStore.create = function (config, _cb) {
|
||||
if (!isValidId(id)) { return void cb("INVALID_ID"); }
|
||||
getActivity(Env, id, cb);
|
||||
},
|
||||
getStats: function (id, _cb) {
|
||||
var cb = Util.once(Util.mkAsync(_cb));
|
||||
if (!isValidId(id)) { return void cb("INVALID_ID"); }
|
||||
getStats(Env, id, cb);
|
||||
},
|
||||
|
||||
list: {
|
||||
blobs: function (handler, _cb) {
|
||||
blobs: function (handler, _cb, fast) {
|
||||
var cb = Util.once(Util.mkAsync(_cb));
|
||||
listBlobs(Env.blobPath, handler, cb);
|
||||
},
|
||||
proofs: function (handler, _cb) {
|
||||
var cb = Util.once(Util.mkAsync(_cb));
|
||||
listProofs(Env.blobPath, handler, cb);
|
||||
listBlobs(Env.blobPath, handler, fast, cb);
|
||||
},
|
||||
archived: {
|
||||
proofs: function (handler, _cb) {
|
||||
blobs: function (handler, _cb, fast) {
|
||||
var cb = Util.once(Util.mkAsync(_cb));
|
||||
listProofs(prependArchive(Env, Env.blobPath), handler, cb);
|
||||
},
|
||||
blobs: function (handler, _cb) {
|
||||
var cb = Util.once(Util.mkAsync(_cb));
|
||||
listBlobs(prependArchive(Env, Env.blobPath), handler, cb);
|
||||
listBlobs(prependArchive(Env, Env.blobPath), handler, fast, cb);
|
||||
},
|
||||
}
|
||||
},
|
||||
|
||||
@ -158,6 +158,12 @@ const isValidOffsetNumber = function (n) {
|
||||
return typeof(n) === 'number' && n >= 0;
|
||||
};
|
||||
|
||||
const updateEnv = data => {
|
||||
const {value} = data;
|
||||
let env = Util.tryParse(value) || {};
|
||||
Env.proofsMigrated = env?.proofsMigrated;
|
||||
};
|
||||
|
||||
const computeIndexFromOffset = function (channelName, offset, cb) {
|
||||
let cpIndex = [];
|
||||
let messageBuf = [];
|
||||
@ -338,7 +344,13 @@ const computeMetadata = function (data, cb) {
|
||||
const ref = {};
|
||||
const lineHandler = Meta.createLineHandler(ref, Env.Log.error);
|
||||
monitoringIncrement('computeMetadata');
|
||||
return void store.readChannelMetadata(data.channel, lineHandler, function (err) {
|
||||
|
||||
let f = store.readChannelMetadata;
|
||||
if (data.channel.length === HK.BLOB_ID_LENGTH) {
|
||||
f = blobStore.readMetadata;
|
||||
}
|
||||
|
||||
return void f(data.channel, lineHandler, function (err) {
|
||||
if (err) {
|
||||
// stream errors?
|
||||
return void cb(err);
|
||||
@ -558,16 +570,33 @@ const removeOwnedBlob = function (data, cb) {
|
||||
if (typeof(data.safeKey) !== 'string') { return void cb("INVALID_KEY"); }
|
||||
const blobId = data.blobId;
|
||||
const safeKey = Util.escapeKeyCharacters(data.safeKey);
|
||||
const unsafeKey = Util.unescapeKeyCharacters(data.safeKey);
|
||||
|
||||
const reason = data.reason || 'ARCHIVE_OWNED';
|
||||
|
||||
nThen(function (w) {
|
||||
// check if you have permissions
|
||||
blobStore.isOwnedBy(safeKey, blobId, w(function (err, owned) {
|
||||
if (err || !owned) {
|
||||
computeMetadata({channel: blobId}, w((err, meta) => {
|
||||
if (err || !meta) {
|
||||
w.abort();
|
||||
return void cb("INSUFFICIENT_PERMISSIONS");
|
||||
}
|
||||
let owners = meta.owners;
|
||||
if (!owners && !Env.proofsMigrated) {
|
||||
// Check old proofs during migration
|
||||
blobStore.isOwnedBy(safeKey, blobId, w((e, owned) => {
|
||||
if (e || !owned) {
|
||||
w.abort();
|
||||
return void cb("INSUFFICIENT_PERMISSIONS");
|
||||
}
|
||||
}));
|
||||
return;
|
||||
}
|
||||
if (!owners || !owners.includes(unsafeKey)) {
|
||||
w.abort();
|
||||
return void cb("INSUFFICIENT_PERMISSIONS");
|
||||
}
|
||||
// Owned, continue
|
||||
}));
|
||||
}).nThen(function (w) {
|
||||
// remove the blob
|
||||
@ -581,20 +610,8 @@ const removeOwnedBlob = function (data, cb) {
|
||||
w.abort();
|
||||
return void cb(err);
|
||||
}
|
||||
}));
|
||||
}).nThen(function () {
|
||||
// archive the proof
|
||||
blobStore.archive.proof(safeKey, blobId, function (err) {
|
||||
Env.Log.info("ARCHIVAL_PROOF_REMOVAL_BY_OWNER_RPC", {
|
||||
safeKey: safeKey,
|
||||
blobId: blobId,
|
||||
status: err? String(err): 'SUCCESS',
|
||||
});
|
||||
if (err) {
|
||||
return void cb("E_PROOF_REMOVAL");
|
||||
}
|
||||
cb(void 0, 'OK');
|
||||
});
|
||||
}));
|
||||
});
|
||||
};
|
||||
|
||||
@ -692,6 +709,7 @@ const getLastChannelTime = function (data, cb) {
|
||||
};
|
||||
|
||||
const COMMANDS = {
|
||||
ENV_UPDATE: updateEnv,
|
||||
COMPUTE_INDEX: computeIndex,
|
||||
COMPUTE_METADATA: computeMetadata,
|
||||
GET_OLDER_HISTORY: getOlderHistory,
|
||||
@ -847,6 +865,9 @@ process.on('message', function (data) {
|
||||
};
|
||||
|
||||
if (!ready) {
|
||||
if (data.env) {
|
||||
updateEnv({value:data.env});
|
||||
}
|
||||
return void init(data.config, function (err) {
|
||||
if (err) { return void cb(Util.serializeError(err)); }
|
||||
ready = true;
|
||||
|
||||
@ -9,6 +9,7 @@ const { fork } = require('child_process');
|
||||
const Workers = module.exports;
|
||||
const PID = process.pid;
|
||||
const Block = require("../storage/block");
|
||||
const Environment = require('../env');
|
||||
|
||||
const DB_PATH = 'lib/workers/db-worker';
|
||||
const MAX_JOBS = 16;
|
||||
@ -256,6 +257,7 @@ Workers.initialize = function (Env, config, _cb) {
|
||||
pid: PID,
|
||||
txid: txid,
|
||||
config: config,
|
||||
env: Environment.serialize(Env)
|
||||
});
|
||||
|
||||
worker.on('message', function (res) {
|
||||
@ -342,7 +344,8 @@ Workers.initialize = function (Env, config, _cb) {
|
||||
type: 'broadcast',
|
||||
pid: PID,
|
||||
command: data.command,
|
||||
txid: data.txid
|
||||
txid: data.txid,
|
||||
value: data.value
|
||||
});
|
||||
});
|
||||
return workers;
|
||||
|
||||
198
scripts/migrations/migrate-blob-proofs.js
Normal file
198
scripts/migrations/migrate-blob-proofs.js
Normal file
@ -0,0 +1,198 @@
|
||||
// SPDX-FileCopyrightText: 2025 XWiki CryptPad Team <contact@cryptpad.org> and contributors
|
||||
//
|
||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
|
||||
const { parentPort } = require('node:worker_threads');
|
||||
const Path = require('node:path');
|
||||
const Fs = require('node:fs');
|
||||
const nThen = require("nthen");
|
||||
const Semaphore = require("saferphore");
|
||||
const Logger = require("../../lib/log");
|
||||
const BlobStorage = require("../../lib/storage/blob");
|
||||
let config = require("../../lib/load-config");
|
||||
|
||||
|
||||
const blobPath = config.blobPath || './blob';
|
||||
let Log = {};
|
||||
|
||||
// NOTE: in cleaning mode, we DON'T migrate
|
||||
// (we suppose data has already been migrated)
|
||||
const start = (clean, dry, cb) => {
|
||||
const DRY_RUN = dry;
|
||||
|
||||
let dirList = [];
|
||||
let blobStore;
|
||||
nThen(w => {
|
||||
Logger.create(config, w(function (_log) {
|
||||
Log = _log;
|
||||
}));
|
||||
}).nThen(w => {
|
||||
config.getSession = function () {};
|
||||
BlobStorage.create(config, w(function (err, _store) {
|
||||
if (err) {
|
||||
w.abort();
|
||||
return void Log.error("ERR_BLOB_STORE", err);
|
||||
}
|
||||
blobStore = _store;
|
||||
}));
|
||||
}).nThen(w => {
|
||||
Fs.readdir(blobPath, w((err, list) => {
|
||||
if (err) {
|
||||
w.abort();
|
||||
return void Log.error("ERR_READING_ROOT", err);
|
||||
}
|
||||
dirList = list;
|
||||
}));
|
||||
}).nThen(() => {
|
||||
let n = nThen;
|
||||
dirList.forEach(dir => {
|
||||
if (dir.length !== 3) { return; }
|
||||
// ./blob/abc
|
||||
const nestedDirPath = Path.join(blobPath, dir);
|
||||
|
||||
if (clean) {
|
||||
n = n(ww => {
|
||||
Log.info("REMOVING_DIR", nestedDirPath);
|
||||
if (DRY_RUN) { return; }
|
||||
Fs.rm(nestedDirPath, {
|
||||
recursive: true, force: true
|
||||
}, ww(err => {
|
||||
if (err) {
|
||||
Log.error("ERR_REMOVE_DIR", {
|
||||
path: nestedDirPath,
|
||||
err
|
||||
});
|
||||
}
|
||||
}));
|
||||
}).nThen;
|
||||
return;
|
||||
}
|
||||
|
||||
n = n(w => {
|
||||
// One user at a time
|
||||
const sema = Semaphore.create(1);
|
||||
let nestedDirList = [];
|
||||
nThen(ww => {
|
||||
Fs.readdir(nestedDirPath, ww((err, list) => {
|
||||
if (err) {
|
||||
w.abort();
|
||||
ww.abort();
|
||||
return Log.error("ERR_READING_DIR", {
|
||||
path: nestedDirPath,
|
||||
err
|
||||
});
|
||||
}
|
||||
nestedDirList = list;
|
||||
}));
|
||||
}).nThen(ww => {
|
||||
nestedDirList.forEach(key => {
|
||||
// ./blob/abc/abcdefg...
|
||||
const keyPath = Path.join(nestedDirPath, key);
|
||||
sema.take(give => {
|
||||
let edPublic = key.replace(/\-/g, '/');
|
||||
let md = JSON.stringify({ owners: [edPublic] });
|
||||
Log.info("START_USER", edPublic);
|
||||
Fs.readdir(keyPath, ww((err, list) => {
|
||||
if (err) {
|
||||
w.abort();
|
||||
ww.abort();
|
||||
return Log.error("ERR_READING_DIR", {
|
||||
path: keyPath,
|
||||
err
|
||||
});
|
||||
}
|
||||
let blobs = [];
|
||||
nThen(www => {
|
||||
list.forEach(dir => {
|
||||
// ./blob/abc/abcdefg.../01
|
||||
const path = Path.join(keyPath, dir);
|
||||
Fs.readdir(path, www((err, blobsList) => {
|
||||
if (err) {
|
||||
w.abort();
|
||||
ww.abort();
|
||||
www.abort();
|
||||
return Log.error("ERR_READING_DIR", {
|
||||
path, err
|
||||
});
|
||||
}
|
||||
Array.prototype.push.apply(blobs, blobsList);
|
||||
}));
|
||||
});
|
||||
}).nThen(www => {
|
||||
// migrate 20 blobs at a time for a given user
|
||||
const sema = Semaphore.create(20);
|
||||
blobs.forEach(blobId => {
|
||||
sema.take(ggive => {
|
||||
blobStore.isBlobAvailable(blobId, www((err, blobExists) => {
|
||||
blobStore.hasMetadata(blobId, www((err, exists) => {
|
||||
// If blob is not available or metadata already
|
||||
// exists, don't write md file
|
||||
if (!blobExists || exists) { return void ggive(); }
|
||||
Log.info('WRITE_METADATA', blobId);
|
||||
if (DRY_RUN) { return void ggive(); }
|
||||
blobStore.writeMetadata(blobId, md, www(e => {
|
||||
if (e) {
|
||||
w.abort();
|
||||
ww.abort();
|
||||
www.abort();
|
||||
return Log.error("ERR_WRITING_MD", { blobId });
|
||||
}
|
||||
ggive();
|
||||
}));
|
||||
|
||||
}));
|
||||
}));
|
||||
});
|
||||
});
|
||||
}).nThen(ww(give(() => {
|
||||
Log.info("END_USER", edPublic);
|
||||
})));
|
||||
}));
|
||||
});
|
||||
|
||||
});
|
||||
}).nThen(w());
|
||||
}).nThen;
|
||||
});
|
||||
n(() => {
|
||||
Log.info("DONE");
|
||||
cb();
|
||||
});
|
||||
});
|
||||
};
|
||||
|
||||
if (parentPort) {
|
||||
// Loaded as worker script
|
||||
config = JSON.parse(JSON.stringify(config));
|
||||
config.logToStdout = false;
|
||||
parentPort.on('message', (message) => {
|
||||
let parsed = message; //JSON.parse(message);
|
||||
if (!parsed?.start) { return; }
|
||||
// Migrate
|
||||
start(false, false, () => {
|
||||
parentPort.postMessage('MIGRATED');
|
||||
// If success, clean
|
||||
start(true, false, () => {
|
||||
parentPort.postMessage('CLEANED');
|
||||
});
|
||||
});
|
||||
});
|
||||
parentPort.postMessage('READY');
|
||||
} else if (require.main === module) {
|
||||
// Loaded from command-line
|
||||
let dry = false;
|
||||
let clean = false;
|
||||
process.argv.forEach(key => {
|
||||
if (key === '--dry') {
|
||||
dry = true;
|
||||
return;
|
||||
}
|
||||
if (key === '--clean') {
|
||||
clean = true;
|
||||
return;
|
||||
}
|
||||
});
|
||||
start(clean, dry, () => {
|
||||
process.exit(0);
|
||||
});
|
||||
}
|
||||
132
scripts/user-statistics.js
Normal file
132
scripts/user-statistics.js
Normal file
@ -0,0 +1,132 @@
|
||||
// SPDX-FileCopyrightText: 2025 XWiki CryptPad Team <contact@cryptpad.org> and contributors
|
||||
//
|
||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
|
||||
const nThen = require("nthen");
|
||||
const Semaphore = require("saferphore");
|
||||
const Logger = require("../lib/log");
|
||||
const Pins = require("../lib/pins");
|
||||
const config = require("../lib/load-config");
|
||||
const BlobStorage = require("../lib/storage/blob");
|
||||
const Store = require("../lib/storage/file");
|
||||
const Fs = require('node:fs');
|
||||
const Quota = require("../lib/commands/quota");
|
||||
const Environment = require('../lib/env');
|
||||
const Env = Environment.create(config);
|
||||
|
||||
const CSV = true;
|
||||
|
||||
config.logPath = false;
|
||||
config.logToStdout = true;
|
||||
|
||||
const start = () => {
|
||||
let time = +new Date();
|
||||
let Log = {};
|
||||
let all = {};
|
||||
let blobStore, store;
|
||||
nThen(w => {
|
||||
Logger.create(config, w(function (_log) {
|
||||
Env.Log = Log = _log;
|
||||
}));
|
||||
}).nThen(w => {
|
||||
config.getSession = function () {};
|
||||
Store.create(config, w(function (err, _store) {
|
||||
if (err) {
|
||||
w.abort();
|
||||
return void Log.error("ERR_PAD_STORE", err);
|
||||
}
|
||||
store = _store;
|
||||
}));
|
||||
BlobStorage.create(config, w(function (err, _store) {
|
||||
if (err) {
|
||||
w.abort();
|
||||
return void Log.error("ERR_BLOB_STORE", err);
|
||||
}
|
||||
blobStore = _store;
|
||||
}));
|
||||
}).nThen(w => {
|
||||
Quota.updateCachedLimits(Env, w((err) => {
|
||||
if (err) {
|
||||
return Env.Log.warn('UPDATE_QUOTA_ERR', err);
|
||||
}
|
||||
Env.Log.info('QUOTA_UPDATED', {});
|
||||
}));
|
||||
}).nThen(w => {
|
||||
Env.Log.info('START_LOADING_PINS');
|
||||
const handlePinLog = (content, id, next) => {
|
||||
const sema = Semaphore.create(20);
|
||||
const data = all[id] = {
|
||||
size: 0,
|
||||
n_pads: 0,
|
||||
n_blobs: 0,
|
||||
n_total: 0,
|
||||
first: content.first,
|
||||
last: content.latest
|
||||
};
|
||||
if (!content.user) {
|
||||
data.maybeTeam = true;
|
||||
}
|
||||
nThen(ww => {
|
||||
Object.keys(content.pins).forEach(id => {
|
||||
sema.take(give => {
|
||||
let addSize = ww(give((e, s) => {
|
||||
if (typeof(s) !== "number") {
|
||||
return; // XXX
|
||||
}
|
||||
data.size += s;
|
||||
data.n_total++;
|
||||
if (id.length === 32) {
|
||||
data.n_pads++;
|
||||
} else {
|
||||
data.n_blobs++;
|
||||
}
|
||||
}));
|
||||
if (id.length === 32) { // PAD
|
||||
return store.getChannelSize(id, addSize);
|
||||
}
|
||||
blobStore.size(id, addSize);
|
||||
});
|
||||
});
|
||||
}).nThen(() => {
|
||||
let key = id.replace(/-/g, '/');
|
||||
if (Env.limits[key]) {
|
||||
let sub = Env.limits[key];
|
||||
data.premium = sub?.plan;
|
||||
}
|
||||
Env.Log.info('PIN_LOG_HANDLED', key);
|
||||
next();
|
||||
});
|
||||
};
|
||||
|
||||
Pins.load(w(() => {
|
||||
let duration = +new Date() - time;
|
||||
Env.Log.info('ALL_PINS_LOADED', duration);
|
||||
}), {
|
||||
pinPath: config.pinPath,
|
||||
handler: handlePinLog,
|
||||
});
|
||||
}).nThen(() => {
|
||||
if (!CSV) { return console.log(all); }
|
||||
let csv = `"User key","Premium plan","Bytes","Number pads","Number blobs","First activity","Last activity","May be a team"\n`;
|
||||
Object.keys(all).sort((a,b) => {
|
||||
return all[b].size - all[a].size;
|
||||
}).forEach(k => {
|
||||
const data = all[k];
|
||||
k = k.replace(/-/g, '/');
|
||||
let first = new Date(data.first).toISOString().slice(0,10);
|
||||
let last = new Date(data.last).toISOString().slice(0,10);
|
||||
let plan = data.premium || '';
|
||||
let t = String(!!data.maybeTeam);
|
||||
csv += `"${k}","${plan}","${data.size}","${data.n_pads}","${data.n_blobs}","${first}","${last}","${t}"\n`;
|
||||
});
|
||||
let filename = `../${new Date().toISOString().slice(0,10)}-stats.csv`;
|
||||
Fs.writeFile(filename, csv, err => {
|
||||
if (err) {
|
||||
console.error(err);
|
||||
} else {
|
||||
console.log('CSV available at', filename);
|
||||
}
|
||||
});
|
||||
});
|
||||
};
|
||||
start();
|
||||
10
server.js
10
server.js
@ -190,7 +190,15 @@ nThen(function (w) {
|
||||
|
||||
var throttledEnvChange = Util.throttle(function () {
|
||||
Env.Log.info('WORKER_ENV_UPDATE', 'Updating HTTP workers with latest state');
|
||||
broadcast('ENV_UPDATE', Environment.serialize(Env));
|
||||
let serialized = Environment.serialize(Env);
|
||||
broadcast('ENV_UPDATE', serialized);
|
||||
if (Env.broadcastWorkerCommand) {
|
||||
Env.broadcastWorkerCommand({
|
||||
command: 'ENV_UPDATE',
|
||||
value: serialized,
|
||||
txid: Util.uid()
|
||||
});
|
||||
}
|
||||
}, 250); // NOTE: changing this value will impact lib/commands/admin-rpc.js#adminDecree callback
|
||||
|
||||
var throttledCacheFlush = Util.throttle(function () {
|
||||
|
||||
Loading…
Reference in New Issue
Block a user