Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions packages/gatekeeper/src/db/abstract-json.ts
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,10 @@ export abstract class AbstractJson implements GatekeeperDb {
}

async addEvent(did: string, event: GatekeeperEvent): Promise<void> {
return this.addEventAndQueue(did, event, []);
}

async addEventAndQueue(did: string, event: GatekeeperEvent, queueRegistries: string[]): Promise<void> {
const suffix = this.splitSuffix(did);
return this.runExclusive(async () => {
const db = this.loadDb();
Expand All @@ -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);
});
}
Expand Down
72 changes: 72 additions & 0 deletions packages/gatekeeper/src/db/mongo.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,12 @@ interface QueueDoc {
ops: Operation[]
}

interface DuplicateQueueDoc {
_id: string
operationBatches: Operation[][]
count: number
}

interface CounterDoc {
id: string;
value: number;
Expand All @@ -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' });
Expand Down Expand Up @@ -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 });
Expand Down Expand Up @@ -204,6 +218,53 @@ export default class DbMongo implements GatekeeperDb {
await collection.createIndex(key, options);
}

private async ensureQueueIndex(): Promise<void> {
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<void> {
if (!this.db) {
throw new Error(MONGO_NOT_STARTED_ERROR);
}

const collection = this.db.collection<QueueDoc>('queue');
const duplicates = await collection.aggregate<DuplicateQueueDoc>([
{ $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;
}
Expand Down Expand Up @@ -325,6 +386,10 @@ export default class DbMongo implements GatekeeperDb {
}

async addEvent(did: string, event: GatekeeperEvent): Promise<number> {
return this.addEventAndQueue(did, event, []);
}

async addEventAndQueue(did: string, event: GatekeeperEvent, queueRegistries: string[]): Promise<number> {
if (!this.db) {
throw new Error(MONGO_NOT_STARTED_ERROR)
}
Expand All @@ -349,6 +414,13 @@ export default class DbMongo implements GatekeeperDb {
did,
event,
}, session);
for (const registry of queueRegistries) {
await this.db!.collection<QueueDoc>('queue').updateOne(
{ id: registry },
{ $push: { ops: event.operation } },
{ upsert: true, session }
);
}
return count
});
}
Expand Down
36 changes: 33 additions & 3 deletions packages/gatekeeper/src/db/postgres.ts
Original file line number Diff line number Diff line change
Expand Up @@ -422,10 +422,14 @@ export default class DbPostgres implements GatekeeperDb {
});
}

async addEvent(did: string, event: GatekeeperEvent): Promise<number> {
const id = this.splitSuffix(did);
private async addEventStrict(
executor: Pool | PoolClient,
did: string,
id: string,
event: GatekeeperEvent
): Promise<number> {
const serializedEvent = JSON.stringify(event);
const result = await this.getPool().query<LengthRow>(
const result = await executor.query<LengthRow>(
`WITH inserted_event AS (
INSERT INTO gatekeeper_events (namespace, id, seq, event)
SELECT $1, $2, COALESCE(MAX(seq), -1) + 1, $3::jsonb
Expand All @@ -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<number> {
const id = this.splitSuffix(did);
return this.addEventStrict(this.getPool(), did, id, event);
}

async addEventAndQueue(did: string, event: GatekeeperEvent, queueRegistries: string[]): Promise<number> {
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<number> {
const id = this.splitSuffix(did);

Expand Down
14 changes: 11 additions & 3 deletions packages/gatekeeper/src/db/redis.ts
Original file line number Diff line number Diff line change
Expand Up @@ -129,6 +129,10 @@ export default class DbRedis implements GatekeeperDb {
}

async addEvent(did: string, event: GatekeeperEvent): Promise<number> {
return this.addEventAndQueue(did, event, []);
}

async addEventAndQueue(did: string, event: GatekeeperEvent, queueRegistries: string[]): Promise<number> {
if (!this.redis) {
throw new Error(REDIS_NOT_STARTED_ERROR)
}
Expand All @@ -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);
Expand Down
35 changes: 22 additions & 13 deletions packages/gatekeeper/src/db/sqlite.ts
Original file line number Diff line number Diff line change
Expand Up @@ -192,6 +192,10 @@ export default class DbSqlite implements GatekeeperDb {
}

async addEvent(did: string, event: GatekeeperEvent): Promise<number> {
return this.addEventAndQueue(did, event, []);
}

async addEventAndQueue(did: string, event: GatekeeperEvent, queueRegistries: string[]): Promise<number> {
if (!did) {
throw new InvalidDIDError();
}
Expand All @@ -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;
})
);
Expand Down Expand Up @@ -323,22 +330,24 @@ export default class DbSqlite implements GatekeeperDb {
);
}

async queueOperation(registry: string, op: Operation): Promise<number> {
private async queueOperationStrict(registry: string, op: Operation): Promise<number> {
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<number> {
return this.runExclusive(() =>
this.withTx(() => this.queueOperationStrict(registry, op))
);
}

Expand Down
43 changes: 26 additions & 17 deletions packages/gatekeeper/src/gatekeeper.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void> {
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);
}
}
}
Expand Down Expand Up @@ -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;
});
}
Expand Down Expand Up @@ -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;
});
Expand Down
1 change: 1 addition & 0 deletions packages/gatekeeper/src/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -205,6 +205,7 @@ export interface GatekeeperDb {
isReady(): Promise<boolean>;
resetDb(): Promise<void | number | JsonDbFile>;
addEvent(did: string, event: GatekeeperEvent): Promise<void | number>;
addEventAndQueue(did: string, event: GatekeeperEvent, queueRegistries: string[]): Promise<void | number>;
getEvents(did: string): Promise<GatekeeperEvent[]>;
setEvents(did: string, events: GatekeeperEvent[], options?: SetEventsOptions): Promise<number | void>;
deleteEvents(did: string): Promise<void | number>;
Expand Down
Loading
Loading