From 9be21c21c46bb65bf92c279f33014f1cd7b80afa Mon Sep 17 00:00:00 2001 From: yflory Date: Tue, 29 Aug 2023 15:19:34 +0200 Subject: [PATCH 1/8] WIP blob metadata --- lib/commands/metadata.js | 6 ++++- lib/hk-util.js | 2 ++ lib/storage/blob.js | 50 +++++++++++++++++++++++++++++++++++++++- lib/workers/db-worker.js | 1 + 4 files changed, 57 insertions(+), 2 deletions(-) diff --git a/lib/commands/metadata.js b/lib/commands/metadata.js index 9b6a23e02..70a811664 100644 --- a/lib/commands/metadata.js +++ b/lib/commands/metadata.js @@ -10,7 +10,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 @@ -76,6 +77,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'); } @@ -137,6 +139,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); @@ -152,6 +155,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/hk-util.js b/lib/hk-util.js index 455964ccb..472e9efb1 100644 --- a/lib/hk-util.js +++ b/lib/hk-util.js @@ -40,6 +40,8 @@ const ADMIN_CHANNEL_LENGTH = HK.ADMIN_CHANNEL_LENGTH = 33; // with a 34 character id const EPHEMERAL_CHANNEL_LENGTH = HK.EPHEMERAL_CHANNEL_LENGTH = 34; +const BLOB_ID_LENGTH = 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/storage/blob.js b/lib/storage/blob.js index 7d55676d7..9126c761c 100644 --- a/lib/storage/blob.js +++ b/lib/storage/blob.js @@ -7,6 +7,11 @@ var BlobStore = module.exports; var nThen = require("nthen"); var Semaphore = require("saferphore"); var Util = require("../common-util"); +var Meta = require("../metadata"); + +const BatchRead = require("../batch-read"); +const readFileBin = require("../stream-file").readFileBin; +const Schedule = require("../schedule"); var isValidSafeKey = function (safeKey) { return typeof(safeKey) === 'string' && !/\//.test(safeKey) && safeKey.length === 44; @@ -26,11 +31,16 @@ 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); }; +// /blob//.metadata.ndjson +var mkMetadataPath = function (env, channelId) { + return Path.join(env.root, channelId.slice(0, 2), channelId) + '.metadata.ndjson'; +}; + // /blobstate// var makeStagePath = function (Env, safeKey) { return Path.join(Env.blobStagingPath, safeKey.slice(0, 2), safeKey); @@ -355,6 +365,43 @@ var restoreProof = function (Env, safeKey, blobId, cb) { Fse.move(archivePath, proofPath, cb); }; +var getDedicatedMetadata = 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); + }); +}; +/* readMetadata + Load the log of metadata amendments. +*/ +var readMetadata = function (Env, blobId, handler, cb) { + getDedicatedMetadata(env, channelId, handler, function (err) { + if (err) { + // stream errors? + return void cb(err); + } + cb(); + }); +}; + var makeWalker = function (n, handleChild, done) { if (!n || typeof(n) !== 'number' || n < 2) { n = 2; } @@ -486,6 +533,7 @@ BlobStore.create = function (config, _cb) { archivePath: config.archivePath || './data/archive', getSession: config.getSession, }; + var schedule = Env.schedule = Schedule(); nThen(function (w) { var CB = Util.both(w.abort, cb); diff --git a/lib/workers/db-worker.js b/lib/workers/db-worker.js index d0b460b11..adf21456a 100644 --- a/lib/workers/db-worker.js +++ b/lib/workers/db-worker.js @@ -307,6 +307,7 @@ const computeIndex = function (data, cb) { const computeMetadata = function (data, cb) { const ref = {}; const lineHandler = Meta.createLineHandler(ref, Env.Log.error); + // XXX BLOBMD use correct store return void store.readChannelMetadata(data.channel, lineHandler, function (err) { if (err) { // stream errors? From 90c95beefee96de3ea61e112dfd5833e1ecc7c2b Mon Sep 17 00:00:00 2001 From: yflory Date: Tue, 17 Dec 2024 16:05:28 +0100 Subject: [PATCH 2/8] Fix issues --- lib/archive-account.js | 18 ++- lib/eviction.js | 59 ---------- lib/storage/blob.js | 245 +++++++++++++++------------------------ lib/workers/db-worker.js | 6 - 4 files changed, 107 insertions(+), 221 deletions(-) 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/eviction.js b/lib/eviction.js index 95f430cdf..8975e3fdd 100644 --- a/lib/eviction.js +++ b/lib/eviction.js @@ -642,64 +642,6 @@ module.exports = function (Env, cb) { })); }; - 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()); - })); - }; - var archiveInactiveChannels = function (w) { var channels = 0; var archived = 0; @@ -802,7 +744,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/storage/blob.js b/lib/storage/blob.js index 8c6b9b2be..629dfd9d3 100644 --- a/lib/storage/blob.js +++ b/lib/storage/blob.js @@ -54,23 +54,10 @@ 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(); } @@ -239,6 +226,11 @@ var archiveMetadata = (Env, blobId, cb) => { // 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 restoreActivity = 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}); @@ -403,7 +395,6 @@ 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 metadata log with an owner @@ -465,19 +456,6 @@ 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); @@ -499,29 +477,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; } @@ -553,7 +513,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. @@ -586,46 +546,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) { @@ -633,32 +553,91 @@ 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); - } - 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; + var metadataName; + + // if the current file is not the channel data, then it must be metadata + if (!/^[0-9a-fA-F]{48}$/.test(item)) { + metadataName = 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; + metadataName = blobName + '.metadata.ndjson'; + } + 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(); }); }; @@ -741,12 +720,6 @@ BlobStore.create = function (config, _cb) { upload_cancel(Env, safeKey, fileSize, cb); }, - isOwnedBy: function (safeKey, blobId, _cb) { - var cb = Util.once(Util.mkAsync(_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); @@ -762,24 +735,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)); @@ -793,12 +754,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: { @@ -807,12 +762,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) { @@ -865,22 +814,14 @@ BlobStore.create = function (config, _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 ddd49ccc0..d1a0fe962 100644 --- a/lib/workers/db-worker.js +++ b/lib/workers/db-worker.js @@ -570,12 +570,6 @@ const removeOwnedBlob = function (data, cb) { nThen(function (w) { // check if you have permissions - blobStore.isOwnedBy(safeKey, blobId, w(function (err, owned) { - if (err || !owned) { - w.abort(); - return void cb("INSUFFICIENT_PERMISSIONS"); - } - })); computeMetadata({channel: blobId}, w((err, meta) => { if (err || !meta) { w.abort(); From 1e2f57a28f30dfc1280cc443b355534bf1391ddf Mon Sep 17 00:00:00 2001 From: yflory Date: Tue, 17 Dec 2024 17:03:27 +0100 Subject: [PATCH 3/8] Blob metadata --- lib/eviction.js | 85 +++++++++++++++++---------------------------- lib/storage/blob.js | 9 +++++ 2 files changed, 40 insertions(+), 54 deletions(-) diff --git a/lib/eviction.js b/lib/eviction.js index 8975e3fdd..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,34 +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()); - })); + }), true); }; var archiveInactiveChannels = function (w) { diff --git a/lib/storage/blob.js b/lib/storage/blob.js index 629dfd9d3..86b3fbeeb 100644 --- a/lib/storage/blob.js +++ b/lib/storage/blob.js @@ -553,6 +553,10 @@ var getActivityStat = function (path, base, cb) { cb(err, stats); }); }; +var getStats = function (Env, blobId, cb) { + var path = makeBlobPath(Env, blobId); + getActivityStat(path, false, cb); +}; let blobRegex = /^[0-9a-fA-F]{48}(\.metadata)*(\.ndjson)*$/; var listBlobs = function (root, handler, fast, cb) { @@ -812,6 +816,11 @@ 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, fast) { From ededf142d93b4b959b527e2c11a4c8848b7201d4 Mon Sep 17 00:00:00 2001 From: yflory Date: Mon, 6 Jan 2025 16:52:59 +0100 Subject: [PATCH 4/8] Blob metadata and migration --- lib/storage/blob.js | 8 +- scripts/migrations/migrate-blob-proofs.js | 164 ++++++++++++++++++++++ 2 files changed, 171 insertions(+), 1 deletion(-) create mode 100644 scripts/migrations/migrate-blob-proofs.js diff --git a/lib/storage/blob.js b/lib/storage/blob.js index 86b3fbeeb..63932a2ce 100644 --- a/lib/storage/blob.js +++ b/lib/storage/blob.js @@ -226,7 +226,7 @@ var archiveMetadata = (Env, blobId, cb) => { // 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 restoreActivity = function (Env, blobId, cb) { +var restoreMetadata = function (Env, blobId, cb) { var path = mkMetadataPath(Env, blobId); var archivePath = prependArchive(Env, path); Fse.move(archivePath, path, cb); @@ -732,6 +732,12 @@ BlobStore.create = function (config, _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) { diff --git a/scripts/migrations/migrate-blob-proofs.js b/scripts/migrations/migrate-blob-proofs.js new file mode 100644 index 000000000..8c977693c --- /dev/null +++ b/scripts/migrations/migrate-blob-proofs.js @@ -0,0 +1,164 @@ +// SPDX-FileCopyrightText: 2025 XWiki CryptPad Team and contributors +// +// SPDX-License-Identifier: AGPL-3.0-or-later + +const Path = require('node:path'); +const nThen = require("nthen"); +const Semaphore = require("saferphore"); +const Logger = require("../../lib/log"); +const config = require("../../lib/load-config"); +const BlobStorage = require("../../lib/storage/blob"); +const Fs = require('node:fs'); + + +const blobPath = config.blobPath || './blob'; +let Log = {}; + +// XXX NOTE: in cleaning mode, we DON'T migrate +// (we suppose data has already been migrated) +const DRY_RUN = true; +const CLEAN_OLD = false; + +const start = (clean) => { + 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"); + process.exit(0); + }); + }); +}; + +start(CLEAN_OLD); From 267f6c56d48b583aa4b0a73b77a8f8012e0c9f25 Mon Sep 17 00:00:00 2001 From: yflory Date: Fri, 10 Jan 2025 15:49:31 +0100 Subject: [PATCH 5/8] User stats script --- lib/pins.js | 9 ++- scripts/user-statistics.js | 133 +++++++++++++++++++++++++++++++++++++ 2 files changed, 141 insertions(+), 1 deletion(-) create mode 100644 scripts/user-statistics.js 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/scripts/user-statistics.js b/scripts/user-statistics.js new file mode 100644 index 000000000..bc6b16227 --- /dev/null +++ b/scripts/user-statistics.js @@ -0,0 +1,133 @@ +// SPDX-FileCopyrightText: 2025 XWiki CryptPad Team and contributors +// +// SPDX-License-Identifier: AGPL-3.0-or-later + +const Path = require('node:path'); +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, pinStore, 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(); From 18e8f057cb2661389bca517ac307af437e777c9c Mon Sep 17 00:00:00 2001 From: yflory Date: Fri, 10 Jan 2025 16:54:37 +0100 Subject: [PATCH 6/8] Start blob proofs migration automatically --- lib/api.js | 25 ++++++++++++ lib/commands/admin-rpc.js | 2 +- lib/decrees.js | 8 ++++ lib/storage/blob.js | 13 ++++++ scripts/migrations/migrate-blob-proofs.js | 50 +++++++++++++++++++---- 5 files changed, 89 insertions(+), 9 deletions(-) diff --git a/lib/api.js b/lib/api.js index df3206b1e..d8313521c 100644 --- a/lib/api.js +++ b/lib/api.js @@ -27,6 +27,31 @@ nThen(function (w) { console.error(err); } })); +}).nThen(function (w) { + 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/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/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/storage/blob.js b/lib/storage/blob.js index 63932a2ce..f3d44c5d7 100644 --- a/lib/storage/blob.js +++ b/lib/storage/blob.js @@ -40,6 +40,7 @@ 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'; }; @@ -116,6 +117,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) { diff --git a/scripts/migrations/migrate-blob-proofs.js b/scripts/migrations/migrate-blob-proofs.js index 8c977693c..c8f87bb3a 100644 --- a/scripts/migrations/migrate-blob-proofs.js +++ b/scripts/migrations/migrate-blob-proofs.js @@ -2,24 +2,24 @@ // // 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 config = require("../../lib/load-config"); const BlobStorage = require("../../lib/storage/blob"); -const Fs = require('node:fs'); +let config = require("../../lib/load-config"); const blobPath = config.blobPath || './blob'; let Log = {}; -// XXX NOTE: in cleaning mode, we DON'T migrate +// NOTE: in cleaning mode, we DON'T migrate // (we suppose data has already been migrated) -const DRY_RUN = true; -const CLEAN_OLD = false; +const start = (clean, dry, cb) => { + const DRY_RUN = dry; -const start = (clean) => { let dirList = []; let blobStore; nThen(w => { @@ -156,9 +156,43 @@ const start = (clean) => { }); n(() => { Log.info("DONE"); - process.exit(0); + cb(); }); }); }; -start(CLEAN_OLD); +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); + }); +} From 5fbcd78ab3b069dac57ab2eeb294c06b92218bb2 Mon Sep 17 00:00:00 2001 From: yflory Date: Fri, 10 Jan 2025 17:23:06 +0100 Subject: [PATCH 7/8] Fallback to owners proofs during migration --- lib/env.js | 1 + lib/storage/blob.js | 5 +++++ lib/workers/db-worker.js | 20 ++++++++++++++++++++ lib/workers/index.js | 5 ++++- server.js | 10 +++++++++- 5 files changed, 39 insertions(+), 2 deletions(-) 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/storage/blob.js b/lib/storage/blob.js index f3d44c5d7..62719cefe 100644 --- a/lib/storage/blob.js +++ b/lib/storage/blob.js @@ -737,6 +737,11 @@ BlobStore.create = function (config, _cb) { upload_cancel(Env, safeKey, fileSize, cb); }, + isOwnedBy: function (safeKey, blobId, _cb) { + var cb = Util.once(Util.mkAsync(_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); diff --git a/lib/workers/db-worker.js b/lib/workers/db-worker.js index d1a0fe962..9a4a39903 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 = []; @@ -576,6 +582,16 @@ const removeOwnedBlob = function (data, cb) { 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"); @@ -693,6 +709,7 @@ const getLastChannelTime = function (data, cb) { }; const COMMANDS = { + ENV_UPDATE: updateEnv, COMPUTE_INDEX: computeIndex, COMPUTE_METADATA: computeMetadata, GET_OLDER_HISTORY: getOlderHistory, @@ -848,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/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 () { From 7186f3e2ef2aa222fae06474f83960dd70693693 Mon Sep 17 00:00:00 2001 From: yflory Date: Fri, 10 Jan 2025 17:28:05 +0100 Subject: [PATCH 8/8] lint compliance --- lib/api.js | 2 +- lib/hk-util.js | 2 +- lib/storage/blob.js | 10 ---------- lib/workers/db-worker.js | 2 +- scripts/migrations/migrate-blob-proofs.js | 4 ++-- scripts/user-statistics.js | 3 +-- 6 files changed, 6 insertions(+), 17 deletions(-) diff --git a/lib/api.js b/lib/api.js index d8313521c..d613df78e 100644 --- a/lib/api.js +++ b/lib/api.js @@ -27,7 +27,7 @@ nThen(function (w) { console.error(err); } })); -}).nThen(function (w) { +}).nThen(function () { if (Env.proofsMigrated) { return; } const { Worker } = require('node:worker_threads'); const Admin = require("./commands/admin-rpc"); diff --git a/lib/hk-util.js b/lib/hk-util.js index 3355f1c09..0ee685bba 100644 --- a/lib/hk-util.js +++ b/lib/hk-util.js @@ -42,7 +42,7 @@ const ADMIN_CHANNEL_LENGTH = HK.ADMIN_CHANNEL_LENGTH = 33; // with a 34 character id const EPHEMERAL_CHANNEL_LENGTH = HK.EPHEMERAL_CHANNEL_LENGTH = 34; -const BLOB_ID_LENGTH = HK.BLOB_ID_LENGTH = 48; +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/storage/blob.js b/lib/storage/blob.js index 62719cefe..a7a053038 100644 --- a/lib/storage/blob.js +++ b/lib/storage/blob.js @@ -10,7 +10,6 @@ var BlobStore = module.exports; var nThen = require("nthen"); var Semaphore = require("saferphore"); var Util = require("../common-util"); -var Meta = require("../metadata"); const PERMISSIVE = 511; const readFileBin = require("../stream-file").readFileBin; @@ -443,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); } @@ -606,11 +599,9 @@ var listBlobs = function (root, handler, fast, cb) { var isLonelyMetadata = false; var blobName; - var metadataName; // if the current file is not the channel data, then it must be metadata if (!/^[0-9a-fA-F]{48}$/.test(item)) { - metadataName = item; blobName = item.replace(/\.metadata\.ndjson/, ''); // check if blob already exists if (list.indexOf(blobName) !== -1) { return; } @@ -619,7 +610,6 @@ var listBlobs = function (root, handler, fast, cb) { isLonelyMetadata = true; } else { blobName = item; - metadataName = blobName + '.metadata.ndjson'; } if (blobName.length !== 48) { return; } diff --git a/lib/workers/db-worker.js b/lib/workers/db-worker.js index 9a4a39903..88fcd5989 100644 --- a/lib/workers/db-worker.js +++ b/lib/workers/db-worker.js @@ -589,7 +589,7 @@ const removeOwnedBlob = function (data, cb) { w.abort(); return void cb("INSUFFICIENT_PERMISSIONS"); } - })) + })); return; } if (!owners || !owners.includes(unsafeKey)) { diff --git a/scripts/migrations/migrate-blob-proofs.js b/scripts/migrations/migrate-blob-proofs.js index c8f87bb3a..d38c94bb0 100644 --- a/scripts/migrations/migrate-blob-proofs.js +++ b/scripts/migrations/migrate-blob-proofs.js @@ -101,7 +101,7 @@ const start = (clean, dry, cb) => { err }); } - let blobs = [] + let blobs = []; nThen(www => { list.forEach(dir => { // ./blob/abc/abcdefg.../01 @@ -142,7 +142,7 @@ const start = (clean, dry, cb) => { })); })); - }) + }); }); }).nThen(ww(give(() => { Log.info("END_USER", edPublic); diff --git a/scripts/user-statistics.js b/scripts/user-statistics.js index bc6b16227..02a74ef89 100644 --- a/scripts/user-statistics.js +++ b/scripts/user-statistics.js @@ -2,7 +2,6 @@ // // SPDX-License-Identifier: AGPL-3.0-or-later -const Path = require('node:path'); const nThen = require("nthen"); const Semaphore = require("saferphore"); const Logger = require("../lib/log"); @@ -24,7 +23,7 @@ const start = () => { let time = +new Date(); let Log = {}; let all = {}; - let blobStore, pinStore, store; + let blobStore, store; nThen(w => { Logger.create(config, w(function (_log) { Env.Log = Log = _log;