Implements #1707
This splits indexing into two phases: * preindex (where only pages, space lua, space styles are indexed) * index (where everything else is indexed) The goal is to get to more-or-less operational state quicker this way, because probably most functionality can already be enabled after the pre-index phase.
This commit is contained in:
@@ -109,6 +109,11 @@ export class EncryptedKvPrimitives implements KvPrimitives {
|
||||
}
|
||||
}
|
||||
|
||||
async countQuery({ prefix }: KvQueryOptions): Promise<number> {
|
||||
const encryptedPrefix = prefix ? await this.encryptKey(prefix) : undefined;
|
||||
return this.wrapped.countQuery({ prefix: encryptedPrefix });
|
||||
}
|
||||
|
||||
close() {
|
||||
this.wrapped.close();
|
||||
}
|
||||
|
||||
@@ -79,6 +79,15 @@ export class IndexedDBKvPrimitives implements KvPrimitives {
|
||||
}
|
||||
}
|
||||
|
||||
countQuery({ prefix }: KvQueryOptions): Promise<number> {
|
||||
const tx = this.db.transaction(objectStoreName, "readonly");
|
||||
prefix = prefix || [];
|
||||
return tx.store.count(IDBKeyRange.bound(
|
||||
this.buildKey([...prefix, ""]),
|
||||
this.buildKey([...prefix, "\uffff"]),
|
||||
));
|
||||
}
|
||||
|
||||
close() {
|
||||
this.db.close();
|
||||
}
|
||||
|
||||
@@ -13,6 +13,11 @@ export interface KvPrimitives {
|
||||
|
||||
query<T = any>(options: KvQueryOptions): AsyncIterableIterator<KV<T>>;
|
||||
|
||||
/**
|
||||
* A more efficient way to do counts
|
||||
*/
|
||||
countQuery(options: KvQueryOptions): Promise<number>;
|
||||
|
||||
// Completely clear all data from this datastore
|
||||
clear(): Promise<void>;
|
||||
|
||||
|
||||
@@ -144,6 +144,14 @@ export class MemoryKvPrimitives implements KvPrimitives {
|
||||
}
|
||||
}
|
||||
|
||||
countQuery({ prefix }: KvQueryOptions): Promise<number> {
|
||||
const prefixStr = prefix ? prefix.join(memoryKeySeparator) : undefined;
|
||||
const keys = [...this.store.keys()];
|
||||
return Promise.resolve(
|
||||
keys.filter((key) => !prefixStr || key.startsWith(prefixStr)).length,
|
||||
);
|
||||
}
|
||||
|
||||
async close(): Promise<void> {
|
||||
// Force immediate persistence when closing
|
||||
if (this.filePath) {
|
||||
|
||||
@@ -5,14 +5,21 @@ import { DataStore } from "./datastore.ts";
|
||||
import { FakeTime } from "@std/testing/time";
|
||||
|
||||
import type { MQMessage } from "../../plug-api/types/datastore.ts";
|
||||
import { EventHook } from "../plugos/hooks/event.ts";
|
||||
import { System } from "../plugos/system.ts";
|
||||
import type { EventHookT } from "@silverbulletmd/silverbullet/type/manifest";
|
||||
|
||||
Deno.test("DataStore MQ", async () => {
|
||||
const time = new FakeTime();
|
||||
const db = new MemoryKvPrimitives(); // In-memory only, no persistence
|
||||
const eventHook = new EventHook();
|
||||
const system = new System<EventHookT>();
|
||||
system.addHook(eventHook);
|
||||
|
||||
try {
|
||||
const mq = new DataStoreMQ(
|
||||
new DataStore(db),
|
||||
eventHook,
|
||||
);
|
||||
|
||||
let messages: MQMessage[];
|
||||
@@ -32,8 +39,8 @@ Deno.test("DataStore MQ", async () => {
|
||||
await time.tickAsync(20);
|
||||
await mq.requeueTimeouts(10);
|
||||
messages = await mq.poll("test", 10);
|
||||
const stats = await mq.getAllQueueStats();
|
||||
assertEquals(stats["test"].processing, 1);
|
||||
const stats = await mq.getQueueStats();
|
||||
assertEquals(stats.processing, 1);
|
||||
assertEquals(messages.length, 1);
|
||||
assertEquals(messages[0].retries, 1);
|
||||
|
||||
@@ -77,10 +84,14 @@ Deno.test("DataStore MQ", async () => {
|
||||
Deno.test("DataStore MQ - Scale test with multiple subscribers", async () => {
|
||||
const time = new FakeTime();
|
||||
const db = new MemoryKvPrimitives();
|
||||
const eventHook = new EventHook();
|
||||
const system = new System<EventHookT>();
|
||||
system.addHook(eventHook);
|
||||
|
||||
try {
|
||||
const mq = new DataStoreMQ(
|
||||
new DataStore(db),
|
||||
eventHook,
|
||||
);
|
||||
|
||||
const queueName = "scale-test";
|
||||
|
||||
+16
-57
@@ -9,6 +9,7 @@ import type {
|
||||
MQSubscribeOptions,
|
||||
} from "../../plug-api/types/datastore.ts";
|
||||
import { race, sleep } from "@silverbulletmd/silverbullet/lib/async";
|
||||
import type { EventHook } from "../plugos/hooks/event.ts";
|
||||
|
||||
export type ProcessingMessage = MQMessage & {
|
||||
ts: number;
|
||||
@@ -49,6 +50,10 @@ export class QueueWorker {
|
||||
await this.callback(messages);
|
||||
} else {
|
||||
// No messages, wait to be woken up or a timeout
|
||||
this.mq.eventHook.dispatchEvent(
|
||||
`mq:emptyQueue:${this.queue}`,
|
||||
this.queue,
|
||||
);
|
||||
try {
|
||||
await race([
|
||||
// Wait to be woken up explicitly
|
||||
@@ -98,6 +103,7 @@ export class DataStoreMQ {
|
||||
|
||||
constructor(
|
||||
private ds: DataStore,
|
||||
public eventHook: EventHook,
|
||||
) {
|
||||
}
|
||||
|
||||
@@ -350,12 +356,16 @@ export class DataStoreMQ {
|
||||
await this.ds.batchDelete(ids);
|
||||
}
|
||||
|
||||
async getQueueStats(queue: string): Promise<MQStats> {
|
||||
const queued =
|
||||
(await (this.ds.luaQuery([...queuedPrefix, queue], {}))).length;
|
||||
const processing =
|
||||
(await (this.ds.luaQuery([...processingPrefix, queue], {}))).length;
|
||||
const dlq = (await (this.ds.luaQuery([...dlqPrefix, queue], {}))).length;
|
||||
async getQueueStats(queue?: string): Promise<MQStats> {
|
||||
const queued = await this.ds.kv.countQuery({
|
||||
prefix: queue ? [...queuedPrefix, queue] : queuedPrefix,
|
||||
});
|
||||
const processing = await this.ds.kv.countQuery({
|
||||
prefix: queue ? [...processingPrefix, queue] : processingPrefix,
|
||||
});
|
||||
const dlq = await this.ds.kv.countQuery({
|
||||
prefix: queue ? [...dlqPrefix, queue] : dlqPrefix,
|
||||
});
|
||||
return {
|
||||
queued,
|
||||
processing,
|
||||
@@ -372,55 +382,4 @@ export class DataStoreMQ {
|
||||
await sleep(200);
|
||||
}
|
||||
}
|
||||
|
||||
async getAllQueueStats(): Promise<Record<string, MQStats>> {
|
||||
const allStatus: Record<string, MQStats> = {};
|
||||
for (
|
||||
const message of await this.ds.luaQuery<MQMessage>(
|
||||
queuedPrefix,
|
||||
{},
|
||||
)
|
||||
) {
|
||||
if (!allStatus[message.queue]) {
|
||||
allStatus[message.queue] = {
|
||||
queued: 0,
|
||||
processing: 0,
|
||||
dlq: 0,
|
||||
};
|
||||
}
|
||||
allStatus[message.queue].queued++;
|
||||
}
|
||||
for (
|
||||
const message of await this.ds.luaQuery<MQMessage>(
|
||||
processingPrefix,
|
||||
{},
|
||||
)
|
||||
) {
|
||||
if (!allStatus[message.queue]) {
|
||||
allStatus[message.queue] = {
|
||||
queued: 0,
|
||||
processing: 0,
|
||||
dlq: 0,
|
||||
};
|
||||
}
|
||||
allStatus[message.queue].processing++;
|
||||
}
|
||||
for (
|
||||
const message of await this.ds.luaQuery<MQMessage>(
|
||||
dlqPrefix,
|
||||
{},
|
||||
)
|
||||
) {
|
||||
if (!allStatus[message.queue]) {
|
||||
allStatus[message.queue] = {
|
||||
queued: 0,
|
||||
processing: 0,
|
||||
dlq: 0,
|
||||
};
|
||||
}
|
||||
allStatus[message.queue].dlq++;
|
||||
}
|
||||
|
||||
return allStatus;
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user