diff --git a/comapeo-schema-PR267+bc403d5.tgz b/comapeo-schema-PR267+bc403d5.tgz new file mode 100644 index 000000000..c3cc3cfd6 Binary files /dev/null and b/comapeo-schema-PR267+bc403d5.tgz differ diff --git a/multi-core-indexer-PR25+2537f7d.tgz b/multi-core-indexer-PR25+2537f7d.tgz new file mode 100644 index 000000000..fe828202f Binary files /dev/null and b/multi-core-indexer-PR25+2537f7d.tgz differ diff --git a/package-lock.json b/package-lock.json index 6c4516853..d3846c374 100644 --- a/package-lock.json +++ b/package-lock.json @@ -9,7 +9,7 @@ "version": "1.0.1", "license": "MIT", "dependencies": { - "@comapeo/schema": "1.0.0", + "@comapeo/schema": "file:comapeo-schema-PR267+bc403d5.tgz", "@digidem/types": "^2.3.0", "@electron/asar": "^3.2.8", "@fastify/error": "^3.4.1", @@ -39,7 +39,7 @@ "magic-bytes.js": "^1.10.0", "map-obj": "^5.0.2", "mime": "^4.0.3", - "multi-core-indexer": "^1.0.0-alpha.10", + "multi-core-indexer": "file:multi-core-indexer-PR25+2537f7d.tgz", "p-defer": "^4.0.0", "p-event": "^6.0.1", "p-timeout": "^6.1.2", @@ -243,8 +243,9 @@ }, "node_modules/@comapeo/schema": { "version": "1.0.0", - "resolved": "https://registry.npmjs.org/@comapeo/schema/-/schema-1.0.0.tgz", - "integrity": "sha512-dK227I+0yg9D2y5/O5NGywx50tgeNYyUkl1uYnSmNAPlbv+r2KX9aaC9m4dEjIja2aR2VFnYn6z537ERZiahqQ==", + "resolved": "file:comapeo-schema-PR267+bc403d5.tgz", + "integrity": "sha512-9NgqM+xxDbDll4LSlK6akfat2qe48sofXOaQHUh9pKT6pVAmEaP3dMLpvixX6ddQBhFcVNz5bK1vmloq3vRdkA==", + "license": "MIT", "dependencies": { "compact-encoding": "^2.12.0", "protobufjs": "^7.2.5", @@ -5238,8 +5239,9 @@ }, "node_modules/multi-core-indexer": { "version": "1.0.0-alpha.10", - "resolved": "https://registry.npmjs.org/multi-core-indexer/-/multi-core-indexer-1.0.0-alpha.10.tgz", - "integrity": "sha512-H9QdpJ/MaelrBZw6jCcsrInE+hwUQmfz/2swtIdQNNh1IHUDGEdPkakjcZAyahpM5iIVz7EqyWO74aC03A3qSA==", + "resolved": "file:multi-core-indexer-PR25+2537f7d.tgz", + "integrity": "sha512-B8Iewe2wajm7Jxa+FHxkHTGTTNbXejci5hsegQ+4ypcsyGa2xDyfRM7xKx/b6mESJPD/NhTmIsPzNfjPA78isQ==", + "license": "MIT", "dependencies": { "@types/node": "^18.16.19", "@types/streamx": "^2.9.1", diff --git a/package.json b/package.json index 657338a8e..d09f4f36d 100644 --- a/package.json +++ b/package.json @@ -151,7 +151,7 @@ "yazl": "^2.5.1" }, "dependencies": { - "@comapeo/schema": "1.0.0", + "@comapeo/schema": "file:comapeo-schema-PR267+bc403d5.tgz", "@digidem/types": "^2.3.0", "@electron/asar": "^3.2.8", "@fastify/error": "^3.4.1", @@ -181,7 +181,7 @@ "magic-bytes.js": "^1.10.0", "map-obj": "^5.0.2", "mime": "^4.0.3", - "multi-core-indexer": "^1.0.0-alpha.10", + "multi-core-indexer": "file:multi-core-indexer-PR25+2537f7d.tgz", "p-defer": "^4.0.0", "p-event": "^6.0.1", "p-timeout": "^6.1.2", diff --git a/src/core-manager/core-index.js b/src/core-manager/core-index.js index 79bb8de6f..d5e713204 100644 --- a/src/core-manager/core-index.js +++ b/src/core-manager/core-index.js @@ -22,13 +22,12 @@ export class CoreIndex { * @param {Object} options * @param {import('hypercore')<"binary", Buffer>} options.core Hypercore instance * @param {Buffer} options.key Buffer containing public key of this core + * @param {string} options.discoveryId discoveryId of this core * @param {Namespace} options.namespace * @param {boolean} [options.writer] Is this a writer core? */ - add({ core, key, namespace, writer = false }) { - const discoveryKey = crypto.discoveryKey(key) - const discoveryId = discoveryKey.toString('hex') - const record = { core, key, namespace } + add({ core, key, discoveryId, namespace, writer = false }) { + const record = { core, key, discoveryId, namespace } if (writer) { this.#writersByNamespace.set(namespace, record) } diff --git a/src/core-manager/index.js b/src/core-manager/index.js index 68c9a3261..3b05badd3 100644 --- a/src/core-manager/index.js +++ b/src/core-manager/index.js @@ -3,6 +3,7 @@ import Corestore from 'corestore' import { debounce } from 'throttle-debounce' import assert from 'node:assert/strict' import { sql, eq } from 'drizzle-orm' +import hCrypto from 'hypercore-crypto' import { HaveExtension, ProjectExtension } from '../generated/extensions.js' import { Logger } from '../logger.js' @@ -20,7 +21,7 @@ const WRITER_CORE_PREHAVES_DEBOUNCE_DELAY = 1000 export const kCoreManagerReplicate = Symbol('replicate core manager') /** @typedef {Hypercore<'binary', Buffer>} Core */ -/** @typedef {{ core: Core, key: Buffer, namespace: Namespace }} CoreRecord */ +/** @typedef {{ core: Core, key: Buffer, discoveryId: string, namespace: Namespace }} CoreRecord */ /** * @typedef {Object} Events * @property {(coreRecord: CoreRecord) => void} add-core @@ -273,6 +274,7 @@ export class CoreManager extends TypedEmitter { if (existingCore) return existingCore const { publicKey: key, secretKey } = keyPair + const discoveryId = hCrypto.discoveryKey(key).toString('hex') const writer = !!secretKey const core = this.#corestore.get({ keyPair, @@ -285,7 +287,7 @@ export class CoreManager extends TypedEmitter { core.setMaxListeners(0) // @ts-ignore - ensure key is defined before hypercore is ready core.key = key - this.#coreIndex.add({ core, key, namespace, writer }) + this.#coreIndex.add({ core, key, discoveryId, namespace, writer }) // **Hack** As soon as a peer is added, eagerly send a "want" for the entire // core. This ensures that the peer sends back its entire bitfield. @@ -341,9 +343,9 @@ export class CoreManager extends TypedEmitter { namespace, key ) - this.emit('add-core', { core, key, namespace }) + this.emit('add-core', { core, key, discoveryId, namespace }) - return { core, key, namespace } + return { core, key, discoveryId, namespace } } /** diff --git a/src/core-ownership.js b/src/core-ownership.js index 38a595d28..5157c8349 100644 --- a/src/core-ownership.js +++ b/src/core-ownership.js @@ -155,12 +155,13 @@ export class CoreOwnership extends TypedEmitter { * the doc with the lowest index (e.g. the first) * * @param {CoreOwnershipWithSignatures} doc - * @param {import('@comapeo/schema').VersionIdObject} version + * @param {import('@comapeo/schema').VersionDiscoveryIdObject} version * @returns {import('@comapeo/schema').CoreOwnership} */ -export function mapAndValidateCoreOwnership(doc, { coreDiscoveryKey }) { +export function mapAndValidateCoreOwnership(doc, { coreDiscoveryId }) { if ( - !coreDiscoveryKey.equals(discoveryKey(Buffer.from(doc.authCoreId, 'hex'))) + coreDiscoveryId !== + discoveryKey(Buffer.from(doc.authCoreId, 'hex')).toString('hex') ) { throw new Error('Invalid coreOwnership record: mismatched authCoreId') } diff --git a/src/datastore/index.js b/src/datastore/index.js index 43712654f..f321d75c5 100644 --- a/src/datastore/index.js +++ b/src/datastore/index.js @@ -2,7 +2,6 @@ import { TypedEmitter } from 'tiny-typed-emitter' import { encode, decode, getVersionId, parseVersionId } from '@comapeo/schema' import MultiCoreIndexer from 'multi-core-indexer' import pDefer from 'p-defer' -import { discoveryKey } from 'hypercore-crypto' import { NAMESPACE_SCHEMAS } from '../constants.js' import { createMap } from '../utils.js' /** @import { MapeoDoc } from '@comapeo/schema' */ @@ -36,7 +35,7 @@ export class DataStore extends TypedEmitter { #coreManager #namespace #batch - #writerCore + #writerCoreRecord #coreIndexer /** @type {Map>} */ #pendingIndex = new Map() @@ -61,7 +60,7 @@ export class DataStore extends TypedEmitter { NAMESPACE_SCHEMAS[namespace], () => new Set() ) - this.#writerCore = coreManager.getWriterCore(namespace).core + this.#writerCoreRecord = coreManager.getWriterCore(namespace) const cores = coreManager.getCores(namespace).map((cr) => cr.core) this.#coreIndexer = new MultiCoreIndexer(cores, { storage, @@ -88,7 +87,7 @@ export class DataStore extends TypedEmitter { } get writerCore() { - return this.#writerCore + return this.#writerCoreRecord.core } getIndexState() { @@ -105,9 +104,11 @@ export class DataStore extends TypedEmitter { // Writes to the writerCore need to wait until the entry is indexed before // returning, so we check if any incoming entry has a pending promise for (const entry of entries) { - if (!entry.key.equals(this.#writerCore.key)) continue + if (entry.discoveryId !== this.#writerCoreRecord.discoveryId) { + continue + } const versionId = getVersionId({ - coreDiscoveryKey: discoveryKey(entry.key), + coreDiscoveryId: entry.discoveryId, index: entry.index, }) const pending = this.#pendingIndex.get(versionId) @@ -153,23 +154,20 @@ export class DataStore extends TypedEmitter { // same tick, so we can't know their index before append resolves. const deferredAppend = pDefer() this.#pendingAppends.add(deferredAppend.promise) - const { length } = await this.#writerCore.append(block) + const { length } = await this.writerCore.append(block) deferredAppend.resolve() this.#pendingAppends.delete(deferredAppend.promise) const index = length - 1 - const coreDiscoveryKey = this.#writerCore.discoveryKey - if (!coreDiscoveryKey) { - throw new Error('Writer core is not ready') - } - const versionId = getVersionId({ coreDiscoveryKey, index }) + const coreDiscoveryId = this.#writerCoreRecord.discoveryId + const versionId = getVersionId({ coreDiscoveryId, index }) /** @type {import('p-defer').DeferredPromise} */ const deferred = pDefer() this.#pendingIndex.set(versionId, deferred) await deferred.promise return /** @type {Extract} */ ( - decode(block, { coreDiscoveryKey, index }) + decode(block, { coreDiscoveryId, index }) ) } @@ -189,13 +187,10 @@ export class DataStore extends TypedEmitter { /** @param {Buffer} buf */ async writeRaw(buf) { - const { length } = await this.#writerCore.append(buf) + const { length } = await this.writerCore.append(buf) const index = length - 1 - const coreDiscoveryKey = this.#writerCore.discoveryKey - if (!coreDiscoveryKey) { - throw new Error('Writer core is not ready') - } - const versionId = getVersionId({ coreDiscoveryKey, index }) + const coreDiscoveryId = this.#writerCoreRecord.discoveryId + const versionId = getVersionId({ coreDiscoveryId, index }) return versionId } diff --git a/src/index-writer/index.js b/src/index-writer/index.js index ae643b99f..813c5302d 100644 --- a/src/index-writer/index.js +++ b/src/index-writer/index.js @@ -2,9 +2,8 @@ import { decode } from '@comapeo/schema' import SqliteIndexer from '@mapeo/sqlite-indexer' import { getTableConfig } from 'drizzle-orm/sqlite-core' import { getBacklinkTableName } from '../schema/utils.js' -import { discoveryKey } from 'hypercore-crypto' import { Logger } from '../logger.js' -/** @import { MapeoDoc, VersionIdObject } from '@comapeo/schema' */ +/** @import { MapeoDoc, VersionDiscoveryIdObject } from '@comapeo/schema' */ /** @import { MapeoDocTables } from '../datatype/index.js' */ /** @@ -32,7 +31,7 @@ export class 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 {(doc: MapeoDocInternal, version: VersionDiscoveryIdObject) => 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] */ @@ -71,13 +70,13 @@ export class IndexWriter { const queued = {} /** @type {IndexedDocIds} */ const indexed = {} - for (const { block, key, index } of entries) { + for (const { block, discoveryId, index } of entries) { /** @type {MapeoDoc} */ let doc try { - const version = { coreDiscoveryKey: discoveryKey(key), index } + const version = { coreDiscoveryId: discoveryId, index } doc = this.#mapDoc(decode(block, version), version) } catch (e) { - this.#l.log('Could not decode entry %d of %h', index, key) + this.#l.log('Could not decode entry %d of %S', index, discoveryId) // Unknown or invalid entry - silently ignore continue } diff --git a/src/mapeo-project.js b/src/mapeo-project.js index dc92ba03b..19c33eb03 100644 --- a/src/mapeo-project.js +++ b/src/mapeo-project.js @@ -461,7 +461,7 @@ export class MapeoProject extends TypedEmitter { projectSettingsEntries.push(entry) } else if (schemaName === 'translation') { const doc = decode(entry.block, { - coreDiscoveryKey: entry.key, + coreDiscoveryId: entry.discoveryId, index: entry.index, }) @@ -966,11 +966,14 @@ function getCoreKeypairs({ projectKey, projectSecretKey, keyManager }) { * e.g. version.coreKey should equal docId * * @param {import('@comapeo/schema').DeviceInfo} doc - * @param {import('@comapeo/schema').VersionIdObject} version + * @param {import('@comapeo/schema').VersionDiscoveryIdObject} version * @returns {import('@comapeo/schema').DeviceInfo} */ -function mapAndValidateDeviceInfo(doc, { coreDiscoveryKey }) { - if (!coreDiscoveryKey.equals(discoveryKey(Buffer.from(doc.docId, 'hex')))) { +function mapAndValidateDeviceInfo(doc, { coreDiscoveryId }) { + if ( + coreDiscoveryId !== + discoveryKey(Buffer.from(doc.docId, 'hex')).toString('hex') + ) { throw new Error( 'Invalid deviceInfo record, cannot write deviceInfo for another device' ) diff --git a/test-e2e/config-import.js b/test-e2e/config-import.js index 55e0420ed..5ef4e12a1 100644 --- a/test-e2e/config-import.js +++ b/test-e2e/config-import.js @@ -3,22 +3,26 @@ import assert from 'node:assert/strict' import { createManager } from './utils.js' import { defaultConfigPath } from '../tests/helpers/default-config.js' -test(' config import - load default config when passed a path to `createProject`', async (t) => { - const manager = createManager('device0', t) - 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' - ) -}) +test( + ' config import - load default config when passed a path to `createProject`', + { only: true }, + async (t) => { + const manager = createManager('device0', t) + 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' + ) + } +) test('config import - load and re-load config manually', async (t) => { const manager = createManager('device0', t) diff --git a/tests/core-ownership.js b/tests/core-ownership.js index bcee185fd..6c449b1ec 100644 --- a/tests/core-ownership.js +++ b/tests/core-ownership.js @@ -7,14 +7,14 @@ import { getWinner, } from '../src/core-ownership.js' import { randomBytes } from 'node:crypto' -import { parseVersionId, getVersionId } from '@comapeo/schema' +import { getVersionId } from '@comapeo/schema' import { discoveryKey } from 'hypercore-crypto' /** @import { Namespace } from '../src/types.js' */ test('Valid coreOwnership record', () => { const validDoc = generateValidDoc() - const version = parseVersionId(validDoc.versionId) + const version = parseVersionIdInternal(validDoc.versionId) const mappedDoc = mapAndValidateCoreOwnership(validDoc, version) @@ -32,7 +32,7 @@ test('Valid coreOwnership record', () => { test('Invalid coreOwnership signatures', () => { const validDoc = generateValidDoc() - const version = parseVersionId(validDoc.versionId) + const version = parseVersionIdInternal(validDoc.versionId) for (const key of Object.keys(validDoc.coreSignatures)) { const invalidDoc = { @@ -54,7 +54,7 @@ test('Invalid coreOwnership signatures', () => { test('Invalid coreOwnership docId and coreIds', () => { const validDoc = generateValidDoc() - const version = parseVersionId(validDoc.versionId) + const version = parseVersionIdInternal(validDoc.versionId) for (const key of Object.keys(validDoc.coreSignatures)) { const invalidDoc = { @@ -73,7 +73,7 @@ test('Invalid coreOwnership docId and coreIds', () => { test('Invalid coreOwnership docId and coreIds (wrong length)', () => { const validDoc = generateValidDoc() - const version = parseVersionId(validDoc.versionId) + const version = parseVersionIdInternal(validDoc.versionId) for (const key of Object.keys(validDoc.coreSignatures)) { const namespace = /** @type {Namespace} */ (key) @@ -94,15 +94,15 @@ test('Invalid coreOwnership docId and coreIds (wrong length)', () => { test('Invalid - different coreKey', () => { const validDoc = generateValidDoc() const version = { - ...parseVersionId(validDoc.versionId), - coreDiscoveryKey: discoveryKey(randomBytes(32)), + ...parseVersionIdInternal(validDoc.versionId), + coreDiscoveryId: discoveryKey(randomBytes(32)).toString('hex'), } assert.throws(() => mapAndValidateCoreOwnership(validDoc, version)) }) test('getWinner (coreOwnership)', () => { const validDoc = generateValidDoc() - const version = parseVersionId(validDoc.versionId) + const version = parseVersionIdInternal(validDoc.versionId) const docA = { ...validDoc, @@ -218,3 +218,17 @@ function generateValidDoc() { } return validDoc } + +/** + * Tests need to parse a versionId to a discoveryId and index, vs the public + * method parseVersionId which parses to a discoveryKey and index. + * + * @param {string} versionId + */ +function parseVersionIdInternal(versionId) { + const [coreDiscoveryId, index] = versionId.split('/') + return { + coreDiscoveryId, + index: Number.parseInt(index), + } +} diff --git a/tests/data-type.js b/tests/data-type.js index fd78752b9..997510a48 100644 --- a/tests/data-type.js +++ b/tests/data-type.js @@ -329,7 +329,7 @@ async function testenv(opts = {}) { try { if (schemaName === 'translation') { const doc = decode(entry.block, { - coreDiscoveryKey: entry.key, + coreDiscoveryId: entry.discoveryId, index: entry.index, }) assert( diff --git a/tests/datastore.js b/tests/datastore.js index 57440d375..780210c14 100644 --- a/tests/datastore.js +++ b/tests/datastore.js @@ -37,9 +37,8 @@ test('read and write', async () => { coreManager: cm, namespace: 'data', batch: async (entries) => { - for (const { index, key } of entries) { - const coreDiscoveryKey = discoveryKey(key) - const versionId = getVersionId({ coreDiscoveryKey, index }) + for (const { index, discoveryId } of entries) { + const versionId = discoveryId + '/' + index indexedVersionIds.push(versionId) } return {}