diff --git a/packages/gatekeeper/src/db/abstract-json.ts b/packages/gatekeeper/src/db/abstract-json.ts index 896961171..41ba755ee 100644 --- a/packages/gatekeeper/src/db/abstract-json.ts +++ b/packages/gatekeeper/src/db/abstract-json.ts @@ -83,6 +83,10 @@ export abstract class AbstractJson implements GatekeeperDb { } async addEvent(did: string, event: GatekeeperEvent): Promise { + return this.addEventAndQueue(did, event, []); + } + + async addEventAndQueue(did: string, event: GatekeeperEvent, queueRegistries: string[]): Promise { const suffix = this.splitSuffix(did); return this.runExclusive(async () => { const db = this.loadDb(); @@ -96,6 +100,11 @@ export abstract class AbstractJson implements GatekeeperDb { did, event, }); + for (const registry of queueRegistries) { + db.queue ??= {}; + db.queue[registry] ??= []; + db.queue[registry].push(event.operation); + } this.writeDb(db); }); } diff --git a/packages/gatekeeper/src/db/mongo.ts b/packages/gatekeeper/src/db/mongo.ts index a114dba23..07c42b106 100644 --- a/packages/gatekeeper/src/db/mongo.ts +++ b/packages/gatekeeper/src/db/mongo.ts @@ -32,6 +32,12 @@ interface QueueDoc { ops: Operation[] } +interface DuplicateQueueDoc { + _id: string + operationBatches: Operation[][] + count: number +} + interface CounterDoc { id: string; value: number; @@ -57,6 +63,13 @@ function isNamespaceNotFoundError(error: unknown): boolean { ); } +function isDuplicateKeyError(error: unknown): boolean { + return typeof error === 'object' + && error !== null + && 'code' in error + && error.code === 11000; +} + const MONGO_NOT_STARTED_ERROR = 'Mongo not started. Call start() first.'; const MONGO_TRANSACTIONS_REQUIRED_ERROR = 'MongoDB transactions require a replica set or sharded cluster. Configure KC_MONGODB_URL to point at a transaction-capable MongoDB deployment.'; const log = childLogger({ service: 'gatekeeper-db', module: 'mongo' }); @@ -91,6 +104,7 @@ export default class DbMongo implements GatekeeperDb { await this.verifyTransactionSupport(); this.db = this.client.db(this.dbName); await this.ensureIndex('dids', { id: 1 }, { name: 'dids_id_unique', unique: true }); + await this.ensureQueueIndex(); await this.ensureIndex('blocks', { registry: 1, height: -1 }, { name: 'blocks_registry_height' }); // for latest and height lookups await this.ensureIndex('blocks', { registry: 1, hash: 1 }, { name: 'blocks_registry_hash_unique', unique: true }); // for hash lookup await this.ensureIndex('counters', { id: 1 }, { name: 'counters_id_unique', unique: true }); @@ -204,6 +218,53 @@ export default class DbMongo implements GatekeeperDb { await collection.createIndex(key, options); } + private async ensureQueueIndex(): Promise { + try { + await this.ensureIndex('queue', { id: 1 }, { name: 'queue_id_unique', unique: true }); + } + catch (error) { + if (!isDuplicateKeyError(error)) { + throw error; + } + + await this.mergeDuplicateQueues(); + await this.ensureIndex('queue', { id: 1 }, { name: 'queue_id_unique', unique: true }); + } + } + + private async mergeDuplicateQueues(): Promise { + if (!this.db) { + throw new Error(MONGO_NOT_STARTED_ERROR); + } + + const collection = this.db.collection('queue'); + const duplicates = await collection.aggregate([ + { $sort: { _id: 1 } }, + { + $group: { + _id: '$id', + operationBatches: { $push: '$ops' }, + count: { $sum: 1 }, + } + }, + { $match: { count: { $gt: 1 } } }, + ]).toArray(); + + if (duplicates.length === 0) { + return; + } + + await this.withTransaction(async session => { + for (const duplicate of duplicates) { + await collection.deleteMany({ id: duplicate._id }, { session }); + await collection.insertOne({ + id: duplicate._id, + ops: duplicate.operationBatches.flat(), + }, { session }); + } + }); + } + private isValidCounterValue(value: unknown): value is number { return typeof value === 'number' && Number.isSafeInteger(value) && value >= 0; } @@ -325,6 +386,10 @@ export default class DbMongo implements GatekeeperDb { } async addEvent(did: string, event: GatekeeperEvent): Promise { + return this.addEventAndQueue(did, event, []); + } + + async addEventAndQueue(did: string, event: GatekeeperEvent, queueRegistries: string[]): Promise { if (!this.db) { throw new Error(MONGO_NOT_STARTED_ERROR) } @@ -349,6 +414,13 @@ export default class DbMongo implements GatekeeperDb { did, event, }, session); + for (const registry of queueRegistries) { + await this.db!.collection('queue').updateOne( + { id: registry }, + { $push: { ops: event.operation } }, + { upsert: true, session } + ); + } return count }); } diff --git a/packages/gatekeeper/src/db/postgres.ts b/packages/gatekeeper/src/db/postgres.ts index b64630bbc..48f2fccce 100644 --- a/packages/gatekeeper/src/db/postgres.ts +++ b/packages/gatekeeper/src/db/postgres.ts @@ -422,10 +422,14 @@ export default class DbPostgres implements GatekeeperDb { }); } - async addEvent(did: string, event: GatekeeperEvent): Promise { - const id = this.splitSuffix(did); + private async addEventStrict( + executor: Pool | PoolClient, + did: string, + id: string, + event: GatekeeperEvent + ): Promise { const serializedEvent = JSON.stringify(event); - const result = await this.getPool().query( + const result = await executor.query( `WITH inserted_event AS ( INSERT INTO gatekeeper_events (namespace, id, seq, event) SELECT $1, $2, COALESCE(MAX(seq), -1) + 1, $3::jsonb @@ -447,6 +451,32 @@ export default class DbPostgres implements GatekeeperDb { return this.toNumber(result.rows[0]?.length ?? 0); } + async addEvent(did: string, event: GatekeeperEvent): Promise { + const id = this.splitSuffix(did); + return this.addEventStrict(this.getPool(), did, id, event); + } + + async addEventAndQueue(did: string, event: GatekeeperEvent, queueRegistries: string[]): Promise { + const id = this.splitSuffix(did); + if (queueRegistries.length === 0) { + return this.addEventStrict(this.getPool(), did, id, event); + } + + return this.withTx(async client => { + const count = await this.addEventStrict(client, did, id, event); + for (const registry of queueRegistries) { + await client.query( + `INSERT INTO gatekeeper_queue (namespace, id, ops) + VALUES ($1, $2, $3::jsonb) + ON CONFLICT (namespace, id) + DO UPDATE SET ops = COALESCE(gatekeeper_queue.ops, '[]'::jsonb) || EXCLUDED.ops`, + [this.dbName, registry, JSON.stringify([event.operation])] + ); + } + return count; + }); + } + async setEvents(did: string, events: GatekeeperEvent[], options?: SetEventsOptions): Promise { const id = this.splitSuffix(did); diff --git a/packages/gatekeeper/src/db/redis.ts b/packages/gatekeeper/src/db/redis.ts index 5fc79bd36..14aa1980e 100644 --- a/packages/gatekeeper/src/db/redis.ts +++ b/packages/gatekeeper/src/db/redis.ts @@ -129,6 +129,10 @@ export default class DbRedis implements GatekeeperDb { } async addEvent(did: string, event: GatekeeperEvent): Promise { + return this.addEventAndQueue(did, event, []); + } + + async addEventAndQueue(did: string, event: GatekeeperEvent, queueRegistries: string[]): Promise { if (!this.redis) { throw new Error(REDIS_NOT_STARTED_ERROR) } @@ -141,21 +145,25 @@ export default class DbRedis implements GatekeeperDb { did, event, }); + const queueKeys = queueRegistries.map(registry => this.queueKey(registry)); const script = ` - ${this.checkRedisTypesScript(['list', 'string', 'zset', 'zset'])} + ${this.checkRedisTypesScript(['list', 'string', 'zset', 'zset', ...queueKeys.map(() => 'list')])} local change = cjson.decode(ARGV[2]) local seq = redis.call('INCR', KEYS[2]) local count = redis.call('RPUSH', KEYS[1], ARGV[1]) change.seq = seq redis.call('ZADD', KEYS[3], seq, cjson.encode(change)) redis.call('ZADD', KEYS[4], 0, ARGV[3]) + for i = 5, #KEYS do + redis.call('RPUSH', KEYS[i], ARGV[4]) + end return count `; const result = await this.evalAtomicMutation( script, - [key, this.indexSeqKey(), this.indexChangesKey(), this.didIndexKey()], - [val, change, id] + [key, this.indexSeqKey(), this.indexChangesKey(), this.didIndexKey(), ...queueKeys], + [val, change, id, ...(queueKeys.length > 0 ? [JSON.stringify(event.operation)] : [])] ); return Number(result ?? 0); diff --git a/packages/gatekeeper/src/db/sqlite.ts b/packages/gatekeeper/src/db/sqlite.ts index 8dd7a2d80..61fb2eb97 100644 --- a/packages/gatekeeper/src/db/sqlite.ts +++ b/packages/gatekeeper/src/db/sqlite.ts @@ -192,6 +192,10 @@ export default class DbSqlite implements GatekeeperDb { } async addEvent(did: string, event: GatekeeperEvent): Promise { + return this.addEventAndQueue(did, event, []); + } + + async addEventAndQueue(did: string, event: GatekeeperEvent, queueRegistries: string[]): Promise { if (!did) { throw new InvalidDIDError(); } @@ -207,6 +211,9 @@ export default class DbSqlite implements GatekeeperDb { did, event, }); + for (const registry of queueRegistries) { + await this.queueOperationStrict(registry, event.operation); + } return changes; }) ); @@ -323,22 +330,24 @@ export default class DbSqlite implements GatekeeperDb { ); } - async queueOperation(registry: string, op: Operation): Promise { + private async queueOperationStrict(registry: string, op: Operation): Promise { if (!this.db) { - throw new Error(SQLITE_NOT_STARTED_ERROR) + throw new Error(SQLITE_NOT_STARTED_ERROR); } - return this.runExclusive(async () => - this.withTx(async () => { - const ops = await this.getQueueStrict(registry); - ops.push(op); - await this.db!.run( - `INSERT OR REPLACE INTO queue(id, ops) VALUES(?, ?)`, - registry, - JSON.stringify(ops) - ); - return ops.length; - }) + const ops = await this.getQueueStrict(registry); + ops.push(op); + await this.db.run( + `INSERT OR REPLACE INTO queue(id, ops) VALUES(?, ?)`, + registry, + JSON.stringify(ops) + ); + return ops.length; + } + + async queueOperation(registry: string, op: Operation): Promise { + return this.runExclusive(() => + this.withTx(() => this.queueOperationStrict(registry, op)) ); } diff --git a/packages/gatekeeper/src/gatekeeper.ts b/packages/gatekeeper/src/gatekeeper.ts index 0f41c115b..591444c8c 100644 --- a/packages/gatekeeper/src/gatekeeper.ts +++ b/packages/gatekeeper/src/gatekeeper.ts @@ -739,21 +739,30 @@ export default class Gatekeeper implements GatekeeperInterface { return this.verifySignature(msgHash, signature.value, publicJwk); } - async queueOperation(registry: string, operation: Operation) { - // Don't distribute local DIDs + private publicationRegistries(registry: string): string[] { if (registry === 'local') { - return; + return []; } + return registry === 'hyperswarm' + ? ['hyperswarm'] + : ['hyperswarm', registry]; + } - // Always distribute on hyperswarm - await this.db.queueOperation('hyperswarm', operation); - - // Distribute on specified registry - if (registry !== 'hyperswarm') { - const queueSize = await this.db.queueOperation(registry, operation); + private async checkQueueCapacity(registry: string, queueSize?: number): Promise { + if (registry === 'local' || registry === 'hyperswarm') { + return; + } + const size = queueSize ?? (await this.db.getQueue(registry)).length; + if (size >= this.maxQueueSize) { + this.supportedRegistries = this.supportedRegistries.filter(reg => reg !== registry); + } + } - if (queueSize >= this.maxQueueSize) { - this.supportedRegistries = this.supportedRegistries.filter(reg => reg !== registry); + async queueOperation(registry: string, operation: Operation) { + for (const target of this.publicationRegistries(registry)) { + const queueSize = await this.db.queueOperation(target, operation); + if (target === registry) { + await this.checkQueueCapacity(registry, queueSize); } } } @@ -781,14 +790,14 @@ export default class Gatekeeper implements GatekeeperInterface { return did; } - await this.mutateDID(did, () => this.db.addEvent(did, { + await this.mutateDID(did, () => this.db.addEventAndQueue(did, { registry: 'local', time: operation.created!, ordinal: [0], operation, did - })); - await this.queueOperation(registry, operation); + }, this.publicationRegistries(registry))); + await this.checkQueueCapacity(registry); return did; }); } @@ -848,14 +857,14 @@ export default class Gatekeeper implements GatekeeperInterface { throw new InvalidOperationError(`registry ${registry} not supported`); } - await this.mutateDID(operation.did!, () => this.db.addEvent(operation.did!, { + await this.mutateDID(operation.did!, () => this.db.addEventAndQueue(operation.did!, { registry: 'local', time: operation.signature?.signed || '', ordinal: [0], operation, did: operation.did - })); - await this.queueOperation(registry, operation); + }, this.publicationRegistries(registry))); + await this.checkQueueCapacity(registry); return true; }); diff --git a/packages/gatekeeper/src/types.ts b/packages/gatekeeper/src/types.ts index faf3e8e08..8906af933 100644 --- a/packages/gatekeeper/src/types.ts +++ b/packages/gatekeeper/src/types.ts @@ -205,6 +205,7 @@ export interface GatekeeperDb { isReady(): Promise; resetDb(): Promise; addEvent(did: string, event: GatekeeperEvent): Promise; + addEventAndQueue(did: string, event: GatekeeperEvent, queueRegistries: string[]): Promise; getEvents(did: string): Promise; setEvents(did: string, events: GatekeeperEvent[], options?: SetEventsOptions): Promise; deleteEvents(did: string): Promise; diff --git a/tests/gatekeeper/crud.test.ts b/tests/gatekeeper/crud.test.ts index de10617ff..df4026c68 100644 --- a/tests/gatekeeper/crud.test.ts +++ b/tests/gatekeeper/crud.test.ts @@ -3,6 +3,7 @@ import Gatekeeper from '@mdip/gatekeeper'; import DbJsonMemory from '@mdip/gatekeeper/db/json-memory.ts'; import { ExpectedExceptionError } from '@mdip/common/errors'; import HeliaClient from '@mdip/ipfs/helia'; +import { jest } from '@jest/globals'; import TestHelper from './helper.ts'; const mockConsole = { @@ -40,6 +41,24 @@ describe('createDID', () => { expect(did.startsWith('did:test:')).toBe(true); }); + it('should not persist a DID when its publication queues cannot be committed', async () => { + const keypair = cipher.generateRandomJwk(); + const agentOp = await helper.createAgentOp(keypair, { registry: 'TFTC' }); + const did = await gatekeeper.generateDID(agentOp); + const addEvent = jest.spyOn(db, 'addEventAndQueue').mockRejectedValueOnce(new Error('queue failure')); + + await expect(gatekeeper.createDID(agentOp)).rejects.toThrow('queue failure'); + addEvent.mockRestore(); + + await expect(gatekeeper.exportDID(did)).resolves.toStrictEqual([]); + await expect(gatekeeper.getQueue('hyperswarm')).resolves.toStrictEqual([]); + await expect(gatekeeper.getQueue('TFTC')).resolves.toStrictEqual([]); + + await expect(gatekeeper.createDID(agentOp)).resolves.toBe(did); + await expect(gatekeeper.getQueue('hyperswarm')).resolves.toStrictEqual([agentOp]); + await expect(gatekeeper.getQueue('TFTC')).resolves.toStrictEqual([agentOp]); + }); + it('should create DID for local registry', async () => { const keypair = cipher.generateRandomJwk(); const agentOp = await helper.createAgentOp(keypair, { version: 1, registry: 'local' }); @@ -704,6 +723,30 @@ describe('updateDID', () => { expect(updatedDoc).toStrictEqual(doc); }); + it('should not persist an update when its publication queues cannot be committed', async () => { + const keypair = cipher.generateRandomJwk(); + const agentOp = await helper.createAgentOp(keypair, { registry: 'TFTC' }); + const did = await gatekeeper.createDID(agentOp); + await gatekeeper.clearQueue('hyperswarm', [agentOp]); + await gatekeeper.clearQueue('TFTC', [agentOp]); + const doc = await gatekeeper.resolveDID(did); + doc.didDocumentData = { mock: 1 }; + const updateOp = await helper.createUpdateOp(keypair, did, doc); + const addEvent = jest.spyOn(db, 'addEventAndQueue').mockRejectedValueOnce(new Error('queue failure')); + + await expect(gatekeeper.updateDID(updateOp)).rejects.toThrow('queue failure'); + addEvent.mockRestore(); + + await expect(gatekeeper.exportDID(did)).resolves.toHaveLength(1); + await expect(gatekeeper.getQueue('hyperswarm')).resolves.toStrictEqual([]); + await expect(gatekeeper.getQueue('TFTC')).resolves.toStrictEqual([]); + + await expect(gatekeeper.updateDID(updateOp)).resolves.toBe(true); + await expect(gatekeeper.exportDID(did)).resolves.toHaveLength(2); + await expect(gatekeeper.getQueue('hyperswarm')).resolves.toStrictEqual([updateOp]); + await expect(gatekeeper.getQueue('TFTC')).resolves.toStrictEqual([updateOp]); + }); + it('should reject an update whose previd is no longer current', async () => { const keypair = cipher.generateRandomJwk(); const did = await gatekeeper.createDID(await helper.createAgentOp(keypair)); diff --git a/tests/gatekeeper/index-export.test.ts b/tests/gatekeeper/index-export.test.ts index d62ee043b..dba4eba88 100644 --- a/tests/gatekeeper/index-export.test.ts +++ b/tests/gatekeeper/index-export.test.ts @@ -121,6 +121,7 @@ function createPostgresFixture(): AdapterFixture { ['z1', [eventA]], ['z2', [eventB]], ]); + const queues = new Map(); const changes = createIndexChanges(); let nextChangeSeq = changes.length; let indexEpoch = 'epoch-test'; @@ -208,6 +209,13 @@ function createPostgresFixture(): AdapterFixture { return { rows: [], rowCount: 1 }; } + if (text.includes('INSERT INTO gatekeeper_queue')) { + const registry = String(params[1]); + const operations = JSON.parse(String(params[2])) as GatekeeperEvent['operation'][]; + queues.set(registry, [...(queues.get(registry) ?? []), ...operations]); + return { rows: [], rowCount: 1 }; + } + if (text.includes('INSERT INTO gatekeeper_events')) { for (let i = 0; i < params.length; i += 4) { const id = String(params[i + 1]); @@ -236,6 +244,14 @@ function createPostgresFixture(): AdapterFixture { }; } + if (text.includes('FROM gatekeeper_queue')) { + const operations = queues.get(String(params[1])); + return { + rows: operations ? [{ ops: operations }] : [], + rowCount: operations ? 1 : 0, + }; + } + if (text.includes('FROM gatekeeper_index_changes')) { const afterSeq = Number(params[1]); const limit = Number(params[2]); @@ -438,10 +454,13 @@ function createRedisFixture(): AdapterFixture { const keys = values.slice(0, keyCount); const args = values.slice(keyCount); - if (keyCount === 4 && args.length === 3 && !/^\d+$/.test(args[0])) { + if (keyCount >= 4 && args.length >= 3 && !/^\d+$/.test(args[0]) && args[1]?.startsWith('{')) { const count = rpush(keys[0], args[0]); recordChange(JSON.parse(args[1]) as Omit); zadd(keys[3], 0, args[2]); + for (const queueKey of keys.slice(4)) { + rpush(queueKey, args[3]); + } return count; } @@ -537,6 +556,7 @@ function createMongoFixture(): AdapterFixture { ['z1', [eventA]], ['z2', [eventB]], ]); + const queues = new Map(); const changes = createIndexChanges(); let nextChangeSeq = changes.length; let indexEpoch = 'epoch-test'; @@ -604,6 +624,22 @@ function createMongoFixture(): AdapterFixture { }; } + if (name === 'queue') { + return { + findOne: async ({ id }: { id: string }) => { + const operations = queues.get(id); + return operations ? { id, ops: operations } : null; + }, + updateOne: async ( + { id }: { id: string }, + update: { $push: { ops: GatekeeperEvent['operation'] } } + ) => { + queues.set(id, [...(queues.get(id) ?? []), update.$push.ops]); + return { modifiedCount: 1, upsertedCount: 0 }; + }, + }; + } + if (name === 'counters') { return { findOneAndUpdate: async () => { @@ -1554,6 +1590,17 @@ describe.each(adapterFactories)('Gatekeeper DB index export adapter: %s', (_name await fixture.cleanup?.(); }); + it('commits an event with its publication queues', async () => { + const didC = 'did:test:z3'; + const eventC = createEvent(didC, '2026-01-01T00:00:03.000Z'); + + await fixture.db.addEventAndQueue(didC, eventC, ['hyperswarm', 'TFTC']); + + await expect(fixture.db.getEvents(didC)).resolves.toStrictEqual([eventC]); + await expect(fixture.db.getQueue('hyperswarm')).resolves.toStrictEqual([eventC.operation]); + await expect(fixture.db.getQueue('TFTC')).resolves.toStrictEqual([eventC.operation]); + }); + it('exports snapshots in stable DID order with cursor paging', async () => { const firstPage = await fixture.db.exportIndexSnapshot({ limit: 1 }); @@ -2370,11 +2417,15 @@ describe('Gatekeeper DB startup and guard behavior', () => { await db.resetDb(); await db.stop(); - expect(createIndex).toHaveBeenCalledTimes(6); + expect(createIndex).toHaveBeenCalledTimes(7); expect(createIndex).toHaveBeenCalledWith( { id: 1 }, { name: 'dids_id_unique', unique: true } ); + expect(createIndex).toHaveBeenCalledWith( + { id: 1 }, + { name: 'queue_id_unique', unique: true } + ); expect(createIndex).toHaveBeenCalledWith( { id: 1 }, { name: 'metadata_id_unique', unique: true } @@ -2398,6 +2449,73 @@ describe('Gatekeeper DB startup and guard behavior', () => { jest.resetModules(); }); + it('merges duplicate Mongo publication queues before creating the unique index', async () => { + const session = { + withTransaction: async (callback: () => Promise) => callback(), + endSession: async () => undefined, + }; + const deleteMany = jest.fn(async () => ({ deletedCount: 2 })); + const insertOne = jest.fn(async () => ({ acknowledged: true })); + let queueIndexAttempts = 0; + const appDb = { + collection: (name: string) => ({ + indexes: async () => [], + createIndex: async () => { + if (name === 'queue' && queueIndexAttempts++ === 0) { + throw { code: 11000 }; + } + }, + aggregate: () => ({ + toArray: async () => name === 'queue' + ? [{ + _id: 'TFTC', + operationBatches: [[eventA.operation], [eventB.operation]], + count: 2, + }] + : [], + }), + deleteMany, + insertOne, + find: () => ({ + sort: () => ({ + limit: () => ({ + next: async () => null, + }), + }), + }), + findOne: async () => null, + updateOne: async () => ({ modifiedCount: 0, upsertedCount: 1 }), + }), + }; + const client = { + connect: jest.fn(async () => undefined), + close: jest.fn(async () => undefined), + startSession: jest.fn(() => session), + db: jest.fn((name: string) => name === 'admin' + ? { command: jest.fn(async () => ({ setName: 'rs0' })) } + : appDb), + }; + jest.resetModules(); + jest.unstable_mockModule('mongodb', () => ({ + MongoClient: jest.fn(() => client), + })); + const { default: MockedDbMongo } = await import('@mdip/gatekeeper/db/mongo.ts'); + const db = new MockedDbMongo('mongo-queue-upgrade'); + + await db.start(); + await db.stop(); + + expect(deleteMany).toHaveBeenCalledWith({ id: 'TFTC' }, { session }); + expect(insertOne).toHaveBeenCalledWith({ + id: 'TFTC', + ops: [eventA.operation, eventB.operation], + }, { session }); + expect(queueIndexAttempts).toBe(2); + + jest.dontMock('mongodb'); + jest.resetModules(); + }); + it('initializes Mongo index sequence counter from existing index changes', async () => { const counterUpdates: unknown[][] = []; const client = { diff --git a/tests/gatekeeper/queue.test.ts b/tests/gatekeeper/queue.test.ts index 697e134f4..ef1480b9b 100644 --- a/tests/gatekeeper/queue.test.ts +++ b/tests/gatekeeper/queue.test.ts @@ -1,8 +1,12 @@ import CipherNode from '@mdip/cipher/node'; import Gatekeeper from '@mdip/gatekeeper'; import DbJsonMemory from '@mdip/gatekeeper/db/json-memory.ts'; +import DbSqlite from '@mdip/gatekeeper/db/sqlite.ts'; import { ExpectedExceptionError } from '@mdip/common/errors'; import HeliaClient from '@mdip/ipfs/helia'; +import fs from 'fs'; +import os from 'os'; +import path from 'path'; import TestHelper from './helper.ts'; const mockConsole = { @@ -30,8 +34,57 @@ beforeEach(async () => { await gatekeeper.resetDb(); // Reset database for each test to ensure isolation }); +describe('atomic event publication', () => { + it('should roll back an event when a publication queue write fails', async () => { + const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), 'gatekeeper-publication-')); + const sqliteDb = new DbSqlite('atomic', tempDir); + await sqliteDb.start(); + + try { + const operation = await helper.createAgentOp(cipher.generateRandomJwk(), { registry: 'TFTC' }); + const did = await gatekeeper.generateDID(operation); + const event = { + registry: 'local', + time: operation.created!, + ordinal: [0], + operation, + did, + }; + await (sqliteDb as any).db.exec(` + CREATE TRIGGER fail_publication_queue + BEFORE INSERT ON queue + BEGIN + SELECT RAISE(FAIL, 'queue insert failed'); + END; + `); + + await expect(sqliteDb.addEventAndQueue( + did, + event, + ['hyperswarm', 'TFTC'] + )).rejects.toThrow('queue insert failed'); + await expect(sqliteDb.getEvents(did)).resolves.toStrictEqual([]); + await expect(sqliteDb.getQueue('hyperswarm')).resolves.toStrictEqual([]); + await expect(sqliteDb.getQueue('TFTC')).resolves.toStrictEqual([]); + } + finally { + await sqliteDb.stop(); + fs.rmSync(tempDir, { recursive: true, force: true }); + } + }); +}); + describe('getQueue', () => { + it('should queue an operation for each publication registry', async () => { + const operation = await helper.createAgentOp(cipher.generateRandomJwk(), { registry: 'TFTC' }); + + await gatekeeper.queueOperation('TFTC', operation); + + await expect(gatekeeper.getQueue('hyperswarm')).resolves.toStrictEqual([operation]); + await expect(gatekeeper.getQueue('TFTC')).resolves.toStrictEqual([operation]); + }); + it('should return empty list when no events in queue', async () => { const registry = 'TFTC'; diff --git a/tests/gatekeeper/verify.test.ts b/tests/gatekeeper/verify.test.ts index fb9e13423..52c5f5638 100644 --- a/tests/gatekeeper/verify.test.ts +++ b/tests/gatekeeper/verify.test.ts @@ -794,9 +794,9 @@ describe('verifyDb', () => { const agentDoc = await gatekeeper.resolveDID(agentDID); agentDoc.mdip!.validUntil = new Date(Date.now() - 1_000).toISOString(); const updateOp = await helper.createUpdateOp(keypair, agentDID, agentDoc); - const originalAddEvent = db.addEvent.bind(db); - const addEvent = jest.spyOn(db, 'addEvent').mockImplementationOnce(async (did, event) => { - await originalAddEvent(did, event); + const originalAddEventAndQueue = db.addEventAndQueue.bind(db); + const addEvent = jest.spyOn(db, 'addEventAndQueue').mockImplementationOnce(async (did, event, registries) => { + await originalAddEventAndQueue(did, event, registries); throw new Error('connection lost after commit'); }); @@ -823,15 +823,15 @@ describe('verifyDb', () => { const agentDoc = await gatekeeper.resolveDID(agentDID); agentDoc.mdip!.validUntil = new Date(Date.now() - 1_000).toISOString(); const updateOp = await helper.createUpdateOp(keypair, agentDID, agentDoc); - const originalAddEvent = db.addEvent.bind(db); + const originalAddEventAndQueue = db.addEventAndQueue.bind(db); let signalWrite: () => void = () => { }; let releaseWrite: () => void = () => { }; const writeStarted = new Promise(resolve => (signalWrite = resolve)); const writeRelease = new Promise(resolve => (releaseWrite = resolve)); - const addEvent = jest.spyOn(db, 'addEvent').mockImplementationOnce(async (did, event) => { + const addEvent = jest.spyOn(db, 'addEventAndQueue').mockImplementationOnce(async (did, event, registries) => { signalWrite(); await writeRelease; - return originalAddEvent(did, event); + return originalAddEventAndQueue(did, event, registries); }); const update = gatekeeper.updateDID(updateOp);