From 9be21c21c46bb65bf92c279f33014f1cd7b80afa Mon Sep 17 00:00:00 2001 From: yflory Date: Tue, 29 Aug 2023 15:19:34 +0200 Subject: [PATCH 01/53] 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 3940b9c4be40402d50ed7de0fce49a2cea75a2cd Mon Sep 17 00:00:00 2001 From: yflory Date: Wed, 15 May 2024 11:17:25 +0200 Subject: [PATCH 02/53] Fix checkpoint issues with integrated OO apps --- www/common/onlyoffice/inner.js | 14 +++++++++----- www/common/outer/async-store.js | 32 ++++++++++++++++++++++++++------ www/cryptpad-api.js | 2 ++ 3 files changed, 37 insertions(+), 11 deletions(-) diff --git a/www/common/onlyoffice/inner.js b/www/common/onlyoffice/inner.js index 9faa7afc2..5184833c7 100644 --- a/www/common/onlyoffice/inner.js +++ b/www/common/onlyoffice/inner.js @@ -392,6 +392,10 @@ define([ }; var onUploaded = function (ev, data, err) { + if (!ev && err) { + console.error(err); + return void UI.warn(Messages.error); + } if (ev.newTemplate) { if (err) { console.error(err); @@ -549,7 +553,7 @@ define([ var saveToServer = function (blob, title) { if (APP.cantCheckpoint) { return; } // TOO_LARGE - var text = getContent(); + var text = !blob && getContent(); if (!text && !blob) { setEditable(false, true); sframeChan.query('Q_CLEAR_CACHE_CHANNELS', [ @@ -3154,9 +3158,9 @@ Uncaught TypeError: Cannot read property 'calculatedType' of null return void UI.errorLoadingScreen(Messages.error); } var blob = new Blob([bin], {type: 'text/plain'}); - var file = getFileType(); - resetData(blob, file); - //saveToServer(blob, title); + //var file = getFileType(); + //resetData(blob, file); + saveToServer(blob, title); Title.updateTitle(title); UI.removeLoadingScreen(); }); @@ -3205,7 +3209,7 @@ Uncaught TypeError: Cannot read property 'calculatedType' of null cb(); }); }); - if (privateData.initialState) { + if (privateData.initialState && (!content || !content.hashes)) { var blob = privateData.initialState; let title = `document.${cfg.fileType}`; console.error(blob, title); diff --git a/www/common/outer/async-store.js b/www/common/outer/async-store.js index cd93675fe..9ba3853e0 100644 --- a/www/common/outer/async-store.js +++ b/www/common/outer/async-store.js @@ -467,6 +467,19 @@ define([ }); }; + var initTempRpc = (clientId, cb) => { + if (store.rpc) { return void cb(store.rpc); } + var kp = Crypto.Nacl.sign.keyPair(); + var keys = { + edPublic: Crypto.Nacl.util.encodeBase64(kp.publicKey), + edPrivate: Crypto.Nacl.util.encodeBase64(kp.secretKey) + }; + Pinpad.create(store.network, keys, function (e, call) { + if (e) { return void cb({error: e}); } + store.rpc = call; + cb(call); + }); + }; var initRpc = function (clientId, data, cb) { if (!store.loggedIn) { return cb(); } if (store.rpc) { return void cb(account); } @@ -3149,7 +3162,7 @@ define([ // If we load CryptPad for the first time from an existing pad, don't create a // drive automatically. - var onNoDrive = function (clientId, cb) { + var onNoDrive = function (clientId, cb, initRpc) { var andThen = function () { // To be able to use all the features inside the pad, we need to // initialize the chat (messenger) and the cursor modules. @@ -3160,9 +3173,16 @@ define([ store.messenger = store.modules['messenger']; // And now we're ready - initAnonRpc(null, null, function () { - cb({}); - }); + let getAnon = () => { + initAnonRpc(null, null, function () { + cb({}); + }); + }; + + if (initRpc) { + return initTempRpc(clientId, getAnon); + } + getAnon(); }; // We need an anonymous RPC to be able to check if the pad exists and to get @@ -3245,7 +3265,7 @@ define([ // First tab, no user hash, no anon hash and this app doesn't need a drive // ==> don't create a drive // Or "neverDrive" (integration into another platform?) - // ==> don't create a drive + // ==> don't create a drive BUT create temp RPC (we may need to upload) if (data.neverDrive || (data.noDrive && !data.userHash && !data.anonHash)) { return void onNoDrive(clientId, function (obj) { if (obj && obj.error) { @@ -3259,7 +3279,7 @@ define([ } Feedback.send("NO_DRIVE", true); callback(obj); - }); + }, !!data.neverDrive); } initialized = true; diff --git a/www/cryptpad-api.js b/www/cryptpad-api.js index e73ef40bd..cdbf7f8dd 100644 --- a/www/cryptpad-api.js +++ b/www/cryptpad-api.js @@ -188,6 +188,7 @@ chan.on('ON_DOWNLOADAS', blob => { let url = URL.createObjectURL(blob); + if (!config.events.onDownloadAs) { return; } config.events.onDownloadAs({ data: { fileType: config.document && config.document.fileType, @@ -198,6 +199,7 @@ chan.on('SAVE', function (data, cb) { blob = data; + if (!config.events.onSave) { return void cb(); } config.events.onSave(data, cb); }); chan.on('RELOAD', function () { From 154b3e66c84d3611e628e3925928282b4d2bd693 Mon Sep 17 00:00:00 2001 From: yflory Date: Wed, 15 May 2024 11:37:03 +0200 Subject: [PATCH 03/53] Get username and lang from API config --- customize.dist/messages.js | 7 +++---- www/common/onlyoffice/inner.js | 8 ++++++-- www/common/sframe-app-outer.js | 4 +++- www/cryptpad-api.js | 3 ++- www/integration/main.js | 10 ++++++++-- 5 files changed, 22 insertions(+), 10 deletions(-) diff --git a/customize.dist/messages.js b/customize.dist/messages.js index 579dcd31e..0fb495bb2 100755 --- a/customize.dist/messages.js +++ b/customize.dist/messages.js @@ -35,11 +35,10 @@ var getStoredLanguage = function () { return localStorage && localStorage.getIte var getBrowserLanguage = function () { return navigator.language || navigator.userLanguage || ''; }; var getLanguage = Messages._getLanguage = function () { if (window.cryptpadLanguage) { return window.cryptpadLanguage; } - try { - if (getStoredLanguage()) { return getStoredLanguage(); } - } catch (e) { console.log(e); } var l = getBrowserLanguage(); - // Edge returns 'fr-FR' --> transform it to 'fr' and check again + try { + l = getStoredLanguage() || getBrowserLanguage(); + } catch (e) { console.log(e); } return map[l] ? l : (map[l.split('-')[0]] ? l.split('-')[0] : (map[l.split('_')[0]] ? l.split('_')[0] : 'en')); diff --git a/www/common/onlyoffice/inner.js b/www/common/onlyoffice/inner.js index 5184833c7..dfa949eb6 100644 --- a/www/common/onlyoffice/inner.js +++ b/www/common/onlyoffice/inner.js @@ -1671,6 +1671,10 @@ define([ var lang = (window.cryptpadLanguage || navigator.language || navigator.userLanguage || '').slice(0,2); + let username = Util.find(privateData, ['integrationConfig', 'user', 'name']) + || metadataMgr.getUserData().name + || Messages.anonymous; + // Config APP.ooconfig = { "document": { @@ -1694,8 +1698,8 @@ define([ }, "user": { "id": String(myOOId), //"c0c3bf82-20d7-4663-bf6d-7fa39c598b1d", - "firstname": metadataMgr.getUserData().name || Messages.anonymous, - "name": metadataMgr.getUserData().name || Messages.anonymous, + "firstname": username, + "name": username }, "mode": "edit", "lang": lang diff --git a/www/common/sframe-app-outer.js b/www/common/sframe-app-outer.js index b1210796f..b634f51b4 100644 --- a/www/common/sframe-app-outer.js +++ b/www/common/sframe-app-outer.js @@ -17,7 +17,9 @@ define([ nThen(function (waitFor) { DomReady.onReady(waitFor()); }).nThen(function (waitFor) { - var obj = SFCommonO.initIframe(waitFor, true, integration.pathname); + let lang = integration && integration.config && integration.config.editorConfig + && integration.config.editorConfig.lang; + var obj = SFCommonO.initIframe(waitFor, true, integration.pathname, lang); href = obj.href; hash = obj.hash; if (isIntegration) { diff --git a/www/cryptpad-api.js b/www/cryptpad-api.js index cdbf7f8dd..d7d2ac7e2 100644 --- a/www/cryptpad-api.js +++ b/www/cryptpad-api.js @@ -116,7 +116,8 @@ application: config.documentType, document: blob, ext: config.document.fileType, - autosave: config.autosave || 10 + autosave: config.autosave || 10, + editorConfig: config.editorConfig || {} }, function (obj) { if (obj && obj.error) { reject(obj.error); return console.error(obj.error); } resolve({}); diff --git a/www/integration/main.js b/www/integration/main.js index fe3ffd010..a09592d3b 100644 --- a/www/integration/main.js +++ b/www/integration/main.js @@ -155,7 +155,12 @@ define([ chan.on('START', function (data) { console.warn('INNER START', data); var href = Hash.hashToHref(data.key, data.application); - console.error(Hash.hrefToHexChannelId(href)); + + if (data.editorConfig.lang) { + var LS_LANG = "CRYPTPAD_LANG"; + localStorage.setItem(LS_LANG, data.editorConfig.lang); + } + window.CP_integration_outer = { pathname: `/${data.application}/`, hash: data.key, @@ -163,7 +168,8 @@ define([ initialState: data.document, config: { fileType: data.ext, - autosave: data.autosave + autosave: data.autosave, + user: data.editorConfig.user }, utils: { onReady: onReady, From 112a222f79f7bf82db96404b1ce73844789494cc Mon Sep 17 00:00:00 2001 From: yflory Date: Tue, 21 May 2024 18:05:03 +0200 Subject: [PATCH 04/53] Fix initial OO content issue with API --- www/common/onlyoffice/inner.js | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/www/common/onlyoffice/inner.js b/www/common/onlyoffice/inner.js index dfa949eb6..9ffbfd47b 100644 --- a/www/common/onlyoffice/inner.js +++ b/www/common/onlyoffice/inner.js @@ -3213,10 +3213,9 @@ Uncaught TypeError: Cannot read property 'calculatedType' of null cb(); }); }); - if (privateData.initialState && (!content || !content.hashes)) { + if (privateData.initialState && (!content || !content.hashes || !Object.keys(content.hashes).length)) { var blob = privateData.initialState; let title = `document.${cfg.fileType}`; - console.error(blob, title); return convertImportBlob(blob, title); } } From 97a806353e911efb3e1c3bc1ce8dca3c33c5d0aa Mon Sep 17 00:00:00 2001 From: yflory Date: Mon, 24 Jun 2024 16:18:27 +0200 Subject: [PATCH 05/53] Fix issue with application names --- www/cryptpad-api.js | 14 ++++++++++---- 1 file changed, 10 insertions(+), 4 deletions(-) diff --git a/www/cryptpad-api.js b/www/cryptpad-api.js index d7d2ac7e2..c72f41f65 100644 --- a/www/cryptpad-api.js +++ b/www/cryptpad-api.js @@ -251,16 +251,22 @@ cryptpadURL = getInstanceURL(); } + config.events = config.events || {}; + // OnlyOffice shim let url = config.document.url; if (/^http:\/\/localhost\/cache\/files\//.test(url)) { url = url.replace(/(http:\/\/localhost\/cache\/files\/)/, getInstanceURL() + 'ooapi/'); } config.document.url = url; - if (config.documentType === "spreadsheet") { + if (config.documentType === "spreadsheet" || config.documentType === "cell") { config.documentType = "sheet"; } - if (config.documentType === "text") { + + if (config.documentType === "slide") { + config.documentType = "presentation"; + } + if (config.documentType === "word" || config.documentType === "text") { config.documentType = "doc"; } @@ -311,8 +317,8 @@ iframe.setAttribute('name', 'frameEditor'); iframe.setAttribute('align', 'top'); iframe.setAttribute("src", url); - iframe.setAttribute("width", config.width); - iframe.setAttribute("height", config.height); + iframe.setAttribute("width", config.width || '100%'); + iframe.setAttribute("height", config.height || '100%'); if (config.editorConfig) { // OnlyOffice container.replaceWith(iframe); container = iframe; From b62ff674ef0088a0adae75856e84c13f41af446b Mon Sep 17 00:00:00 2001 From: yflory Date: Fri, 6 Sep 2024 11:31:09 +0200 Subject: [PATCH 06/53] API: download document from our server --- www/cryptpad-api.js | 6 +++ www/integration/main.js | 90 +++++++++++++++++++++++++++++------------ 2 files changed, 70 insertions(+), 26 deletions(-) diff --git a/www/cryptpad-api.js b/www/cryptpad-api.js index e73ef40bd..f4a10d147 100644 --- a/www/cryptpad-api.js +++ b/www/cryptpad-api.js @@ -18,7 +18,9 @@ var scripts = document.getElementsByTagName('script'); for (var i = scripts.length - 1; i >= 0; i--) { var match = scripts[i].src.match(/(.*)web-apps\/apps\/api\/documents\/api.js/i); + var match2 = scripts[i].src.match(/(.*)\/cryptpad-api.js/i); if (match) { return match[1]; } + else if (match2) { return match2[1]; } } }; @@ -114,6 +116,7 @@ chan.send('START', { key: key, application: config.documentType, + url: config.document.url, document: blob, ext: config.document.fileType, autosave: config.autosave || 10 @@ -130,6 +133,9 @@ blob = config.document.blob; return start(); } + return start(); + // XXX use server only when not zero knowledge? i.e. no save handler? + // XXX or when error with client? getBlob(function (err, _blob) { if (err) { reject(err); return console.error(err); } _blob.name = `document.${config.document.fileType}`; diff --git a/www/integration/main.js b/www/integration/main.js index fe3ffd010..eb137c039 100644 --- a/www/integration/main.js +++ b/www/integration/main.js @@ -3,9 +3,10 @@ // SPDX-License-Identifier: AGPL-3.0-or-later define([ + '/api/config', '/common/sframe-common-outer.js', '/common/common-hash.js', -], function (SCO, Hash) { +], function (Config, SCO, Hash) { var getTxid = function () { return Math.random().toString(16).replace('0.', ''); @@ -152,36 +153,73 @@ define([ chan.send('ON_DOWNLOADAS', blob); }; - chan.on('START', function (data) { + + let getInstanceURL = function () { + return Config.httpUnsafeOrigin; + }; + let getBlobServer = function (documentURL, cb) { + let xhr = new XMLHttpRequest(); + let data = encodeURIComponent(documentURL); + let url = getInstanceURL() + '/ooapidl?url=' + data; + console.log(url); + xhr.open('GET', url, true); + xhr.responseType = 'blob'; + //xhr.setRequestHeader('Content-Type', 'application/json'); + xhr.onload = function () { + if (this.status === 200) { + var blob = this.response; + // myBlob is now the blob that the object URL pointed to. + cb(null, blob); + } else { + cb(this.status); + } + }; + xhr.onerror = function (e) { + cb(e.message); + }; + xhr.send(); + }; + chan.on('START', function (data, cb) { console.warn('INNER START', data); var href = Hash.hashToHref(data.key, data.application); console.error(Hash.hrefToHexChannelId(href)); - window.CP_integration_outer = { - pathname: `/${data.application}/`, - hash: data.key, - href: href, - initialState: data.document, - config: { - fileType: data.ext, - autosave: data.autosave - }, - utils: { - onReady: onReady, - onDownloadAs, - setDownloadAs, - save: save, - reload: reload, - onHasUnsavedChanges: onHasUnsavedChanges, - onInsertImage: onInsertImage + let startApp = function (blob) { + window.CP_integration_outer = { + pathname: `/${data.application}/`, + hash: data.key, + href: href, + initialState: blob, + config: { + fileType: data.ext, + autosave: data.autosave + }, + utils: { + onReady: onReady, + onDownloadAs, + setDownloadAs, + save: save, + reload: reload, + onHasUnsavedChanges: onHasUnsavedChanges, + onInsertImage: onInsertImage + } + }; + let path = "/common/sframe-app-outer.js"; + if (['sheet', 'doc', 'presentation'].includes(data.application)) { + path = '/common/onlyoffice/main.js'; } + require([path], function () { + console.warn('SAO REQUIRED'); + delete window.CP_integration_outer; + cb(); + }); }; - let path = "/common/sframe-app-outer.js"; - if (['sheet', 'doc', 'presentation'].includes(data.application)) { - path = '/common/onlyoffice/main.js'; - } - require([path], function () { - console.warn('SAO REQUIRED'); - delete window.CP_integration_outer; + + if (data.document) { return void startApp(data.document); } + getBlobServer(data.url, (err, blob) => { + if (err) { + return void cb({error: err}); + } + startApp(blob); }); }); From 062c1550d884079f4ee7e6f2d8719219de4521d3 Mon Sep 17 00:00:00 2001 From: yflory Date: Mon, 30 Sep 2024 17:27:15 +0200 Subject: [PATCH 07/53] Fix issue with non-base64 keys --- www/integration/main.js | 15 +++++++++++++-- 1 file changed, 13 insertions(+), 2 deletions(-) diff --git a/www/integration/main.js b/www/integration/main.js index 4d252f659..1cbf1f241 100644 --- a/www/integration/main.js +++ b/www/integration/main.js @@ -6,8 +6,10 @@ define([ '/api/config', '/common/sframe-common-outer.js', '/common/common-hash.js', + '/components/tweetnacl/nacl-fast.min.js' ], function (Config, SCO, Hash) { + let Nacl = window.nacl; var getTxid = function () { return Math.random().toString(16).replace('0.', ''); }; @@ -96,9 +98,17 @@ define([ }; http.send(); }; + let sanitizeKey = key => { + try { + Nacl.util.decodeBase64(key); + return key; + } catch (e) { + return Nacl.util.encodeBase64(Nacl.util.decodeUTF8(key)); + } + }; chan.on('GET_SESSION', function (data, cb) { - if (data.keepOld) { - var key = data.key + "000000000000000000000000000000000"; + if (data.keepOld) { // they provide their own key, we must turn it into a hash + var key = sanitizeKey(data.key) + "000000000000000000000000000000000"; console.warn('KEY', key); return void cb({ key: `/2/integration/edit/${key.slice(0,24)}/` @@ -181,6 +191,7 @@ define([ }; chan.on('START', function (data, cb) { console.warn('INNER START', data); + // data.key is a hash var href = Hash.hashToHref(data.key, data.application); if (data.editorConfig.lang) { var LS_LANG = "CRYPTPAD_LANG"; From 66047fed03220962f6b98ea34ced9f48a7b4fdd7 Mon Sep 17 00:00:00 2001 From: yflory Date: Tue, 1 Oct 2024 12:31:23 +0200 Subject: [PATCH 08/53] Fix onlyoffice out of sync with integration API --- www/common/outer/async-store.js | 1 + 1 file changed, 1 insertion(+) diff --git a/www/common/outer/async-store.js b/www/common/outer/async-store.js index af20c0722..704038b2c 100644 --- a/www/common/outer/async-store.js +++ b/www/common/outer/async-store.js @@ -2684,6 +2684,7 @@ define([ }; */ var loadOnlyOffice = function () { + if (store.onlyoffice) { return; } store.onlyoffice = OnlyOffice.init(store, function (ev, data, clients) { clients.forEach(function (cId) { postMessage(cId, 'OO_EVENT', { From 8edf6001f2dcc2a5b31dca56ed1d83d3757b22e0 Mon Sep 17 00:00:00 2001 From: yflory Date: Thu, 3 Oct 2024 16:59:55 +0200 Subject: [PATCH 09/53] First version of OO save --- www/common/onlyoffice/inner.js | 12 ++++++ www/common/sframe-common-integration.js | 2 +- www/cryptpad-api.js | 9 ++++- www/integration/main.js | 49 +++++++++++++++++++++---- 4 files changed, 62 insertions(+), 10 deletions(-) diff --git a/www/common/onlyoffice/inner.js b/www/common/onlyoffice/inner.js index a29c10c6b..4e878005b 100644 --- a/www/common/onlyoffice/inner.js +++ b/www/common/onlyoffice/inner.js @@ -3230,11 +3230,23 @@ Uncaught TypeError: Cannot read property 'calculatedType' of null }); } integrationChannel.on('Q_INTEGRATION_NEEDSAVE', function (data, cb) { + if (!cfg.autosave) { return; } integrationSave(function (obj) { if (obj && obj.error) { console.error(obj.error); } cb(); }); }); + + if (!cfg.autosave) { + let $save = common.createButton('save', true, {}, function () { + $save.attr('disabled', 'disabled'); + integrationSave(err => { + $save.removeAttr('disabled'); + }); + }); + $('body').prepend($save); + } + if (privateData.initialState && (!content || !content.hashes || !Object.keys(content.hashes).length)) { var blob = privateData.initialState; let title = `document.${cfg.fileType}`; diff --git a/www/common/sframe-common-integration.js b/www/common/sframe-common-integration.js index 8b1ab5374..41d81a7b9 100644 --- a/www/common/sframe-common-integration.js +++ b/www/common/sframe-common-integration.js @@ -17,7 +17,7 @@ define([ var privateData = metadataMgr.getPrivateData(); var config = privateData.integrationConfig; - if (!config.autosave) { return; } + if (!config.autosave) { return void exp; } if (typeof(saveHandler) !== "function") { throw new Error("Incorrect save handler"); } diff --git a/www/cryptpad-api.js b/www/cryptpad-api.js index 9e873cba3..efe065618 100644 --- a/www/cryptpad-api.js +++ b/www/cryptpad-api.js @@ -88,6 +88,7 @@ var start = function (config, chan) { return new Promise(function (resolve, reject) { setTimeout(function () { + var docID = config.document.key; var key = config.document.key; var blob; @@ -116,10 +117,12 @@ chan.send('START', { key: key, application: config.documentType, + name: config.document.title, url: config.document.url, + documentKey: docID, document: blob, ext: config.document.fileType, - autosave: config.autosave || 10, + autosave: config.events.onSave && (config.autosave || 10), editorConfig: config.editorConfig || {} }, function (obj) { if (obj && obj.error) { reject(obj.error); return console.error(obj.error); } @@ -134,9 +137,11 @@ blob = config.document.blob; return start(); } - return start(); // XXX use server only when not zero knowledge? i.e. no save handler? // XXX or when error with client? + if (!config.events.onSave) { + return start(); + } getBlob(function (err, _blob) { if (err) { reject(err); return console.error(err); } _blob.name = `document.${config.document.fileType}`; diff --git a/www/integration/main.js b/www/integration/main.js index 1cbf1f241..cffd22a33 100644 --- a/www/integration/main.js +++ b/www/integration/main.js @@ -129,12 +129,6 @@ define([ }); }); - var save = function (obj, cb) { - chan.send('SAVE', obj.blob, function (err) { - if (err) { return cb({error: err}); } - cb(); - }); - }; var reload = function (data) { chan.send('RELOAD', data); }; @@ -171,7 +165,6 @@ define([ let xhr = new XMLHttpRequest(); let data = encodeURIComponent(documentURL); let url = getInstanceURL() + '/ooapidl?url=' + data; - console.log(url); xhr.open('GET', url, true); xhr.responseType = 'blob'; //xhr.setRequestHeader('Content-Type', 'application/json'); @@ -189,6 +182,30 @@ define([ }; xhr.send(); }; + let saveBlobServer = function (cfg, blob, cb) { + let {callbackUrl, name, key} = cfg; + let xhr = new XMLHttpRequest(); + name = encodeURIComponent(name); + callbackUrl = encodeURIComponent(callbackUrl); + key = encodeURIComponent(key); + let query = `name=${name}&cb=${callbackUrl}&key=${key}` + let url = getInstanceURL() + `/oosave?${query}`; + xhr.open('POST', url, true); + xhr.responseType = 'blob'; + //xhr.setRequestHeader('Content-Type', 'application/json'); + xhr.onload = function () { + console.error(this.status); + if (this.status === 200) { + cb(); + } else { + cb(this.status); + } + }; + xhr.onerror = function (e) { + cb(e.message); + }; + xhr.send(blob); + }; chan.on('START', function (data, cb) { console.warn('INNER START', data); // data.key is a hash @@ -198,6 +215,23 @@ define([ localStorage.setItem(LS_LANG, data.editorConfig.lang); } + let fileName = data.name || `document.${data.ext}`; + var save = function (obj, cb) { + let cbUrl = data.editorConfig.callbackUrl; + if (!data.autosave && cbUrl) { + saveBlobServer({ + callbackUrl: cbUrl, + name: fileName, + key: data.documentKey + }, obj.blob, cb); + return; + } + chan.send('SAVE', obj.blob, function (err) { + if (err) { return cb({error: err}); } + cb(); + }); + }; + console.error(Hash.hrefToHexChannelId(href)); let startApp = function (blob) { window.CP_integration_outer = { @@ -206,6 +240,7 @@ define([ href: href, initialState: blob, config: { + fileName: data.name, fileType: data.ext, autosave: data.autosave, user: data.editorConfig.user From 235c04f2ad9148d226d270b8753ddcaebc0533e4 Mon Sep 17 00:00:00 2001 From: yflory Date: Tue, 29 Oct 2024 11:20:12 +0100 Subject: [PATCH 10/53] Fix placeholder issues --- lib/hk-util.js | 3 ++- www/cryptpad-api.js | 5 ++++- 2 files changed, 6 insertions(+), 2 deletions(-) diff --git a/lib/hk-util.js b/lib/hk-util.js index 93c3f9c0e..069be04a3 100644 --- a/lib/hk-util.js +++ b/lib/hk-util.js @@ -143,7 +143,8 @@ const dropChannel = HK.dropChannel = function (Env, chanName) { delete Env.channel_cache[chanName]; if (meta && meta.selfdestruct && Env.selfDestructTo) { Env.selfDestructTo[chanName] = setTimeout(function () { - expireChannel(Env, chanName); + if (!Env.store) { return; } + Env.store.archiveChannel(chanName, false, () => {}); }, TEMPORARY_CHANNEL_LIFETIME); } if (Env.store) { Env.store.closeChannel(chanName, function () {}); } diff --git a/www/cryptpad-api.js b/www/cryptpad-api.js index e73ef40bd..9dfcc55b9 100644 --- a/www/cryptpad-api.js +++ b/www/cryptpad-api.js @@ -131,7 +131,10 @@ return start(); } getBlob(function (err, _blob) { - if (err) { reject(err); return console.error(err); } + if (err) { // Can't get blob from client, try from server + console.warn(err); + return void start(); + } _blob.name = `document.${config.document.fileType}`; blob = _blob; start(); From 1fb3dbc2ceb421beddc9d7f4addfde518a945e36 Mon Sep 17 00:00:00 2001 From: yflory Date: Tue, 29 Oct 2024 11:23:17 +0100 Subject: [PATCH 11/53] Fix merge error --- www/cryptpad-api.js | 5 ----- 1 file changed, 5 deletions(-) diff --git a/www/cryptpad-api.js b/www/cryptpad-api.js index f41aad8ec..6fd6d3be4 100644 --- a/www/cryptpad-api.js +++ b/www/cryptpad-api.js @@ -137,11 +137,6 @@ blob = config.document.blob; return start(); } - // XXX use server only when not zero knowledge? i.e. no save handler? - // XXX or when error with client? - if (!config.events.onSave) { - return start(); - } getBlob(function (err, _blob) { if (err) { // Can't get blob from client, try from server console.warn(err); From 4a881555817333798133a18073f49fc3de0a92a0 Mon Sep 17 00:00:00 2001 From: yflory Date: Tue, 29 Oct 2024 11:33:45 +0100 Subject: [PATCH 12/53] Fix expire selfdestruct channel --- lib/historyKeeper.js | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/lib/historyKeeper.js b/lib/historyKeeper.js index 7178b5f76..b8ac2ca58 100644 --- a/lib/historyKeeper.js +++ b/lib/historyKeeper.js @@ -67,7 +67,10 @@ module.exports.create = function (Env, cb) { } if (metadata && metadata.selfdestruct && metadata.selfdestruct !== Env.id) { - HK.expireChannel(Env, channelName); + if (Env.store) { + Env.store.archiveChannel(channelName, false, + () => {}); + } return void cb('ESELFDESTRUCT'); } From 3de5784be76daba5a9d0135ddf11d816e94fe279 Mon Sep 17 00:00:00 2001 From: yflory Date: Thu, 7 Nov 2024 16:23:07 +0100 Subject: [PATCH 13/53] Fix issues with temporary documents --- lib/historyKeeper.js | 5 +---- lib/hk-util.js | 13 ++++++++++--- www/common/outer/async-store.js | 3 ++- 3 files changed, 13 insertions(+), 8 deletions(-) diff --git a/lib/historyKeeper.js b/lib/historyKeeper.js index b8ac2ca58..e0a9fea2c 100644 --- a/lib/historyKeeper.js +++ b/lib/historyKeeper.js @@ -67,10 +67,7 @@ module.exports.create = function (Env, cb) { } if (metadata && metadata.selfdestruct && metadata.selfdestruct !== Env.id) { - if (Env.store) { - Env.store.archiveChannel(channelName, false, - () => {}); - } + HK.removeChannel(Env, channelName); return void cb('ESELFDESTRUCT'); } diff --git a/lib/hk-util.js b/lib/hk-util.js index 069be04a3..2ba656376 100644 --- a/lib/hk-util.js +++ b/lib/hk-util.js @@ -134,6 +134,13 @@ const expireChannel = HK.expireChannel = function (Env, channel) { }); }; +const removeChannel = HK.removeChannel = function (Env, channel) { + if (!Env.store) { return; } + Env.store.archiveChannel(channel, void 0, () => {}); + delete Env.metadata_cache[channel]; + delete Env.channel_cache[channel]; +}; + /* dropChannel * cleans up memory structures which are managed entirely by the historyKeeper */ @@ -143,8 +150,7 @@ const dropChannel = HK.dropChannel = function (Env, chanName) { delete Env.channel_cache[chanName]; if (meta && meta.selfdestruct && Env.selfDestructTo) { Env.selfDestructTo[chanName] = setTimeout(function () { - if (!Env.store) { return; } - Env.store.archiveChannel(chanName, false, () => {}); + removeChannel(Env, chanName); }, TEMPORARY_CHANNEL_LIFETIME); } if (Env.store) { Env.store.closeChannel(chanName, function () {}); } @@ -693,6 +699,7 @@ const handleGetHistory = function (Env, Server, seq, userId, parsed) { }, (err, reason) => { // Any error but ENOENT: abort // ENOENT is allowed in case we want to create a new pad + if (err && err.error) { err = err.error; } if (err && err.code !== 'ENOENT') { if (err.message === "EUNKNOWN") { Log.error("HK_GET_HISTORY", { @@ -708,7 +715,7 @@ const handleGetHistory = function (Env, Server, seq, userId, parsed) { stack: err && err.stack, }); } // FIXME err.message isn't useful for users - const parsedMsg = {error:err.message, channel: channelName, txid: txid}; + const parsedMsg = {error:err.message || 'ERROR', channel: channelName, txid: txid}; Server.send(userId, [0, HISTORY_KEEPER_ID, 'MSG', userId, JSON.stringify(parsedMsg)]); return; } diff --git a/www/common/outer/async-store.js b/www/common/outer/async-store.js index 704038b2c..f4568dbc9 100644 --- a/www/common/outer/async-store.js +++ b/www/common/outer/async-store.js @@ -1843,7 +1843,7 @@ define([ Store.leavePad(null, data, function () {}); }; var conf = { - Cache: Cache, // ICE pad cache + Cache: store.neverCache ? undefined : Cache, // ICE pad cache onCacheStart: function () { postMessage(clientId, "PAD_CACHE"); }, @@ -3271,6 +3271,7 @@ define([ // ==> don't create a drive // Or "neverDrive" (integration into another platform?) // ==> don't create a drive BUT create temp RPC (we may need to upload) + if (data.neverDrive) { store.neverCache = true; } if (data.neverDrive || (data.noDrive && !data.userHash && !data.anonHash)) { return void onNoDrive(clientId, function (obj) { if (obj && obj.error) { From ced24e655beb511fae94d8047fa3559cd202dc9a Mon Sep 17 00:00:00 2001 From: yflory Date: Tue, 3 Dec 2024 16:34:10 +0100 Subject: [PATCH 14/53] Fix API issues --- lib/hk-util.js | 2 +- www/common/outer/async-store.js | 2 +- www/cryptpad-api.js | 3 ++- 3 files changed, 4 insertions(+), 3 deletions(-) diff --git a/lib/hk-util.js b/lib/hk-util.js index 2ba656376..61374d7a4 100644 --- a/lib/hk-util.js +++ b/lib/hk-util.js @@ -135,7 +135,7 @@ const expireChannel = HK.expireChannel = function (Env, channel) { }; const removeChannel = HK.removeChannel = function (Env, channel) { - if (!Env.store) { return; } + if (!Env.store) { return; } Env.store.archiveChannel(channel, void 0, () => {}); delete Env.metadata_cache[channel]; delete Env.channel_cache[channel]; diff --git a/www/common/outer/async-store.js b/www/common/outer/async-store.js index f4568dbc9..6e37a97c2 100644 --- a/www/common/outer/async-store.js +++ b/www/common/outer/async-store.js @@ -1843,7 +1843,7 @@ define([ Store.leavePad(null, data, function () {}); }; var conf = { - Cache: store.neverCache ? undefined : Cache, // ICE pad cache + Cache: store.neverCache ? undefined : Cache, onCacheStart: function () { postMessage(clientId, "PAD_CACHE"); }, diff --git a/www/cryptpad-api.js b/www/cryptpad-api.js index 6fd6d3be4..c398cdb54 100644 --- a/www/cryptpad-api.js +++ b/www/cryptpad-api.js @@ -113,7 +113,7 @@ }; var start = function () { - config.document.key = key; + //config.document.key = key; chan.send('START', { key: key, application: config.documentType, @@ -137,6 +137,7 @@ blob = config.document.blob; return start(); } + // XXX Nextcloud will log us out if we try from the client getBlob(function (err, _blob) { if (err) { // Can't get blob from client, try from server console.warn(err); From 90c95beefee96de3ea61e112dfd5833e1ecc7c2b Mon Sep 17 00:00:00 2001 From: yflory Date: Tue, 17 Dec 2024 16:05:28 +0100 Subject: [PATCH 15/53] 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 16/53] 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 17/53] 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 18/53] 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 19/53] 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 20/53] 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 21/53] 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; From b135e815e7bd353ac7240ecffbdfb904c5cf7711 Mon Sep 17 00:00:00 2001 From: yflory Date: Tue, 14 Jan 2025 18:13:45 +0100 Subject: [PATCH 22/53] Allow new types of decrees in different files --- lib/commands/admin-rpc.js | 13 ++ lib/decrees-core.js | 141 +++++++++++++++++ lib/decrees.js | 324 ++++++++++++-------------------------- lib/env.js | 2 +- lib/http-worker.js | 15 +- lib/load-config.js | 7 - package-lock.json | 16 +- package.json | 2 +- www/admin/inner.js | 22 ++- 9 files changed, 287 insertions(+), 255 deletions(-) create mode 100644 lib/decrees-core.js diff --git a/lib/commands/admin-rpc.js b/lib/commands/admin-rpc.js index 9030159f8..b682bbf62 100644 --- a/lib/commands/admin-rpc.js +++ b/lib/commands/admin-rpc.js @@ -1138,6 +1138,19 @@ Admin.command = function (Env, safeKey, data, _cb, Server) { var command = commands[data[0]]; + Object.keys(Env.plugins || {}).forEach(name => { + let plugin = Env.plugins[name]; + if (!plugin.addAdminCommands) { return; } + try { + let c = plugin.addAdminCommands(Env); + Object.keys(c || {}).forEach(cmd => { + if (typeof(c[cmd]) !== "function") { return; } + if (commands[cmd]) { return; } + commands[cmd] = c[cmd]; + }); + } catch (e) {} + }); + if (typeof(command) === 'function') { return void command(Env, Server, cb, data, unsafeKey); } diff --git a/lib/decrees-core.js b/lib/decrees-core.js new file mode 100644 index 000000000..a8320a922 --- /dev/null +++ b/lib/decrees-core.js @@ -0,0 +1,141 @@ +// SPDX-FileCopyrightText: 2023 XWiki CryptPad Team and contributors +// +// SPDX-License-Identifier: AGPL-3.0-or-later + +var Decrees = module.exports; +var Util = require("./common-util"); +var Fs = require("fs"); +var Path = require("path"); +var readFileBin = require("./stream-file").readFileBin; +var Schedule = require("./schedule"); +var Fse = require("fs-extra"); +var nThen = require("nthen"); + + +const Utils = Decrees.Utils = {}; +var isString = (str) => { + return typeof(str) === "string"; +}; +var isInteger = function (n) { + return !(typeof(n) !== 'number' || isNaN(n) || (n % 1) !== 0); +}; +Utils.args_isBoolean = function (args) { + return !(!Array.isArray(args) || typeof(args[0]) !== 'boolean'); +}; +Utils.args_isString = function (args) { + return !(!Array.isArray(args) || !isString(args[0])); +}; +Utils.args_isInteger = function (args) { + return !(!Array.isArray(args) || !isInteger(args[0])); +}; +Utils.args_isPositiveInteger = function (args) { + return Array.isArray(args) && isInteger(args[0]) && args[0] > 0; +}; + + +Decrees.create = (name, commands) => { + // [, , ,