diff --git a/lib/api.js b/lib/api.js index df3206b1e..d613df78e 100644 --- a/lib/api.js +++ b/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 || []; diff --git a/lib/archive-account.js b/lib/archive-account.js index 43ff1357e..44928f6be 100644 --- a/lib/archive-account.js +++ b/lib/archive-account.js @@ -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 diff --git a/lib/commands/admin-rpc.js b/lib/commands/admin-rpc.js index 9030159f8..d029fc079 100644 --- a/lib/commands/admin-rpc.js +++ b/lib/commands/admin-rpc.js @@ -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'); } diff --git a/lib/commands/metadata.js b/lib/commands/metadata.js index 8eeca31ce..a468f6f6e 100644 --- a/lib/commands/metadata.js +++ b/lib/commands/metadata.js @@ -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; diff --git a/lib/decrees.js b/lib/decrees.js index 0124d9d90..c52476865 100644 --- a/lib/decrees.js +++ b/lib/decrees.js @@ -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"); diff --git a/lib/env.js b/lib/env.js index d3748750f..a3593c7a1 100644 --- a/lib/env.js +++ b/lib/env.js @@ -415,6 +415,7 @@ const BAD = [ 'limits', 'customLimits', 'scheduleDecree', + 'plugins', 'httpServer', diff --git a/lib/eviction.js b/lib/eviction.js index 95f430cdf..b75fc78a6 100644 --- a/lib/eviction.js +++ b/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(); diff --git a/lib/hk-util.js b/lib/hk-util.js index e9d177968..0ee685bba 100644 --- a/lib/hk-util.js +++ b/lib/hk-util.js @@ -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; diff --git a/lib/pins.js b/lib/pins.js index 2201e561c..d2504077f 100644 --- a/lib/pins.js +++ b/lib/pins.js @@ -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 + } + }; }; diff --git a/lib/storage/blob.js b/lib/storage/blob.js index aeb9c1f70..a7a053038 100644 --- a/lib/storage/blob.js +++ b/lib/storage/blob.js @@ -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//// +// /blob// 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//.metadata.ndjson +var mkMetadataPath = function (Env, blobId) { + return Path.join(Env.blobPath, blobId.slice(0, 2), blobId) + '.metadata.ndjson'; +}; + // /blobstate// var makeStagePath = function (Env, safeKey) { return Path.join(Env.blobStagingPath, safeKey.slice(0, 2), safeKey); }; -// /blob//// -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//// +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); }, } }, diff --git a/lib/workers/db-worker.js b/lib/workers/db-worker.js index 64376236f..88fcd5989 100644 --- a/lib/workers/db-worker.js +++ b/lib/workers/db-worker.js @@ -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; diff --git a/lib/workers/index.js b/lib/workers/index.js index 5bdbd5d7a..bc1780cc5 100644 --- a/lib/workers/index.js +++ b/lib/workers/index.js @@ -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; diff --git a/scripts/migrations/migrate-blob-proofs.js b/scripts/migrations/migrate-blob-proofs.js new file mode 100644 index 000000000..d38c94bb0 --- /dev/null +++ b/scripts/migrations/migrate-blob-proofs.js @@ -0,0 +1,198 @@ +// SPDX-FileCopyrightText: 2025 XWiki CryptPad Team 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); + }); +} diff --git a/scripts/user-statistics.js b/scripts/user-statistics.js new file mode 100644 index 000000000..02a74ef89 --- /dev/null +++ b/scripts/user-statistics.js @@ -0,0 +1,132 @@ +// SPDX-FileCopyrightText: 2025 XWiki CryptPad Team 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(); diff --git a/server.js b/server.js index 407f10f43..ca2b2ae77 100644 --- a/server.js +++ b/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 () {