diff --git a/drizzle/client/0005_mature_mephistopheles.sql b/drizzle/client/0005_mature_mephistopheles.sql new file mode 100644 index 000000000..9e004053f --- /dev/null +++ b/drizzle/client/0005_mature_mephistopheles.sql @@ -0,0 +1,5 @@ +CREATE TABLE `cores` ( + `projectPublicId` text NOT NULL, + `publicKey` blob NOT NULL, + `namespace` text NOT NULL +); diff --git a/drizzle/client/meta/0005_snapshot.json b/drizzle/client/meta/0005_snapshot.json new file mode 100644 index 000000000..4f9039abf --- /dev/null +++ b/drizzle/client/meta/0005_snapshot.json @@ -0,0 +1,265 @@ +{ + "version": "5", + "dialect": "sqlite", + "id": "aabe981d-9cf2-48d5-b207-68cb9043dd1d", + "prevId": "2dec4710-70ac-4857-b889-2bc3586695ee", + "tables": { + "cores": { + "name": "cores", + "columns": { + "projectPublicId": { + "name": "projectPublicId", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "publicKey": { + "name": "publicKey", + "type": "blob", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "namespace": { + "name": "namespace", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": {}, + "foreignKeys": {}, + "compositePrimaryKeys": {}, + "uniqueConstraints": {} + }, + "deviceSettings": { + "name": "deviceSettings", + "columns": { + "deviceId": { + "name": "deviceId", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "deviceInfo": { + "name": "deviceInfo", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "isArchiveDevice": { + "name": "isArchiveDevice", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + } + }, + "indexes": { + "deviceSettings_deviceId_unique": { + "name": "deviceSettings_deviceId_unique", + "columns": [ + "deviceId" + ], + "isUnique": true + } + }, + "foreignKeys": {}, + "compositePrimaryKeys": {}, + "uniqueConstraints": {} + }, + "projectSettings_backlink": { + "name": "projectSettings_backlink", + "columns": { + "versionId": { + "name": "versionId", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + } + }, + "indexes": {}, + "foreignKeys": {}, + "compositePrimaryKeys": {}, + "uniqueConstraints": {} + }, + "projectKeys": { + "name": "projectKeys", + "columns": { + "projectId": { + "name": "projectId", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "projectPublicId": { + "name": "projectPublicId", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "projectInviteId": { + "name": "projectInviteId", + "type": "blob", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "keysCipher": { + "name": "keysCipher", + "type": "blob", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "projectInfo": { + "name": "projectInfo", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false, + "default": "'{\"sendStats\":false}'" + }, + "hasLeftProject": { + "name": "hasLeftProject", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false, + "default": false + } + }, + "indexes": {}, + "foreignKeys": {}, + "compositePrimaryKeys": {}, + "uniqueConstraints": {} + }, + "projectSettings": { + "name": "projectSettings", + "columns": { + "docId": { + "name": "docId", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "versionId": { + "name": "versionId", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "originalVersionId": { + "name": "originalVersionId", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "schemaName": { + "name": "schemaName", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "createdAt": { + "name": "createdAt", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "updatedAt": { + "name": "updatedAt", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "links": { + "name": "links", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "deleted": { + "name": "deleted", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "name": { + "name": "name", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "projectDescription": { + "name": "projectDescription", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "projectColor": { + "name": "projectColor", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "sendStats": { + "name": "sendStats", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "defaultPresets": { + "name": "defaultPresets", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "configMetadata": { + "name": "configMetadata", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "forks": { + "name": "forks", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": {}, + "foreignKeys": {}, + "compositePrimaryKeys": {}, + "uniqueConstraints": {} + } + }, + "enums": {}, + "_meta": { + "schemas": {}, + "tables": {}, + "columns": {} + } +} \ No newline at end of file diff --git a/drizzle/client/meta/_journal.json b/drizzle/client/meta/_journal.json index d6dd71db7..efc9f6959 100644 --- a/drizzle/client/meta/_journal.json +++ b/drizzle/client/meta/_journal.json @@ -36,6 +36,13 @@ "when": 1758874242014, "tag": "0004_glorious_shape", "breakpoints": true + }, + { + "idx": 5, + "version": "5", + "when": 1759256715782, + "tag": "0005_mature_mephistopheles", + "breakpoints": true } ] } \ No newline at end of file diff --git a/src/core-manager/index.js b/src/core-manager/index.js index 534d2ca31..ecfbc633a 100644 --- a/src/core-manager/index.js +++ b/src/core-manager/index.js @@ -2,7 +2,7 @@ import { TypedEmitter } from 'tiny-typed-emitter' import Corestore from 'corestore' import { debounce } from 'throttle-debounce' import assert from 'node:assert/strict' -import { sql, eq } from 'drizzle-orm' +import { sql, eq, and } from 'drizzle-orm' import { HaveExtension, @@ -11,8 +11,8 @@ import { } from '../generated/extensions.js' import { Logger } from '../logger.js' import { NAMESPACES } from '../constants.js' -import { noop } from '../utils.js' -import { coresTable } from '../schema/project.js' +import { noop, projectKeyToPublicId } from '../utils.js' +import { coresTable } from '../schema/client.js' import * as rle from './bitfield-rle.js' import { CoreIndex } from './core-index.js' import mapObject from 'map-obj' @@ -59,7 +59,7 @@ export class CoreManager extends TypedEmitter { /** * @param {Object} options - * @param {import('drizzle-orm/better-sqlite3').BetterSQLite3Database} options.db Drizzle better-sqlite3 database instance + * @param {import('drizzle-orm/better-sqlite3').BetterSQLite3Database} options.db Drizzle better-sqlite3 database instance * @param {import('@mapeo/crypto').KeyManager} options.keyManager mapeo/crypto KeyManager instance * @param {Buffer} options.projectKey 32-byte public key of the project creator core * @param {Buffer} [options.projectSecretKey] 32-byte secret key of the project creator core @@ -95,11 +95,14 @@ export class CoreManager extends TypedEmitter { this.#encryptionKeys = encryptionKeys this.#autoDownload = autoDownload + const projectPublicId = projectKeyToPublicId(projectKey) + // Pre-prepare SQL statement for better performance this.#queries = { addCore: db .insert(coresTable) .values({ + projectPublicId, publicKey: sql.placeholder('publicKey'), namespace: sql.placeholder('namespace'), }) @@ -107,7 +110,12 @@ export class CoreManager extends TypedEmitter { .prepare(), removeCores: db .delete(coresTable) - .where(eq(coresTable.namespace, sql.placeholder('namespace'))) + .where( + and( + eq(coresTable.namespace, sql.placeholder('namespace')), + eq(coresTable.projectPublicId, projectPublicId) + ) + ) .prepare(), } diff --git a/src/core-ownership.js b/src/core-ownership.js index 6eee2960a..16d6ad834 100644 --- a/src/core-ownership.js +++ b/src/core-ownership.js @@ -1,6 +1,4 @@ import { verifySignature, sign } from '@mapeo/crypto' -import { parseVersionId } from '@comapeo/schema' -import { defaultGetWinner } from '@mapeo/sqlite-indexer' import assert from 'node:assert/strict' import sodium from 'sodium-universal' import { @@ -11,11 +9,9 @@ import { } from './datatype/index.js' import { eq, or } from 'drizzle-orm' import mapObject from 'map-obj' -import { discoveryKey } from 'hypercore-crypto' import pDefer from 'p-defer' import { NAMESPACES } from './constants.js' import { TypedEmitter } from 'tiny-typed-emitter' -import { omit } from './lib/omit.js' import { NotFoundError } from './errors.js' /** * @import { @@ -145,35 +141,6 @@ export class CoreOwnership extends TypedEmitter { } } -/** - * - Validate that the doc is written to the core identified by doc.authCoreId - * - Verify the signatures - * - Remove the signatures (we don't add them to the indexer) - * - Set doc.links to an empty array - this forces the indexer to treat every - * document as a fork, so getWinner is called for every doc, which resolves to - * the doc with the lowest index (e.g. the first) - * - * @param {CoreOwnershipWithSignatures} doc - * @param {import('@comapeo/schema').VersionIdObject} version - * @returns {import('@comapeo/schema').CoreOwnership} - */ -export function mapAndValidateCoreOwnership(doc, { coreDiscoveryKey }) { - if ( - !coreDiscoveryKey.equals(discoveryKey(Buffer.from(doc.authCoreId, 'hex'))) - ) { - throw new Error('Invalid coreOwnership record: mismatched authCoreId') - } - if (!verifyCoreOwnership(doc)) { - throw new Error('Invalid coreOwnership record: signatures are invalid') - } - const docWithoutSignatures = omit(doc, [ - 'identitySignature', - 'coreSignatures', - ]) - docWithoutSignatures.links = [] - return docWithoutSignatures -} - /** * Verify the signatures of a coreOwnership record, which verify that the device * with the identityKey matching the docIds does own (e.g. can write to) cores @@ -182,7 +149,7 @@ export function mapAndValidateCoreOwnership(doc, { coreDiscoveryKey }) { * @param {CoreOwnershipWithSignatures} doc * @returns {boolean} */ -function verifyCoreOwnership(doc) { +export function verifyCoreOwnership(doc) { const { coreSignatures, identitySignature } = doc for (const namespace of NAMESPACES) { const signature = coreSignatures[namespace] @@ -211,26 +178,3 @@ function verifyCoreOwnership(doc) { if (!isValidIdentitySignature) return false return true } - -/** - * For coreOwnership records, we only trust the first record written to the core. - * - * @type {NonNullable[0]['getWinner']>} - */ -export function getWinner(docA, docB) { - if ( - 'schemaName' in docA && - docA.schemaName === 'coreOwnership' && - 'schemaName' in docB && - docB.schemaName === 'coreOwnership' - ) { - // Assumes docA and docB have same coreKey, so we choose the first one - // written to the core - const docAindex = parseVersionId(docA.versionId).index - const docBindex = parseVersionId(docB.versionId).index - if (docAindex < docBindex) return docA - return docB - } else { - return defaultGetWinner(docA, docB) - } -} diff --git a/src/datastore/index.js b/src/datastore/index.js index dab7458f2..4077a04ac 100644 --- a/src/datastore/index.js +++ b/src/datastore/index.js @@ -94,6 +94,19 @@ export class DataStore extends TypedEmitter { return this.#writerCore } + /** + * Wait for any data to be fully flushed to the indexer + * @returns {Promise} + */ + async waitIdle() { + while (this.#pendingAppends.size || this.#pendingIndex.size) { + const pendingIndexes = [...this.#pendingIndex.values()].map( + ({ promise }) => promise + ) + await Promise.all([...this.#pendingAppends].concat(pendingIndexes)) + } + } + /** * * @param {MultiCoreIndexer.Entry<'binary'>[]} entries diff --git a/src/datatype/index.js b/src/datatype/index.js index 4bb0c5fd2..2555101d4 100644 --- a/src/datatype/index.js +++ b/src/datatype/index.js @@ -104,7 +104,7 @@ export class DataType extends TypedEmitter { * @param {object} opts * @param {TTable} opts.table * @param {TDataStore} opts.dataStore - * @param {import('drizzle-orm/better-sqlite3').BetterSQLite3Database} opts.db + * @param {import('drizzle-orm/better-sqlite3').BetterSQLite3Database>} opts.db * @param {import('../translation-api.js').default['get']} [opts.getTranslations] * @param {(versionId: string) => Promise} opts.getDeviceIdForVersionId */ diff --git a/src/discovery/local-discovery.js b/src/discovery/local-discovery.js index 24bb0df5a..125f80ff0 100644 --- a/src/discovery/local-discovery.js +++ b/src/discovery/local-discovery.js @@ -290,9 +290,8 @@ export class LocalDiscovery extends TypedEmitter { async #stop({ force = false, timeout = 0 } = {}) { this.#log('stopping') const port = this.#port - this.#server.close() const closePromise = once(this.#server, 'close') - + this.#server.close() const forceClose = () => { for (const socket of this.#noiseConnections.values()) { socket.destroy() @@ -311,6 +310,7 @@ export class LocalDiscovery extends TypedEmitter { fallback: forceClose, }) } + this.#log(`stopped for ${port}`) } } diff --git a/src/index-writer/get-winner.js b/src/index-writer/get-winner.js new file mode 100644 index 000000000..5b19db9e0 --- /dev/null +++ b/src/index-writer/get-winner.js @@ -0,0 +1,24 @@ +import { parseVersionId } from '@comapeo/schema' +import { defaultGetWinner } from '@mapeo/sqlite-indexer' + +/** + * For coreOwnership records, we only trust the first record written to the core. + * + * @type {typeof import('@mapeo/sqlite-indexer').defaultGetWinner} + */ +export function getWinner(docA, docB) { + // Written "backwards" to minimize conditional checks if not 'coreOwnership' + if ( + !('schemaName' in docA && docA.schemaName === 'coreOwnership') || + !('schemaName' in docB && docB.schemaName === 'coreOwnership') + ) { + return defaultGetWinner(docA, docB) + } else { + // Assumes docA and docB have same coreKey, so we choose the first one + // written to the core + const docAindex = parseVersionId(docA.versionId).index + const docBindex = parseVersionId(docB.versionId).index + if (docAindex < docBindex) return docA + return docB + } +} diff --git a/src/index-writer/index-worker.js b/src/index-writer/index-worker.js new file mode 100644 index 000000000..f9bd693b4 --- /dev/null +++ b/src/index-writer/index-worker.js @@ -0,0 +1,91 @@ +import ensureError from 'ensure-error' +import { workerData, parentPort } from 'worker_threads' +import { IndexWriter } from './index-writer.js' +import { Logger } from '../logger.js' +import Database from 'better-sqlite3' + +/** @import { + WorkerRequest, + BatchRequestData, DeleteSchemaRequestData, CloseRequestData, + BatchResponseData, DeleteSchemaResponseData, CloseResponseData, + WorkerData, WorkerResponse + } from './index-writer-proxy.js' */ + +/** @typedef { + WorkerRequest<'batch', BatchRequestData> | + WorkerRequest<'close', CloseRequestData> | + WorkerRequest<'deleteSchema', DeleteSchemaRequestData> + } ExpectedRequestMessage +*/ + +const { schemas, dbPath, parentLoggerNamespace, deviceId } = + /** @type {WorkerData} */ (workerData) + +const sqlite = new Database(dbPath) +sqlite.pragma('journal_mode=WAL') + +const indexWriter = new IndexWriter({ + schemas, + sqlite, + logger: new Logger({ + ns: parentLoggerNamespace, + deviceId: deviceId || '', + }), +}) + +if (!parentPort) { + throw new Error('This module must be run in a worker thread.') +} + +parentPort.on('message', handleMessage) +parentPort.start() +// Message queuing on the parentPort isn't happening +// This lets the worker proxy know we can get messages +const readyData = { id: -1, data: null } +parentPort.postMessage(readyData) + +/** + * @param {ExpectedRequestMessage} msg + * @returns {Promise} + */ +async function handleMessage({ id, type, data }) { + if (!parentPort) { + throw new Error('This module must be run in a worker thread.') + } + try { + switch (type) { + case 'batch': { + // Need to convert b + for (const item of data) { + item.block = Buffer.from(item.block) + } + const result = await indexWriter.batch(data) + /** @type {WorkerResponse} */ + const msg = { id, data: result } + parentPort.postMessage(msg) + return + } + case 'deleteSchema': { + await indexWriter.deleteSchema(data) + /** @type {WorkerResponse} */ + const msg = { id, data: null } + parentPort.postMessage(msg) + return + } + case 'close': { + sqlite.close() + parentPort.postMessage({ id, data: 'ok' }) + process.exit(0) + return + } + default: + throw new Error(`Unknown message type: ${type}`) + } + } catch (error) { + parentPort.postMessage({ + id, + error: ensureError(error).message, + errorStack: ensureError(error).stack, + }) + } +} diff --git a/src/index-writer/index-writer-proxy.js b/src/index-writer/index-writer-proxy.js new file mode 100644 index 000000000..34b56c70e --- /dev/null +++ b/src/index-writer/index-writer-proxy.js @@ -0,0 +1,134 @@ +import { Worker } from 'worker_threads' +import { pEvent } from 'p-event' +/** @import { MapeoDoc } from '@comapeo/schema' */ +/** @import { IndexedDocIds, SchemaName } from './index-writer.js'*/ + +/** + * @template {string} [TType=string] + * @template [TData=any] + * @typedef {{ id: number, type: TType, data: TData }} WorkerRequest + */ +/** + * @template T + * @typedef {{ id: number, data: T } | { id: number, error: string, errorStack: string }} WorkerResponse + */ +/** @typedef {SchemaName} DeleteSchemaRequestData */ +/** @typedef {null} DeleteSchemaResponseData */ +/** @typedef {null} CloseRequestData */ +/** @typedef {string} CloseResponseData */ +/** @typedef {import('multi-core-indexer').Entry[]} BatchRequestData */ +/** @typedef {IndexedDocIds} BatchResponseData */ +/** + * @typedef {object} WorkerData + * @property {SchemaName[]} schemas + * @property {string} dbPath + * @property {string} [parentLoggerNamespace] + * @property {string} [deviceId] + */ + +/** + * Proxy calls to the IndexWriter class to a worker thread. + */ +export class IndexWriterProxy { + #worker + #nextId = 0 + #onLoaded + + /** + * + * @param {object} opts + * @param {import('better-sqlite3').Database} opts.sqlite + * @param {Iterable} opts.schemas + * @param {import('../logger.js').Logger} [opts.logger] + */ + constructor({ schemas, sqlite, logger }) { + if (sqlite.memory) { + throw new Error( + 'Cannot use IndexWriterProxy with an in-memory SQLite database' + ) + } + /** @type {WorkerData} */ + const workerData = { + schemas: Array.from(schemas), + dbPath: sqlite.name, + parentLoggerNamespace: logger?.ns, + deviceId: logger?.deviceId, + } + this.#worker = new Worker(new URL('./index-worker.js', import.meta.url), { + workerData, + }) + this.#onLoaded = pEvent(this.#worker, 'message', { + // Signifies "ready" + filter: (msg) => msg.id === -1, + }) + this.#worker.unref() + } + + /** + * @overload + * @param {'batch'} type + * @param {BatchRequestData} data + * @param {import('worker_threads').TransferListItem[]} [transferList] + * @returns {Promise} + */ + /** + * @overload + * @param {'deleteSchema'} type + * @param {DeleteSchemaRequestData} data + * @param {import('worker_threads').TransferListItem[]} [transferList] + * @returns {Promise} + */ + /** + * @overload + * @param {'close'} type + * @param {CloseRequestData} data + * @param {import('worker_threads').TransferListItem[]} [transferList] + * @returns {Promise} + */ + /** + * @param {string} type + * @param {any} data + * @param {import('worker_threads').TransferListItem[]} [transferList] + * @returns {Promise} resolves with the response from the worker + */ + async #workerRequest(type, data, transferList) { + const id = this.#nextId++ + const responsePromise = pEvent(this.#worker, 'message', { + filter: (msg) => msg.id === id, + }) + /** @type {WorkerRequest} */ + const request = { id, type, data } + this.#worker.postMessage(request, transferList) + const response = /** @type {WorkerResponse} */ (await responsePromise) + if ('error' in response) { + throw new Error(response.error + '\n\n' + response.errorStack) + } + return response.data + } + + /** + * @param {import('multi-core-indexer').Entry[]} entries + * @returns {Promise} map of indexed docIds by schemaName + */ + async batch(entries) { + const transferList = entries.map((entry) => entry.block.buffer) + return this.#workerRequest('batch', entries, transferList) + } + + /** + * @param {SchemaName} schemaName + * @return {Promise} + */ + async deleteSchema(schemaName) { + await this.#workerRequest('deleteSchema', schemaName) + } + + /** + * Clean up any remaining index writer resources + * @returns {Promise} + */ + async close() { + await this.#onLoaded + await this.#workerRequest('close', null) + } +} diff --git a/src/index-writer/index-writer.js b/src/index-writer/index-writer.js new file mode 100644 index 000000000..b89d9476a --- /dev/null +++ b/src/index-writer/index-writer.js @@ -0,0 +1,113 @@ +import { decode } from '@comapeo/schema' +import SqliteIndexer from '@mapeo/sqlite-indexer' +import { getBacklinkTableName } from '../schema/comapeo-to-drizzle.js' +import { discoveryKey } from 'hypercore-crypto' +import { Logger } from '../logger.js' +import { mapDoc } from './map-doc.js' +import { getWinner } from './get-winner.js' +/** @import { MapeoDoc } from '@comapeo/schema' */ +/** + * @typedef {{ [K in MapeoDoc['schemaName']]?: string[] }} IndexedDocIds + */ +/** + * @typedef {MapeoDoc['schemaName']} SchemaName + */ + +export class IndexWriter { + /** @type {Map} */ + #indexers = new Map() + #mapDoc + #l + /** + * + * @param {object} opts + * @param {import('better-sqlite3').Database} opts.sqlite + * @param {Iterable} opts.schemas + * @param {Logger} [opts.logger] + */ + constructor({ schemas, sqlite, logger }) { + this.#l = Logger.create('indexWriter', logger) + this.#mapDoc = mapDoc + + for (const schemaName of schemas) { + const indexer = new SqliteIndexer(sqlite, { + docTableName: schemaName, + backlinkTableName: getBacklinkTableName(schemaName), + getWinner, + }) + this.#indexers.set(schemaName, indexer) + } + } + + /** + * @param {import('multi-core-indexer').Entry[]} entries + * @returns {Promise} map of indexed docIds by schemaName + */ + async batch(entries) { + // sqlite-indexer is _significantly_ faster when batching even <10 at a + // time, so best to queue docs here before calling sliteIndexer.batch() + /** @type {Record} */ + const queued = {} + /** @type {IndexedDocIds} */ + const indexed = {} + for (const { block, key, index } of entries) { + /** @type {MapeoDoc} */ let doc + try { + const version = { coreDiscoveryKey: discoveryKey(key), index } + doc = this.#mapDoc(decode(block, version), version) + } catch (e) { + this.#l.log('Could not decode entry %d of %h', index, key) + // Unknown or invalid entry - silently ignore + continue + } + // Don't have an indexer for this type - silently ignore + if (!this.#indexers.has(doc.schemaName)) continue + if (queued[doc.schemaName]) { + queued[doc.schemaName].push(doc) + // @ts-expect-error - we know this is defined, TS doesn't + indexed[doc.schemaName].push(doc.docId) + } else { + queued[doc.schemaName] = [doc] + indexed[doc.schemaName] = [doc.docId] + } + } + for (const [schemaName, docs] of Object.entries(queued)) { + // @ts-expect-error + const indexer = this.#indexers.get(schemaName) + if (!indexer) continue // Won't happen, but TS doesn't know that + indexer.batch(docs) + // TODO: selectively turn this on when log level is 'trace' or 'debug' + // Otherwise this has a big performance overhead because this is all synchronous + // if (this.#l.log.enabled) { + // for (const doc of docs) { + // this.#l.log( + // 'Indexed %s %S @ %S', + // doc.schemaName, + // doc.docId, + // doc.versionId + // ) + // } + // } + } + return indexed + } + + /** + * @param {SchemaName} schemaName + */ + async deleteSchema(schemaName) { + const indexer = this.#indexers.get(schemaName) + if (!indexer) { + throw new Error(`IndexWriter doesn't know a schema named "${schemaName}"`) + } + await indexer.deleteAll() + } + + /** + * Clean up any remaining index writer resources + * @returns {Promise} + */ + async close() { + // Nothing to do, everything is synchronous + } +} diff --git a/src/index-writer/index.js b/src/index-writer/index.js index 8dd145672..67507fca8 100644 --- a/src/index-writer/index.js +++ b/src/index-writer/index.js @@ -1,63 +1,64 @@ -import { decode } from '@comapeo/schema' -import SqliteIndexer from '@mapeo/sqlite-indexer' import { getTableConfig } from 'drizzle-orm/sqlite-core' -import { getBacklinkTableName } from '../schema/comapeo-to-drizzle.js' -import { discoveryKey } from 'hypercore-crypto' -import { Logger } from '../logger.js' -/** @import { MapeoDoc, VersionIdObject } from '@comapeo/schema' */ +import { IndexWriter } from './index-writer.js' +import { IndexWriterProxy } from './index-writer-proxy.js' + +export { IndexWriter, IndexWriterProxy } + +/** @import { MapeoDoc } from '@comapeo/schema' */ /** @import { MapeoDocTables } from '../datatype/index.js' */ /** * @typedef {{ [K in MapeoDoc['schemaName']]?: string[] }} IndexedDocIds */ -/** - * @typedef {ReturnType} MapeoDocInternal - */ /** * @template {MapeoDocTables} [TTables=MapeoDocTables] */ -export class IndexWriter { +export class IndexWriterWrapper { /** * @internal * @typedef {TTables['_']['name']} SchemaName */ - - /** @type {Map} */ - #indexers = new Map() - #mapDoc - #l + /** @type {SchemaName[]} */ + #schemas = [] + #indexWriter /** * * @param {object} opts * @param {import('better-sqlite3').Database} opts.sqlite * @param {TTables[]} opts.tables - * @param {(doc: MapeoDocInternal, version: VersionIdObject) => MapeoDoc} [opts.mapDoc] optionally transform a document prior to indexing. Can also validate, if an error is thrown then the document will not be indexed - * @param {typeof import('@mapeo/sqlite-indexer').defaultGetWinner} [opts.getWinner] custom function to determine the "winner" of two forked documents. Defaults to choosing the document with the most recent `updatedAt` - * @param {Logger} [opts.logger] + * @param {boolean} [opts.useWorker] if true, create a worker thread for indexing + * @param {import('../logger.js').Logger} [opts.logger] */ - constructor({ tables, sqlite, mapDoc = (d) => d, getWinner, logger }) { - this.#l = Logger.create('indexWriter', logger) - this.#mapDoc = mapDoc + constructor({ tables, sqlite, logger, useWorker = false }) { for (const table of tables) { const config = getTableConfig(table) const schemaName = /** @type {(typeof table)['_']['name']} */ ( config.name ) - const indexer = new SqliteIndexer(sqlite, { - docTableName: config.name, - backlinkTableName: getBacklinkTableName(config.name), - getWinner, + this.#schemas.push(schemaName) + } + + if (useWorker && !sqlite.memory) { + this.#indexWriter = new IndexWriterProxy({ + schemas: this.schemas, + sqlite, + logger, + }) + } else { + this.#indexWriter = new IndexWriter({ + schemas: this.schemas, + sqlite, + logger, }) - this.#indexers.set(schemaName, indexer) } } /** - * @returns {Iterable} + * @returns {Array} */ get schemas() { - return this.#indexers.keys() + return this.#schemas } /** @@ -65,62 +66,22 @@ export class IndexWriter { * @returns {Promise} map of indexed docIds by schemaName */ async batch(entries) { - // sqlite-indexer is _significantly_ faster when batching even <10 at a - // time, so best to queue docs here before calling sliteIndexer.batch() - /** @type {Record} */ - const queued = {} - /** @type {IndexedDocIds} */ - const indexed = {} - for (const { block, key, index } of entries) { - /** @type {MapeoDoc} */ let doc - try { - const version = { coreDiscoveryKey: discoveryKey(key), index } - doc = this.#mapDoc(decode(block, version), version) - } catch (e) { - this.#l.log('Could not decode entry %d of %h', index, key) - // Unknown or invalid entry - silently ignore - continue - } - // Don't have an indexer for this type - silently ignore - if (!this.#indexers.has(doc.schemaName)) continue - if (queued[doc.schemaName]) { - queued[doc.schemaName].push(doc) - // @ts-expect-error - we know this is defined, TS doesn't - indexed[doc.schemaName].push(doc.docId) - } else { - queued[doc.schemaName] = [doc] - indexed[doc.schemaName] = [doc.docId] - } - } - for (const [schemaName, docs] of Object.entries(queued)) { - // @ts-expect-error - const indexer = this.#indexers.get(schemaName) - if (!indexer) continue // Won't happen, but TS doesn't know that - indexer.batch(docs) - // TODO: selectively turn this on when log level is 'trace' or 'debug' - // Otherwise this has a big performance overhead because this is all synchronous - // if (this.#l.log.enabled) { - // for (const doc of docs) { - // this.#l.log( - // 'Indexed %s %S @ %S', - // doc.schemaName, - // doc.docId, - // doc.versionId - // ) - // } - // } - } - return indexed + return this.#indexWriter.batch(entries) } /** * @param {SchemaName} schemaName + * @return {Promise} */ - deleteSchema(schemaName) { - const indexer = this.#indexers.get(schemaName) - if (!indexer) { - throw new Error(`IndexWriter doesn't know a schema named "${schemaName}"`) - } - indexer.deleteAll() + async deleteSchema(schemaName) { + return this.#indexWriter.deleteSchema(schemaName) + } + + /** + * Clean up any remaining index writer resources + * @returns {Promise} + */ + async close() { + await this.#indexWriter.close() } } diff --git a/src/index-writer/map-doc.js b/src/index-writer/map-doc.js new file mode 100644 index 000000000..8efa166cb --- /dev/null +++ b/src/index-writer/map-doc.js @@ -0,0 +1,50 @@ +import { discoveryKey } from 'hypercore-crypto' +import { omit } from '../lib/omit.js' +import { verifyCoreOwnership } from '../core-ownership.js' + +/** + * @typedef {ReturnType} MapeoDocInternal + */ + +/** + * @param {MapeoDocInternal} doc + * @param {import('@comapeo/schema').VersionIdObject} version + * @returns {import('@comapeo/schema').MapeoDoc} + */ +export function mapDoc(doc, version) { + switch (doc.schemaName) { + case 'coreOwnership': { + if ( + !version.coreDiscoveryKey.equals( + discoveryKey(Buffer.from(doc.authCoreId, 'hex')) + ) + ) { + throw new Error('Invalid coreOwnership record: mismatched authCoreId') + } + if (!verifyCoreOwnership(doc)) { + throw new Error('Invalid coreOwnership record: signatures are invalid') + } + const docWithoutSignatures = omit(doc, [ + 'identitySignature', + 'coreSignatures', + ]) + docWithoutSignatures.links = [] + return docWithoutSignatures + } + case 'deviceInfo': + // Validate that a deviceInfo record is written by the device that is it + // about, e.g. version.coreKey should equal docId + if ( + !version.coreDiscoveryKey.equals( + discoveryKey(Buffer.from(doc.docId, 'hex')) + ) + ) { + throw new Error( + 'Invalid deviceInfo record, cannot write deviceInfo for another device' + ) + } + return doc + default: + return doc + } +} diff --git a/src/lib/drizzle-helpers.js b/src/lib/drizzle-helpers.js index f4aaad827..61b3fb1ff 100644 --- a/src/lib/drizzle-helpers.js +++ b/src/lib/drizzle-helpers.js @@ -1,9 +1,11 @@ -import { sql } from 'drizzle-orm' +import { count, eq, sql } from 'drizzle-orm' import fs from 'node:fs' import path from 'node:path' import { assert } from '../utils.js' import { migrate as drizzleMigrate } from 'drizzle-orm/better-sqlite3/migrator' import { DRIZZLE_MIGRATIONS_TABLE } from '../constants.js' +import { coresTable } from '../schema/client.js' +import { coresTable_Deprecated } from '../schema/project.js' /** @import { BetterSQLite3Database } from 'drizzle-orm/better-sqlite3' */ /** @@ -25,7 +27,8 @@ const getNumberResult = (queryResult) => { * Get the latest migration time, or 0 if no migrations have been run or the * migrations table has not been created yet. * - * @param {BetterSQLite3Database} db + * @template {Record} TSchema + * @param {BetterSQLite3Database} db * @returns {number} */ const safeGetLatestMigrationMillis = (db) => @@ -61,10 +64,11 @@ const safeGetLatestMigrationMillis = (db) => * * Returns what happened during migration; did a migration occur? * - * @param {BetterSQLite3Database} db + * @template {Record} TSchema + * @param {BetterSQLite3Database} db * @param {object} options * @param {string} options.migrationsFolder - * @param {Record void>} [options.migrationFns] + * @param {Record) => void>} [options.migrationFns] * @returns {MigrationResult} */ export function migrate(db, { migrationsFolder, migrationFns = {} }) { @@ -102,6 +106,54 @@ export function migrate(db, { migrationsFolder, migrationFns = {} }) { return 'migrated' } +/** + * Copy the `cores` table from the project database to the client database. + * @param {object} options + * @param {BetterSQLite3Database} options.clientDb + * @param {BetterSQLite3Database} options.projectDb + * @param {string} options.projectPublicId + * @returns {void} + */ +export function migrateCoresTable({ clientDb, projectDb, projectPublicId }) { + const migratedRowCount = + clientDb + .select({ count: count() }) + .from(coresTable) + .where(eq(coresTable.projectPublicId, projectPublicId)) + .get()?.count || 0 + if (migratedRowCount > 0) { + // Already migrated, nothing to do + return + } + const projectCores = projectDb.select().from(coresTable_Deprecated).all() + if (projectCores.length === 0) { + // No cores to migrate + return + } + + clientDb.transaction((tx) => { + for (const core of projectCores) { + tx.insert(coresTable) + .values({ + ...core, + projectPublicId, + }) + .run() + } + }) + // Verify that the migration was successful + const migratedCount = + clientDb + .select({ count: count() }) + .from(coresTable) + .where(eq(coresTable.projectPublicId, projectPublicId)) + .get()?.count || 0 + assert( + migratedCount === projectCores.length, + `Expected to migrate ${projectCores.length} cores, but migrated ${migratedCount}` + ) +} + /** * Assert that the migration journal is the expected format. * @param {unknown} journal diff --git a/src/logger.js b/src/logger.js index 9043402e8..6fd3a0d34 100644 --- a/src/logger.js +++ b/src/logger.js @@ -118,6 +118,9 @@ export class Logger { this.#log = log } } + get ns() { + return this.#baseLogger.namespace + } get log() { return this.#log } diff --git a/src/mapeo-manager.js b/src/mapeo-manager.js index f1f0dfbc0..20f7dc2df 100644 --- a/src/mapeo-manager.js +++ b/src/mapeo-manager.js @@ -8,8 +8,8 @@ import Hypercore from 'hypercore' import { TypedEmitter } from 'tiny-typed-emitter' import pTimeout from 'p-timeout' import { createRequire } from 'module' - -import { IndexWriter } from './index-writer/index.js' +import * as clientSchema from './schema/client.js' +import { IndexWriterWrapper as IndexWriter } from './index-writer/index.js' import { MapeoProject, kBlobStore, @@ -128,6 +128,7 @@ export class MapeoManager extends TypedEmitter { #l #defaultConfigPath #makeWebsocket + #useIndexWorkers #defaultIsArchiveDevice /** @@ -144,6 +145,7 @@ export class MapeoManager extends TypedEmitter { * @param {string} [opts.defaultOnlineStyleUrl] URL for an online-hosted StyleJSON asset. * @param {boolean} [opts.defaultIsArchiveDevice] Whether the node is an archive device by default * @param {(url: string) => WebSocket} [opts.makeWebsocket] + * @param {boolean} [opts.useIndexWorkers] if true, use a worker thread for each project for indexing cores to sqlite */ constructor({ rootKey, @@ -158,6 +160,7 @@ export class MapeoManager extends TypedEmitter { defaultOnlineStyleUrl = DEFAULT_ONLINE_STYLE_URL, defaultIsArchiveDevice = DEFAULT_IS_ARCHIVE_DEVICE, makeWebsocket = (url) => new WebSocket(url), + useIndexWorkers = false, }) { super() this.#keyManager = new KeyManager(rootKey) @@ -169,13 +172,15 @@ export class MapeoManager extends TypedEmitter { this.#l = Logger.create('manager', logger) this.#dbFolder = dbFolder this.#projectMigrationsFolder = projectMigrationsFolder + this.#useIndexWorkers = useIndexWorkers + const sqlite = new Database( dbFolder === ':memory:' ? ':memory:' : path.join(dbFolder, CLIENT_SQLITE_FILE_NAME) ) sqlite.pragma('journal_mode=WAL') - this.#db = drizzle(sqlite) + this.#db = drizzle(sqlite, { schema: clientSchema }) migrate(this.#db, { migrationsFolder: clientMigrationsFolder, migrationFns: { @@ -574,6 +579,7 @@ export class MapeoManager extends TypedEmitter { getMediaBaseUrl: this.#getMediaBaseUrl.bind(this), isArchiveDevice, makeWebsocket: this.#makeWebsocket, + useIndexWorkers: this.#useIndexWorkers, getFallbackProjectInfo: () => { return this.#db .select({ projectInfo: projectKeysTable.projectInfo }) @@ -761,7 +767,9 @@ export class MapeoManager extends TypedEmitter { try { await this.#waitForInitialSync(project) } catch (e) { - this.#l.log('ERROR: could not do initial project sync', e) + // Needed for TS to allow e.stack 🙄 + if (!(e instanceof Error)) throw e + this.#l.log('ERROR: could not do initial project sync', e.stack) } } this.#l.log('Added project %h, public ID: %S', projectKey, projectPublicId) @@ -813,7 +821,7 @@ export class MapeoManager extends TypedEmitter { } else { this.#l.log( 'Pending initial sync: role %s, projectSettings %o, auth %o, config %o', - isRoleSynced, + ownRole.name, isProjectSettingsSynced, isAuthSynced, isConfigSynced @@ -1076,8 +1084,7 @@ export class MapeoManager extends TypedEmitter { * @returns {Promise} */ async close() { - // This added for workers PR - // await this.#projectSettingsIndexWriter.close() + await this.#projectSettingsIndexWriter.close() await Promise.all( [...this.#activeProjects.values()].map((project) => project.close()) ) diff --git a/src/mapeo-project.js b/src/mapeo-project.js index 65a3132ac..f2bf81ed8 100644 --- a/src/mapeo-project.js +++ b/src/mapeo-project.js @@ -3,7 +3,6 @@ import Database from 'better-sqlite3' import { decodeBlockPrefix, decode, parseVersionId } from '@comapeo/schema' import { drizzle } from 'drizzle-orm/better-sqlite3' import { sql, count, eq } from 'drizzle-orm' -import { discoveryKey } from 'hypercore-crypto' import { TypedEmitter } from 'tiny-typed-emitter' import ZipArchive from 'zip-stream-promise' import * as b4a from 'b4a' @@ -17,7 +16,7 @@ import { DataStore } from './datastore/index.js' import { DataType, kCreateWithDocId } from './datatype/index.js' import { BlobStore } from './blob-store/index.js' import { BlobApi } from './blob-api.js' -import { IndexWriter } from './index-writer/index.js' +import { IndexWriterWrapper as IndexWriter } from './index-writer/index.js' import { projectSettingsTable } from './schema/client.js' import { coreOwnershipTable, @@ -31,11 +30,7 @@ import { translationTable, remoteDetectionAlertTable, } from './schema/project.js' -import { - CoreOwnership, - getWinner, - mapAndValidateCoreOwnership, -} from './core-ownership.js' +import { CoreOwnership } from './core-ownership.js' import { BLOCKED_ROLE_ID, COORDINATOR_ROLE_ID, @@ -51,7 +46,7 @@ import { projectKeyToPublicId, valueOf, } from './utils.js' -import { migrate } from './lib/drizzle-helpers.js' +import { migrate, migrateCoresTable } from './lib/drizzle-helpers.js' import { omit } from './lib/omit.js' import { MemberApi } from './member-api.js' import { @@ -66,6 +61,7 @@ import TranslationApi from './translation-api.js' import { NotFoundError, nullIfNotFound } from './errors.js' import { WebSocket } from 'ws' import { createWriteStream } from 'fs' +import * as projectSchema from './schema/project.js' import ensureError from 'ensure-error' /** @import { ProjectSettingsValue, Observation, Track } from '@comapeo/schema' */ /** @import { Attachment, CoreStorage, BlobFilter, BlobId, BlobStoreEntriesStream, KeyPair, Namespace, ReplicationStream, GenericBlobFilter, MapeoValueMap, MapeoDocMap } from './types.js' */ @@ -103,6 +99,7 @@ export const kClearData = Symbol('clear project data') export const kSetIsArchiveDevice = Symbol('set isArchiveDevice') export const kIsArchiveDevice = Symbol('isArchiveDevice (temp - test only)') export const kGeoJSONFileName = Symbol('geoJSONFileName') +export const kWaitForDataStoresIdle = Symbol('waitForDataStoresIdle') const EMPTY_PROJECT_SETTINGS = Object.freeze({ sendStats: false }) @@ -156,7 +153,7 @@ export class MapeoProject extends TypedEmitter { * @param {Buffer} opts.projectKey 32-byte public key of the project creator core * @param {Buffer} [opts.projectSecretKey] 32-byte secret key of the project creator core * @param {import('./generated/keys.js').EncryptionKeys} opts.encryptionKeys Encryption keys for each namespace - * @param {import('drizzle-orm/better-sqlite3').BetterSQLite3Database} opts.sharedDb + * @param {import('drizzle-orm/better-sqlite3').BetterSQLite3Database} opts.sharedDb * @param {IndexWriter} opts.sharedIndexWriter * @param {CoreStorage} opts.coreStorage Folder to store all hypercore data * @param {(mediaType: 'blobs' | 'icons') => Promise} opts.getMediaBaseUrl @@ -165,6 +162,7 @@ export class MapeoProject extends TypedEmitter { * @param {boolean} opts.isArchiveDevice Whether this device is an archive device * @param {() => import('./schema/client.js').ProjectInfo | undefined} opts.getFallbackProjectInfo * @param {Logger} [opts.logger] + * @param {boolean} [opts.useIndexWorkers] if true, use a worker thread for each project for indexing cores to sqlite * */ constructor({ @@ -182,6 +180,7 @@ export class MapeoProject extends TypedEmitter { localPeers, logger, isArchiveDevice, + useIndexWorkers = false, getFallbackProjectInfo, }) { super() @@ -198,11 +197,22 @@ export class MapeoProject extends TypedEmitter { this.#sqlite = new Database(dbPath) this.#sqlite.pragma('journal_mode=WAL') - const db = drizzle(this.#sqlite) + + let db = drizzle(this.#sqlite, { schema: projectSchema }) this.#db = db + const migrationResult = migrate(db, { migrationsFolder: projectMigrationsFolder, }) + + // Above v4.1.4 we moved the cores table from the project db to the shared + // client db, so the project db can be opened read-only in the main thread. + migrateCoresTable({ + clientDb: sharedDb, + projectDb: db, + projectPublicId: this.#projectPublicId, + }) + let reindex switch (migrationResult) { case 'initialized database': @@ -240,6 +250,15 @@ export class MapeoProject extends TypedEmitter { .run() } + if (useIndexWorkers && !this.#sqlite.memory) { + // Re-open the db as read-only, because all writes will be done in the worker thread + this.#sqlite.close() + this.#sqlite = new Database(dbPath, { readonly: true }) + this.#sqlite.pragma('journal_mode=WAL') + db = drizzle(this.#sqlite, { schema: projectSchema }) + this.#db = db + } + ///////// 3. Setup random-access-storage functions /** @type {ConstructorParameters[0]['storage']} */ @@ -258,25 +277,15 @@ export class MapeoProject extends TypedEmitter { projectKey, keyManager, storage: coreManagerStorage, - db, + db: sharedDb, logger: this.#l, }) this.#indexWriter = new IndexWriter({ tables: indexedTables, sqlite: this.#sqlite, - getWinner, - mapDoc: (doc, version) => { - switch (doc.schemaName) { - case 'coreOwnership': - return mapAndValidateCoreOwnership(doc, version) - case 'deviceInfo': - return mapAndValidateDeviceInfo(doc, version) - default: - return doc - } - }, logger: this.#l, + useWorker: useIndexWorkers, }) this.#dataStores = { @@ -593,7 +602,7 @@ export class MapeoProject extends TypedEmitter { } await Promise.all(dataStorePromises) await this.#coreManager.close() - + await this.#indexWriter.close() this.#sqlite.close() this.emit('close') @@ -1379,10 +1388,20 @@ export class MapeoProject extends TypedEmitter { const authSchemas = new Set(NAMESPACE_SCHEMAS.auth) for (const schemaName of this.#indexWriter.schemas) { const isAuthSchema = authSchemas.has(schemaName) - if (!isAuthSchema) this.#indexWriter.deleteSchema(schemaName) + if (!isAuthSchema) await this.#indexWriter.deleteSchema(schemaName) } } + /** + * Wait for the datastore to flush data to the indexer fully + * @returns {Promise} + */ + async [kWaitForDataStoresIdle]() { + await Promise.all( + Object.values(this.#dataStores).map((store) => store.waitIdle()) + ) + } + /** * @deprecated * @param {object} opts @@ -1455,23 +1474,6 @@ function getCoreKeypairs({ projectKey, projectSecretKey, keyManager }) { return keypairs } -/** - * Validate that a deviceInfo record is written by the device that is it about, - * e.g. version.coreKey should equal docId - * - * @param {import('@comapeo/schema').DeviceInfo} doc - * @param {import('@comapeo/schema').VersionIdObject} version - * @returns {import('@comapeo/schema').DeviceInfo} - */ -function mapAndValidateDeviceInfo(doc, { coreDiscoveryKey }) { - if (!coreDiscoveryKey.equals(discoveryKey(Buffer.from(doc.docId, 'hex')))) { - throw new Error( - 'Invalid deviceInfo record, cannot write deviceInfo for another device' - ) - } - return doc -} - /** * * @param {string} baseUrl @@ -1486,7 +1488,7 @@ export function baseUrlToWS(baseUrl, projectPublicId) { } /** - * @param {import('drizzle-orm/better-sqlite3').BetterSQLite3Database} db + * @param {import('drizzle-orm/better-sqlite3').BetterSQLite3Database} db * @param {import('./datatype/index.js').MapeoDocTables} table * @returns {Stats} */ diff --git a/src/roles.js b/src/roles.js index 7fb84ffbc..887a64b48 100644 --- a/src/roles.js +++ b/src/roles.js @@ -102,7 +102,7 @@ export const CREATOR_ROLE = { /** * @type {Role} */ -const BLOCKED_ROLE = { +export const BLOCKED_ROLE = { roleId: BLOCKED_ROLE_ID, name: 'Blocked', docs: mapObject(currentSchemaVersions, (key) => { diff --git a/src/schema/client.js b/src/schema/client.js index 91a12c90c..dc6bcb44a 100644 --- a/src/schema/client.js +++ b/src/schema/client.js @@ -3,6 +3,7 @@ // device import { blob, sqliteTable, text, int } from 'drizzle-orm/sqlite-core' import { dereferencedDocSchemas as schemas } from '@comapeo/schema' +import { NAMESPACES } from '../constants.js' import { comapeoSchemaToDrizzleTable as toDrizzle, backlinkTable, @@ -49,3 +50,9 @@ export const deviceSettingsTable = sqliteTable('deviceSettings', { (text('deviceInfo', { mode: 'json' })), isArchiveDevice: int('isArchiveDevice', { mode: 'boolean' }), }) + +export const coresTable = sqliteTable('cores', { + projectPublicId: text('projectPublicId').notNull(), + publicKey: blob('publicKey', { mode: 'buffer' }).notNull(), + namespace: text('namespace', { enum: NAMESPACES }).notNull(), +}) diff --git a/src/schema/project.js b/src/schema/project.js index a4bacffd7..dc363aa31 100644 --- a/src/schema/project.js +++ b/src/schema/project.js @@ -32,7 +32,7 @@ export const roleBacklinkTable = backlinkTable('role') export const deviceInfoBacklinkTable = backlinkTable('deviceInfo') export const iconBacklinkTable = backlinkTable('icon') -export const coresTable = sqliteTable('cores', { +export const coresTable_Deprecated = sqliteTable('cores', { publicKey: blob('publicKey', { mode: 'buffer' }).notNull(), namespace: text('namespace', { enum: NAMESPACES }).notNull(), }) diff --git a/src/sync/sync-api.js b/src/sync/sync-api.js index 3c55d9388..e69cf2cac 100644 --- a/src/sync/sync-api.js +++ b/src/sync/sync-api.js @@ -160,7 +160,9 @@ export class SyncApi extends TypedEmitter { }) ) ) - .catch(noop) + .catch((e) => { + this.#l.log('ERROR: Unable to init sync API', e) + }) } /** @type {import('../local-peers.js').LocalPeersEvents['discovery-key']} */ diff --git a/test-e2e/index-worker.js b/test-e2e/index-worker.js new file mode 100644 index 000000000..b870d745d --- /dev/null +++ b/test-e2e/index-worker.js @@ -0,0 +1,94 @@ +import test from 'node:test' +import assert from 'node:assert/strict' +import { performance } from 'node:perf_hooks' +import { createManager } from './utils.js' +import { defaultConfigPath } from '../test/helpers/default-config.js' + +test('indexer-worker - create project', async (t) => { + const manager = createManager('device0', t, { useIndexWorkers: true }) + await manager.getProject( + await manager.createProject({ configPath: defaultConfigPath }) + ) +}) + +test('indexer-worker - perf comparison for config imports', async (t) => { + performance.mark('pre-nonworker', { detail: performance.nodeTiming }) + const manager = createManager('device1', t, { useIndexWorkers: false }) + const project = await manager.getProject( + await manager.createProject({ configPath: defaultConfigPath }) + ) + const presets = await project.preset.getMany() + const fields = await project.field.getMany() + const translations = await project.$translation.dataType.getMany() + assert.equal(presets.length, 28, 'correct number of loaded presets') + assert.equal(fields.length, 11, 'correct number of loaded fields') + assert.equal( + translations.length, + 870, + 'correct number of loaded translations' + ) + performance.mark('post-nonworker', { detail: performance.nodeTiming }) + + performance.mark('pre-worker', { detail: performance.nodeTiming }) + const workerManager = createManager('device2', t, { useIndexWorkers: true }) + const workerProject = await workerManager.getProject( + await workerManager.createProject({ configPath: defaultConfigPath }) + ) + const workerPresets = await workerProject.preset.getMany() + const workerFields = await workerProject.field.getMany() + const workerTranslations = await workerProject.$translation.dataType.getMany() + assert.equal(workerPresets.length, 28, 'correct number of loaded presets') + assert.equal(workerFields.length, 11, 'correct number of loaded fields') + assert.equal( + workerTranslations.length, + 870, + 'correct number of loaded translations' + ) + performance.mark('post-worker', { detail: performance.nodeTiming }) + + const nonworkerDiff = calcDiff( + /** @type Mark */ (performance.getEntriesByName('pre-nonworker')[0]), + /** @type Mark */ (performance.getEntriesByName('post-nonworker')[0]) + ) + + const workerDiff = calcDiff( + /** @type Mark */ (performance.getEntriesByName('pre-worker')[0]), + /** @type Mark */ (performance.getEntriesByName('post-worker')[0]) + ) + + assert( + nonworkerDiff.idleTime < workerDiff.idleTime, + 'Worker idles main thread more' + ) + // TODO: Should add this back when we upgrade from node 18 + // assert(nonworkerDiff.duration > workerDiff.duration, 'Worker is faster') +}) + +/** + * @typedef {Object} MarkDetail + * @property {number} idleTime + */ + +/** + * @typedef {Object} Mark + * @property {number} startTime + * @property {MarkDetail} detail + */ + +/** + * @typedef {Object} Diff + * @property {number} duration + * @property {number} idleTime + */ + +/** + * @param {Mark} mark1 + * @param {Mark} mark2 + * @returns {Diff} + */ +function calcDiff(mark1, mark2) { + const idleTime = mark2.detail.idleTime - mark1.detail.idleTime + const duration = mark2.startTime - mark1.startTime + + return { idleTime, duration } +} diff --git a/test-e2e/manager-basic.js b/test-e2e/manager-basic.js index 785ebf0a9..1992cef08 100644 --- a/test-e2e/manager-basic.js +++ b/test-e2e/manager-basic.js @@ -222,7 +222,8 @@ describe('Consistent loading of config', async () => { }) test('Managing added projects', async (t) => { - const manager = createManager('test', t) + // Workers init slower than this test and lead to race conditions + const manager = createManager('test', t, { useIndexWorkers: false }) const project1Id = await manager.addProject( { @@ -339,6 +340,7 @@ test('Managing both created and added projects', async (t) => { test('Manager cannot add project that already exists', async (t) => { const manager = createManager('test', t) + const existingProjectId = await manager.createProject() const existingProjectsCountBefore = (await manager.listProjects()).length @@ -393,7 +395,9 @@ test('Consistent storage folders', async (t) => { test('Reusing port after start/stop of discovery', async (t) => { const manager = createManager('test', t) - t.after(() => manager.stopLocalPeerDiscoveryServer()) + t.after(() => + manager.stopLocalPeerDiscoveryServer({ force: true, timeout: 0 }) + ) const { port } = await manager.startLocalPeerDiscoveryServer() diff --git a/test-e2e/project-settings.js b/test-e2e/project-settings.js index b276da067..4129b1917 100644 --- a/test-e2e/project-settings.js +++ b/test-e2e/project-settings.js @@ -1,6 +1,5 @@ import test from 'node:test' import assert from 'node:assert/strict' - import { COORDINATOR_ROLE_ID } from '../src/roles.js' import { diff --git a/test-e2e/sync.js b/test-e2e/sync.js index 0aa95872b..a597383a8 100644 --- a/test-e2e/sync.js +++ b/test-e2e/sync.js @@ -910,6 +910,7 @@ test('no sync capabilities === no namespaces sync apart from auth', async (t) => projectId, roleId: BLOCKED_ROLE_ID, }) + await invite({ invitor, invitees: [invitee], diff --git a/test-e2e/utils.js b/test-e2e/utils.js index c4f7baa06..f06d82fe9 100644 --- a/test-e2e/utils.js +++ b/test-e2e/utils.js @@ -15,6 +15,7 @@ import { setTimeout as delay } from 'node:timers/promises' import { getProperty } from 'dot-prop-extra' import { MapeoManager, roles } from '../src/index.js' +import { kWaitForDataStoresIdle } from '../src/mapeo-project.js' import { generate } from '@mapeo/mock-data' import { ExhaustivenessError, valueOf } from '../src/utils.js' import { createHash, randomBytes, randomInt } from 'node:crypto' @@ -280,6 +281,7 @@ export function createManager(seed, t, overrides = {}) { dbFolder, coreStorage, fastify, + useIndexWorkers: true, ...overrides, }) @@ -594,13 +596,17 @@ async function waitForProjectSync(project, peerIds, type = 'initial') { if (hasPeerIds(state.auth.remoteStates, peerIds)) { return project.$sync.waitForSync(type) } - return new Promise((res) => { + const result = await new Promise((res) => { project.$sync[kSyncState].on('state', function onState(state) { if (!hasPeerIds(state.auth.remoteStates, peerIds)) return project.$sync[kSyncState].off('state', onState) res(project.$sync.waitForSync(type)) }) }) + + await project[kWaitForDataStoresIdle]() + + return result } /** diff --git a/test-types/data-types.ts b/test-types/data-types.ts index c723d028a..3b1ef96c7 100644 --- a/test-types/data-types.ts +++ b/test-types/data-types.ts @@ -14,7 +14,7 @@ import { import Database from 'better-sqlite3' import { drizzle } from 'drizzle-orm/better-sqlite3' import RAM from 'random-access-memory' -import { IndexWriter } from '../dist/index-writer/index.js' +import { IndexWriterWrapper as IndexWriter } from '../dist/index-writer/index.js' import { DerivedDocFields } from '../dist/datatype/index.js' import { projectSettingsTable } from '../dist/schema/client.js' import { LocalPeers } from '../dist/local-peers.js' diff --git a/test/core-manager.js b/test/core-manager.js index 363e78649..c222d1e24 100644 --- a/test/core-manager.js +++ b/test/core-manager.js @@ -24,7 +24,8 @@ import path from 'path' import { Transform } from 'streamx' import { waitForCores } from './helpers/core-manager.js' import { drizzle } from 'drizzle-orm/better-sqlite3' -import { coresTable } from '../src/schema/project.js' +import { coresTable } from '../src/schema/client.js' +import * as clientSchema from '../src/schema/client.js' import { eq } from 'drizzle-orm' /** @import { Namespace } from '../src/types.js' */ @@ -216,7 +217,7 @@ test('Added cores are persisted', async () => { const keyManager = new KeyManager(randomBytes(16)) const projectKey = randomBytes(32) - const db = drizzle(new Sqlite(':memory:')) + const db = drizzle(new Sqlite(':memory:'), { schema: clientSchema }) const cm1 = createCoreManager({ db, @@ -525,7 +526,7 @@ test('deleteOthersData()', async (t) => { const peer1TempPath = path.join(tempPath, 'peer1') /// Set up core managers - const db1 = drizzle(new Sqlite(':memory:')) + const db1 = drizzle(new Sqlite(':memory:'), { schema: clientSchema }) const cm1 = createCoreManager({ db: db1, projectKey, @@ -536,7 +537,7 @@ test('deleteOthersData()', async (t) => { autoDownload: true, }) - const db2 = drizzle(new Sqlite(':memory:')) + const db2 = drizzle(new Sqlite(':memory:'), { schema: clientSchema }) const cm2 = createCoreManager({ db: db2, projectKey, diff --git a/test/core-ownership.js b/test/core-ownership.js index bcee185fd..6e575027c 100644 --- a/test/core-ownership.js +++ b/test/core-ownership.js @@ -2,10 +2,8 @@ import test from 'node:test' import assert from 'node:assert/strict' import { KeyManager, sign } from '@mapeo/crypto' import sodium from 'sodium-universal' -import { - mapAndValidateCoreOwnership, - getWinner, -} from '../src/core-ownership.js' +import { mapDoc } from '../src/index-writer/map-doc.js' +import { getWinner } from '../src/index-writer/get-winner.js' import { randomBytes } from 'node:crypto' import { parseVersionId, getVersionId } from '@comapeo/schema' import { discoveryKey } from 'hypercore-crypto' @@ -16,7 +14,7 @@ test('Valid coreOwnership record', () => { const validDoc = generateValidDoc() const version = parseVersionId(validDoc.versionId) - const mappedDoc = mapAndValidateCoreOwnership(validDoc, version) + const mappedDoc = mapDoc(validDoc, version) assert(validDoc.links.length > 0, 'original doc has links') assert.deepEqual(mappedDoc.links, [], 'links are stripped from mapped doc') @@ -42,14 +40,14 @@ test('Invalid coreOwnership signatures', () => { [key]: randomBytes(sodium.crypto_sign_BYTES), }, } - assert.throws(() => mapAndValidateCoreOwnership(invalidDoc, version)) + assert.throws(() => mapDoc(invalidDoc, version)) } const invalidDoc = { ...validDoc, identitySignature: randomBytes(sodium.crypto_sign_BYTES), } - assert.throws(() => mapAndValidateCoreOwnership(invalidDoc, version)) + assert.throws(() => mapDoc(invalidDoc, version)) }) test('Invalid coreOwnership docId and coreIds', () => { @@ -61,14 +59,14 @@ test('Invalid coreOwnership docId and coreIds', () => { ...validDoc, [`${key}CoreId`]: randomBytes(32).toString('hex'), } - assert.throws(() => mapAndValidateCoreOwnership(invalidDoc, version)) + assert.throws(() => mapDoc(invalidDoc, version)) } const invalidDoc = { ...validDoc, docId: randomBytes(32).toString('hex'), } - assert.throws(() => mapAndValidateCoreOwnership(invalidDoc, version)) + assert.throws(() => mapDoc(invalidDoc, version)) }) test('Invalid coreOwnership docId and coreIds (wrong length)', () => { @@ -81,14 +79,14 @@ test('Invalid coreOwnership docId and coreIds (wrong length)', () => { ...validDoc, [`${namespace}CoreId`]: validDoc[`${namespace}CoreId`].slice(0, -1), } - assert.throws(() => mapAndValidateCoreOwnership(invalidDoc, version)) + assert.throws(() => mapDoc(invalidDoc, version)) } const invalidDoc = { ...validDoc, docId: validDoc.docId.slice(0, -1), } - assert.throws(() => mapAndValidateCoreOwnership(invalidDoc, version)) + assert.throws(() => mapDoc(invalidDoc, version)) }) test('Invalid - different coreKey', () => { @@ -97,7 +95,7 @@ test('Invalid - different coreKey', () => { ...parseVersionId(validDoc.versionId), coreDiscoveryKey: discoveryKey(randomBytes(32)), } - assert.throws(() => mapAndValidateCoreOwnership(validDoc, version)) + assert.throws(() => mapDoc(validDoc, version)) }) test('getWinner (coreOwnership)', () => { diff --git a/test/data-type.js b/test/data-type.js index de9d4438a..583bb9d8e 100644 --- a/test/data-type.js +++ b/test/data-type.js @@ -13,8 +13,9 @@ import { trackTable, translationTable, } from '../src/schema/project.js' +import * as projectSchema from '../src/schema/project.js' import { DataType, kCreateWithDocId } from '../src/datatype/index.js' -import { IndexWriter } from '../src/index-writer/index.js' +import { IndexWriterWrapper } from '../src/index-writer/index.js' import { NotFoundError } from '../src/errors.js' import Database from 'better-sqlite3' @@ -56,13 +57,13 @@ const trackFixture = valueOf(generate('track')[0]) test('private createWithDocId() method', async () => { const sqlite = new Database(':memory:') - const db = drizzle(sqlite) + const db = drizzle(sqlite, { schema: projectSchema }) migrate(db, { migrationsFolder: new URL('../drizzle/project', import.meta.url).pathname, }) const coreManager = createCoreManager() - const indexWriter = new IndexWriter({ + const indexWriter = new IndexWriterWrapper({ tables: [observationTable], sqlite, }) @@ -95,13 +96,13 @@ test('private createWithDocId() method', async () => { test('private createWithDocId() method throws when doc exists', async () => { const sqlite = new Database(':memory:') - const db = drizzle(sqlite) + const db = drizzle(sqlite, { schema: projectSchema }) migrate(db, { migrationsFolder: new URL('../drizzle/project', import.meta.url).pathname, }) const coreManager = createCoreManager() - const indexWriter = new IndexWriter({ + const indexWriter = new IndexWriterWrapper({ tables: [observationTable], sqlite, }) @@ -359,14 +360,14 @@ test('translation', async () => { */ async function testenv(opts = {}) { const sqlite = new Database(':memory:') - const db = drizzle(sqlite) + const db = drizzle(sqlite, { schema: projectSchema }) migrate(db, { migrationsFolder: new URL('../drizzle/project', import.meta.url).pathname, }) - const coreManager = createCoreManager({ ...opts, db }) + const coreManager = createCoreManager({ ...opts }) - const indexWriter = new IndexWriter({ + const indexWriter = new IndexWriterWrapper({ tables: [observationTable, trackTable, translationTable], sqlite, }) diff --git a/test/helpers/core-manager.js b/test/helpers/core-manager.js index 88c78c44a..bfc4b2270 100644 --- a/test/helpers/core-manager.js +++ b/test/helpers/core-manager.js @@ -10,6 +10,7 @@ import NoiseSecretStream from '@hyperswarm/secret-stream' import { drizzle } from 'drizzle-orm/better-sqlite3' import { migrate } from 'drizzle-orm/better-sqlite3/migrator' import { NAMESPACES } from '../../src/constants.js' +import * as clientSchema from '../../src/schema/client.js' /** @import { Namespace } from '../../src/types.js' */ /** @@ -20,12 +21,11 @@ import { NAMESPACES } from '../../src/constants.js' export function createCoreManager({ rootKey = randomBytes(16), projectKey = randomBytes(32), - db = drizzle(new Sqlite(':memory:')), + db = drizzle(new Sqlite(':memory:'), { schema: clientSchema }), ...opts } = {}) { migrate(db, { - migrationsFolder: new URL('../../drizzle/project', import.meta.url) - .pathname, + migrationsFolder: new URL('../../drizzle/client', import.meta.url).pathname, }) const keyManager = new KeyManager(rootKey) diff --git a/test/icon-api.js b/test/icon-api.js index 33b81994b..be1db589f 100644 --- a/test/icon-api.js +++ b/test/icon-api.js @@ -16,7 +16,9 @@ import { DataType } from '../src/datatype/index.js' import { DataStore } from '../src/datastore/index.js' import { createCoreManager } from './helpers/core-manager.js' import { iconTable } from '../src/schema/project.js' -import { IndexWriter } from '../src/index-writer/index.js' +import { IndexWriterWrapper } from '../src/index-writer/index.js' +import * as clientSchema from '../src/schema/client.js' +import * as projectSchema from '../src/schema/project.js' test('create()', async () => { const { iconApi, iconDataType } = setup() @@ -665,18 +667,23 @@ test('constructIconPath() - good inputs', () => { function setup({ getMediaBaseUrl = async () => 'http://127.0.0.1:8080/icons', } = {}) { - const sqlite = new Database(':memory:') - const db = drizzle(sqlite) + const clientSqlite = new Database(':memory:') + const projectSqlite = new Database(':memory:') + const clientDb = drizzle(clientSqlite, { schema: clientSchema }) + const projectDb = drizzle(projectSqlite, { schema: projectSchema }) - migrate(db, { + migrate(clientDb, { + migrationsFolder: new URL('../drizzle/client', import.meta.url).pathname, + }) + migrate(projectDb, { migrationsFolder: new URL('../drizzle/project', import.meta.url).pathname, }) - const cm = createCoreManager({ db }) + const cm = createCoreManager({ db: clientDb }) - const indexWriter = new IndexWriter({ + const indexWriter = new IndexWriterWrapper({ tables: [iconTable], - sqlite, + sqlite: projectSqlite, }) const iconDataStore = new DataStore({ @@ -690,7 +697,7 @@ function setup({ const iconDataType = new DataType({ dataStore: iconDataStore, table: iconTable, - db, + db: projectDb, getTranslations() { throw new Error('Translations should not be fetched in this test') }, diff --git a/test/translation-api.js b/test/translation-api.js index 75146b98a..ae20bfc5c 100644 --- a/test/translation-api.js +++ b/test/translation-api.js @@ -10,7 +10,9 @@ import Database from 'better-sqlite3' import { drizzle } from 'drizzle-orm/better-sqlite3' import { migrate } from 'drizzle-orm/better-sqlite3/migrator' import { createCoreManager } from './helpers/core-manager.js' -import { IndexWriter } from '../src/index-writer/index.js' +import { IndexWriterWrapper } from '../src/index-writer/index.js' +import * as clientSchema from '../src/schema/client.js' +import * as projectSchema from '../src/schema/project.js' import RAM from 'random-access-memory' import { hashObject } from '../src/utils.js' import { omit } from '../src/lib/omit.js' @@ -118,18 +120,23 @@ test('translation api - put() and get()', async () => { }) function setup() { - const sqlite = new Database(':memory:') - const db = drizzle(sqlite) + const clientSqlite = new Database(':memory:') + const projectSqlite = new Database(':memory:') + const clientDb = drizzle(clientSqlite, { schema: clientSchema }) + const projectDb = drizzle(projectSqlite, { schema: projectSchema }) - migrate(db, { + migrate(clientDb, { + migrationsFolder: new URL('../drizzle/client', import.meta.url).pathname, + }) + migrate(projectDb, { migrationsFolder: new URL('../drizzle/project', import.meta.url).pathname, }) - const cm = createCoreManager({ db }) + const cm = createCoreManager({ db: clientDb }) - const indexWriter = new IndexWriter({ + const indexWriter = new IndexWriterWrapper({ tables: [table], - sqlite, + sqlite: projectSqlite, }) const dataStore = new DataStore({ @@ -143,7 +150,7 @@ function setup() { const dataType = new DataType({ dataStore, table, - db, + db: projectDb, getTranslations() { throw new Error('Cannot get translations from translations') },