Code restructuring: eliminated the top-level lib/ directory

This commit is contained in:
Zef Hemel
2025-09-23 17:00:34 +02:00
parent 4fbdb4c9dd
commit f4336bf09f
175 changed files with 353 additions and 327 deletions
+45
View File
@@ -0,0 +1,45 @@
import "fake-indexeddb/auto";
import { IndexedDBKvPrimitives } from "./indexeddb_kv_primitives.ts";
import { MemoryKvPrimitives } from "./memory_kv_primitives.ts";
import type { KvPrimitives } from "./kv_primitives.ts";
import { assertEquals } from "@std/assert";
import { PrefixedKvPrimitives } from "./prefixed_kv_primitives.ts";
import { DataStore } from "./datastore.ts";
import { LuaEnv, LuaStackFrame } from "../space_lua/runtime.ts";
import { parseExpressionString } from "../space_lua/parse.ts";
async function test(db: KvPrimitives) {
const datastore = new DataStore(new PrefixedKvPrimitives(db, ["ds"]));
await datastore.set(["user", "peter"], { name: "Peter" });
await datastore.set(["user", "hank"], { name: "Hank" });
const env = new LuaEnv();
const sf = LuaStackFrame.lostFrame;
// Basic test, fancier tests are done in common/space_lua/query_collection.test.ts
const results = await datastore.luaQuery(
["user"],
{
objectVariable: "user",
where: parseExpressionString("user.name == 'Peter'"),
},
env,
sf,
);
assertEquals(results, [{ name: "Peter" }]);
}
Deno.test("Test Memory KV DataStore", async () => {
const db = new MemoryKvPrimitives(); // In-memory only, no persistence
await test(db);
await db.close();
});
Deno.test("Test IndexDB DataStore", {
sanitizeResources: false,
sanitizeOps: false,
}, async () => {
const db = new IndexedDBKvPrimitives("test");
await db.init();
await test(db);
db.close();
});
+85
View File
@@ -0,0 +1,85 @@
import {
type LuaCollectionQuery,
queryLua,
} from "../space_lua/query_collection.ts";
import { LuaEnv, LuaStackFrame } from "../space_lua/runtime.ts";
import type { KvPrimitives, KvQueryOptions } from "./kv_primitives.ts";
import type { KV, KvKey } from "../../plug-api/types/datastore.ts";
/**
* This is the data store class you'll actually want to use, wrapping the primitives
* in a more user-friendly way
*/
export class DataStore {
constructor(
readonly kv: KvPrimitives,
) {
}
async get<T = any>(key: KvKey): Promise<T | null> {
return (await this.batchGet([key]))[0];
}
batchGet<T = any>(keys: KvKey[]): Promise<(T | null)[]> {
if (keys.length === 0) {
return Promise.resolve([]);
}
return this.kv.batchGet(keys);
}
set(key: KvKey, value: any): Promise<void> {
return this.batchSet([{ key, value }]);
}
batchSet<T = any>(entries: KV<T>[]): Promise<void> {
if (entries.length === 0) {
return Promise.resolve();
}
const allKeyStrings = new Set<string>();
const uniqueEntries: KV[] = [];
for (const { key, value } of entries) {
const keyString = JSON.stringify(key);
if (allKeyStrings.has(keyString)) {
console.warn(`Duplicate key ${keyString} in batchSet, skipping`);
} else {
allKeyStrings.add(keyString);
uniqueEntries.push({ key, value });
}
}
return this.kv.batchSet(uniqueEntries);
}
delete(key: KvKey): Promise<void> {
return this.batchDelete([key]);
}
batchDelete(keys: KvKey[]): Promise<void> {
if (keys.length === 0) {
return Promise.resolve();
}
return this.kv.batchDelete(keys);
}
async batchDeletePrefix(prefix: KvKey): Promise<void> {
const keys: KvKey[] = [];
for await (const { key } of this.kv.query({ prefix })) {
keys.push(key);
}
return this.batchDelete(keys);
}
query(options: KvQueryOptions): AsyncIterableIterator<KV> {
return this.kv.query(options);
}
luaQuery<T = any>(
prefix: KvKey,
query: LuaCollectionQuery,
env: LuaEnv = new LuaEnv(),
sf: LuaStackFrame = LuaStackFrame.lostFrame,
enricher?: (key: KvKey, item: any) => any,
): Promise<T[]> {
return queryLua(this.kv, prefix, query, env, sf, enricher);
}
}
@@ -0,0 +1,13 @@
import "fake-indexeddb/auto";
import { IndexedDBKvPrimitives } from "./indexeddb_kv_primitives.ts";
import { allTests } from "./kv_primitives.test.ts";
Deno.test("Test IDB key primitives", {
sanitizeResources: false,
sanitizeOps: false,
}, async () => {
const db = new IndexedDBKvPrimitives("test");
await db.init();
await allTests(db);
db.close();
});
+96
View File
@@ -0,0 +1,96 @@
import type { KvPrimitives, KvQueryOptions } from "./kv_primitives.ts";
import { type IDBPDatabase, openDB } from "idb";
import type { KV, KvKey } from "../../plug-api/types/datastore.ts";
// Separator character to use for key serialization
const sep = "\0";
const objectStoreName = "data";
export class IndexedDBKvPrimitives implements KvPrimitives {
db!: IDBPDatabase<any>;
constructor(
private dbName: string,
) {
}
async init() {
this.db = await openDB(this.dbName, 1, {
upgrade: (db) => {
db.createObjectStore(objectStoreName);
},
});
}
async clear(): Promise<void> {
const objectStoreNames = this.db.objectStoreNames;
// Create a transaction that includes all object stores
const tx = this.db.transaction(objectStoreNames, "readwrite");
// Clear each object store in parallel
const clearPromises = Array.from(objectStoreNames).map((storeName) =>
tx.objectStore(storeName).clear()
);
// Wait for all clears to complete
await Promise.all(clearPromises);
// Complete the transaction
await tx.done;
}
batchGet(keys: KvKey[]): Promise<any[]> {
const tx = this.db.transaction(objectStoreName, "readonly");
return Promise.all(keys.map((key) => tx.store.get(this.buildKey(key))));
}
async batchSet(entries: KV[]): Promise<void> {
const tx = this.db.transaction(objectStoreName, "readwrite");
await Promise.all([
...entries.map(({ key, value }) =>
tx.store.put(value, this.buildKey(key))
),
tx.done,
]);
}
async batchDelete(keys: KvKey[]): Promise<void> {
const tx = this.db.transaction(objectStoreName, "readwrite");
await Promise.all([
...keys.map((key) => tx.store.delete(this.buildKey(key))),
tx.done,
]);
}
async *query({ prefix }: KvQueryOptions): AsyncIterableIterator<KV> {
const tx = this.db.transaction(objectStoreName, "readonly");
prefix = prefix || [];
for await (
const entry of tx.store.iterate(IDBKeyRange.bound(
this.buildKey([...prefix, ""]),
this.buildKey([...prefix, "\uffff"]),
))
) {
yield { key: this.extractKey(entry.key), value: entry.value };
}
}
close() {
this.db.close();
}
private buildKey(key: KvKey): string {
for (const k of key) {
if (k.includes(sep)) {
throw new Error(`Key cannot contain ${sep}`);
}
}
return key.join(sep);
}
private extractKey(key: string): KvKey {
return key.split(sep);
}
}
+70
View File
@@ -0,0 +1,70 @@
import type { KvPrimitives } from "./kv_primitives.ts";
import { assertEquals } from "@std/assert";
import type { KV } from "../../plug-api/types/datastore.ts";
export async function allTests(db: KvPrimitives) {
await db.batchSet([
{ key: ["kv", "test2"], value: "Hello2" },
{ key: ["kv", "test1"], value: "Hello1" },
{ key: ["other", "random"], value: "Hello3" },
]);
const result = await db.batchGet([["kv", "test1"], ["kv", "test2"], [
"kv",
"test3",
]]);
assertEquals(result.length, 3);
assertEquals(result[0], "Hello1");
assertEquals(result[1], "Hello2");
assertEquals(result[2], undefined);
let counter = 0;
// Query all
for await (const _entry of db.query({})) {
counter++;
}
assertEquals(counter, 3);
counter = 0;
// Query prefix
for await (const _entry of db.query({ prefix: ["kv"] })) {
counter++;
console.log(_entry);
}
assertEquals(counter, 2);
// Delete a few keys
await db.batchDelete([["kv", "test1"], ["other", "random"]]);
const result2 = await db.batchGet([["kv", "test1"], ["kv", "test2"], [
"other",
"random",
]]);
assertEquals(result2.length, 3);
assertEquals(result2[0], undefined);
assertEquals(result2[1], "Hello2");
assertEquals(result2[2], undefined);
// Update a key
await db.batchSet([{ key: ["kv", "test2"], value: "Hello2.1" }]);
const [val] = await db.batchGet([["kv", "test2"]]);
assertEquals(val, "Hello2.1");
// Set a large batch
const largeBatch: KV[] = [];
for (let i = 0; i < 50; i++) {
largeBatch.push({ key: ["test", "test" + i], value: "Hello" });
}
await db.batchSet(largeBatch);
const largeBatchResult: KV[] = [];
for await (const entry of db.query({ prefix: ["test"] })) {
largeBatchResult.push(entry);
}
assertEquals(largeBatchResult.length, 50);
// Delete the large batch
await db.batchDelete(largeBatch.map((e) => e.key));
// Make sure they're gone
for await (const _entry of db.query({ prefix: ["test"] })) {
throw new Error("This should not happen");
}
}
+20
View File
@@ -0,0 +1,20 @@
import type { KV, KvKey } from "../../plug-api/types/datastore.ts";
export type KvQueryOptions = {
prefix?: KvKey;
};
export interface KvPrimitives {
batchGet(keys: KvKey[]): Promise<(any | undefined)[]>;
batchSet(entries: KV[]): Promise<void>;
batchDelete(keys: KvKey[]): Promise<void>;
query(options: KvQueryOptions): AsyncIterableIterator<KV>;
// Completely clear all data from this datastore
clear(): Promise<void>;
close(): void;
}
+198
View File
@@ -0,0 +1,198 @@
import { assertEquals } from "@std/assert";
import { MemoryKvPrimitives } from "./memory_kv_primitives.ts";
import { allTests } from "./kv_primitives.test.ts";
import type { KV } from "../../plug-api/types/datastore.ts";
Deno.test("MemoryKvPrimitives loads from non-existent file without error", async () => {
const tempPath = await Deno.makeTempFile() + "_nonexistent";
// Disable throttling for tests
const store = new MemoryKvPrimitives(tempPath, { throttleMs: 0 });
await store.init();
// Should create an empty store
const result = await store.batchGet([["test"]]);
assertEquals(result, [undefined]);
});
Deno.test("MemoryKvPrimitives passes all KvPrimitives tests", async () => {
const tempPath = await Deno.makeTempFile();
try {
const store = new MemoryKvPrimitives(tempPath, { throttleMs: 0 });
await store.init();
await allTests(store);
await store.close();
} finally {
// Clean up
try {
await Deno.remove(tempPath);
} catch (_) {
// Ignore errors during cleanup
}
}
});
Deno.test("MemoryKvPrimitives persists and loads data", async () => {
const tempPath = await Deno.makeTempFile();
try {
// Create and populate first instance
const store1 = new MemoryKvPrimitives(tempPath, { throttleMs: 0 });
await store1.init();
await store1.batchSet([
{ key: ["test", "key1"], value: "value1" },
{ key: ["test", "key2"], value: "value2" },
]);
// Force persistence
await store1.close();
// Create second instance that loads from the same file
const store2 = new MemoryKvPrimitives(tempPath, { throttleMs: 0 });
await store2.init();
// Check if data was loaded correctly
const results = await store2.batchGet([["test", "key1"], ["test", "key2"]]);
assertEquals(results, ["value1", "value2"]);
} finally {
// Clean up
try {
await Deno.remove(tempPath);
} catch (_) {
// Ignore errors during cleanup
}
}
});
Deno.test("MemoryKvPrimitives mutations trigger persistence", async () => {
const tempPath = await Deno.makeTempFile();
try {
// Create and populate first instance
const store1 = new MemoryKvPrimitives(tempPath, { throttleMs: 0 });
await store1.init();
await store1.batchSet([{ key: ["test", "key"], value: "value" }]);
// No need to wait for throttled persistence since we disabled it
// Create second instance without closing the first one
const store2 = new MemoryKvPrimitives(tempPath, { throttleMs: 0 });
await store2.init();
// Check if data was persisted
const results = await store2.batchGet([["test", "key"]]);
assertEquals(results, ["value"]);
} finally {
// Clean up
try {
await Deno.remove(tempPath);
} catch (_) {
// Ignore errors during cleanup
}
}
});
Deno.test("MemoryKvPrimitives persists delete operations", async () => {
const tempPath = await Deno.makeTempFile();
try {
// Create and populate store
const store1 = new MemoryKvPrimitives(tempPath, { throttleMs: 0 });
await store1.init();
await store1.batchSet([
{ key: ["test", "key1"], value: "value1" },
{ key: ["test", "key2"], value: "value2" },
]);
// Delete one key
await store1.batchDelete([["test", "key1"]]);
// No need to wait for throttled persistence since we disabled it
// Create second instance
const store2 = new MemoryKvPrimitives(tempPath, { throttleMs: 0 });
await store2.init();
// Check if delete was persisted
const results = await store2.batchGet([["test", "key1"], ["test", "key2"]]);
assertEquals(results, [undefined, "value2"]);
} finally {
// Clean up
try {
await Deno.remove(tempPath);
} catch (_) {
// Ignore errors during cleanup
}
}
});
Deno.test("MemoryKvPrimitives.fromFile creates and initializes store", async () => {
const tempPath = await Deno.makeTempFile();
try {
// Create JSON file with initial data
const initialData = {
"test\0key": "value",
};
await Deno.writeTextFile(tempPath, JSON.stringify(initialData));
// Use factory method with throttling disabled
const store = await MemoryKvPrimitives.fromFile(tempPath, {
throttleMs: 0,
});
// Check if data was loaded
const result = await store.batchGet([["test", "key"]]);
assertEquals(result, ["value"]);
} finally {
// Clean up
try {
await Deno.remove(tempPath);
} catch (_) {
// Ignore errors during cleanup
}
}
});
Deno.test("MemoryKvPrimitives query works with persisted data", async () => {
const tempPath = await Deno.makeTempFile();
try {
// Create and populate store
const store1 = new MemoryKvPrimitives(tempPath, { throttleMs: 0 });
await store1.init();
await store1.batchSet([
{ key: ["test", "key1"], value: "value1" },
{ key: ["test", "key2"], value: "value2" },
{ key: ["other", "key"], value: "value3" },
]);
// Force persistence
await store1.close();
// Create second instance
const store2 = await MemoryKvPrimitives.fromFile(tempPath, {
throttleMs: 0,
});
// Test query with prefix
const results: KV[] = [];
for await (const item of store2.query({ prefix: ["test"] })) {
results.push(item);
}
assertEquals(results.length, 2);
assertEquals(results[0].key, ["test", "key1"]);
assertEquals(results[0].value, "value1");
assertEquals(results[1].key, ["test", "key2"]);
assertEquals(results[1].value, "value2");
} finally {
// Clean up
try {
await Deno.remove(tempPath);
} catch (_) {
// Ignore errors during cleanup
}
}
});
+170
View File
@@ -0,0 +1,170 @@
import type { KvPrimitives, KvQueryOptions } from "./kv_primitives.ts";
import { throttle } from "@silverbulletmd/silverbullet/lib/async";
import type { KV, KvKey } from "@silverbulletmd/silverbullet/type/datastore";
const memoryKeySeparator = "\0";
export class MemoryKvPrimitives implements KvPrimitives {
protected store = new Map<string, any>();
private throttledPersist?: () => void;
constructor(
protected filePath?: string,
options: { throttleMs?: number } = {},
) {
// Set up throttled persistence if throttleMs is provided or default to 1000ms
if (this.filePath) {
const throttleMs = options.throttleMs !== undefined
? options.throttleMs
: 1000;
// If throttleMs is 0, persistence will happen immediately without throttling
if (throttleMs > 0) {
this.throttledPersist = throttle(() => {
this.persistToDisk().catch((err) =>
console.error(`Error persisting to disk: ${err}`)
);
}, throttleMs);
}
}
}
static fromJSON(json: Record<string, any>): MemoryKvPrimitives {
const result = new MemoryKvPrimitives();
for (const key of Object.keys(json)) {
result.store.set(key, json[key]);
}
return result;
}
/**
* Create a new MemoryKvPrimitives instance from a file and initialize it
*/
static async fromFile(
filePath: string,
options: { throttleMs?: number } = {},
): Promise<MemoryKvPrimitives> {
const instance = new MemoryKvPrimitives(filePath, options);
await instance.init();
return instance;
}
clear(): Promise<void> {
this.store.clear();
return Promise.resolve();
}
/**
* Initialize the store by loading data from disk if a file path was provided
*/
async init(): Promise<void> {
if (!this.filePath) return;
try {
const text = await Deno.readTextFile(this.filePath);
// Handle empty files gracefully to prevent "SyntaxError: Unexpected end of JSON input"
if (text.trim() === "") {
return;
}
const jsonData = JSON.parse(text);
for (const key of Object.keys(jsonData)) {
this.store.set(key, jsonData[key]);
}
} catch (error) {
// Handle specific errors more gracefully
if (error instanceof Deno.errors.NotFound) {
// File doesn't exist yet, nothing to load
return;
}
// Other errors (like invalid JSON) should be logged
console.warn(`Failed to load KV store from ${this.filePath}:`, error);
}
}
batchGet(keys: KvKey[]): Promise<any[]> {
return Promise.resolve(
keys.map((key) => this.store.get(key.join(memoryKeySeparator))),
);
}
async batchSet(entries: KV[]): Promise<void> {
for (const { key, value } of entries) {
this.store.set(key.join(memoryKeySeparator), value);
}
// Trigger persistence
if (this.throttledPersist) {
this.throttledPersist();
} else if (this.filePath) {
// If no throttling is set up but we have a filePath, persist immediately
await this.persistToDisk();
}
return Promise.resolve();
}
async batchDelete(keys: KvKey[]): Promise<void> {
for (const key of keys) {
this.store.delete(key.join(memoryKeySeparator));
}
// Trigger persistence
if (this.throttledPersist) {
this.throttledPersist();
} else if (this.filePath) {
// If no throttling is set up but we have a filePath, persist immediately
await this.persistToDisk();
}
return Promise.resolve();
}
toJSON(): Record<string, any> {
const result: Record<string, any> = {};
for (const [key, value] of this.store) {
result[key] = value;
}
return result;
}
async *query(options: KvQueryOptions): AsyncIterableIterator<KV> {
const prefix = options.prefix?.join(memoryKeySeparator);
const sortedKeys = [...this.store.keys()].sort();
for (const key of sortedKeys) {
if (prefix && !key.startsWith(prefix)) {
continue;
}
yield {
key: key.split(memoryKeySeparator),
value: this.store.get(key),
};
}
}
async close(): Promise<void> {
// Force immediate persistence when closing
if (this.filePath) {
await this.persistToDisk();
}
}
/**
* Persist the current state to disk
*/
private async persistToDisk(): Promise<void> {
if (!this.filePath) return;
try {
const jsonData = this.toJSON();
await Deno.writeTextFile(
this.filePath,
JSON.stringify(jsonData, null, 2),
);
} catch (error) {
console.error(`Failed to persist KV store to ${this.filePath}:`, error);
}
}
}
+229
View File
@@ -0,0 +1,229 @@
import { DataStoreMQ } from "./mq.datastore.ts";
import { assertEquals } from "@std/assert";
import { MemoryKvPrimitives } from "./memory_kv_primitives.ts";
import { DataStore } from "./datastore.ts";
import { PrefixedKvPrimitives } from "./prefixed_kv_primitives.ts";
import { FakeTime } from "@std/testing/time";
import type { MQMessage } from "../../plug-api/types/datastore.ts";
Deno.test("DataStore MQ", async () => {
const time = new FakeTime();
const db = new MemoryKvPrimitives(); // In-memory only, no persistence
try {
const mq = new DataStoreMQ(
new DataStore(new PrefixedKvPrimitives(db, ["mq"])),
);
let messages: MQMessage[];
// Send and ack
await mq.send("test", "Hello World");
messages = await mq.poll("test", 10);
assertEquals(messages.length, 1);
await mq.ack("test", messages[0].id);
assertEquals([], await mq.poll("test", 10));
// Timeout
await mq.send("test", "Hello World");
messages = await mq.poll("test", 10);
assertEquals(messages.length, 1);
assertEquals([], await mq.poll("test", 10));
await time.tickAsync(20);
await mq.requeueTimeouts(10);
messages = await mq.poll("test", 10);
const stats = await mq.getAllQueueStats();
assertEquals(stats["test"].processing, 1);
assertEquals(messages.length, 1);
assertEquals(messages[0].retries, 1);
// Max retries
await time.tickAsync(20);
await mq.requeueTimeouts(10, 1);
assertEquals((await mq.fetchDLQMessages()).length, 1);
// Batch send and ack
await mq.batchSend("test", ["Hello", "World"]);
const messageBatch1 = await mq.poll("test", 1);
assertEquals(messageBatch1.length, 1);
assertEquals(messageBatch1[0].body, "Hello");
const messageBatch2 = await mq.poll("test", 1);
assertEquals(messageBatch2.length, 1);
assertEquals(messageBatch2[0].body, "World");
await mq.batchAck("test", [messageBatch1[0].id, messageBatch2[0].id]);
assertEquals(await mq.fetchProcessingMessages(), []);
// Subscribe
let receivedMessage = false;
const worker = mq.subscribe("test123", {}, async (messages) => {
assertEquals(messages.length, 1);
receivedMessage = true;
await mq.ack("test123", messages[0].id);
});
await mq.send("test123", "Hello World");
// Wait for message to be processed by checking queue stats
while ((await mq.getQueueStats("test123")).queued > 0) {
await time.tickAsync(100);
}
assertEquals(receivedMessage, true);
worker.stop();
assertEquals(mq.queueWaiters.size, 0);
} finally {
await db.close();
time.restore();
}
});
Deno.test("DataStore MQ - Scale test with multiple subscribers", async () => {
const time = new FakeTime();
const db = new MemoryKvPrimitives();
try {
const mq = new DataStoreMQ(
new DataStore(new PrefixedKvPrimitives(db, ["mq"])),
);
const queueName = "scale-test";
const totalMessages = 1000;
const batchSize = 7;
const numSubscribers = 3;
// Track processed messages across all subscribers
const processedMessages = new Set<string>();
const subscriberStats = new Map<number, number>();
const processingLock = new Set<string>(); // Track which messages are being processed
// Initialize subscriber stats
for (let i = 0; i < numSubscribers; i++) {
subscriberStats.set(i, 0);
}
// Create 3 subscribers, each processing batches of 7 messages
const workers = [];
for (let subscriberId = 0; subscriberId < numSubscribers; subscriberId++) {
const worker = mq.subscribe(
queueName,
{ batchSize },
async (messages) => {
assertEquals(
messages.length <= batchSize,
true,
`Batch size should not exceed ${batchSize}`,
);
console.log(
`[Subscriber ${subscriberId}] Processing batch of ${messages.length} messages`,
);
// Process each message in the batch
const messageIds = [];
for (const message of messages) {
// Check for concurrent processing (this shouldn't happen due to MQ design)
if (processingLock.has(message.body)) {
console.warn(
`Message ${message.body} is being processed concurrently by subscriber ${subscriberId}`,
);
continue;
}
processingLock.add(message.body);
// Ensure no duplicate processing
if (processedMessages.has(message.body)) {
console.warn(
`Message ${message.body} already processed by another subscriber`,
);
processingLock.delete(message.body);
continue;
}
processedMessages.add(message.body);
messageIds.push(message.id);
// Update subscriber stats
const currentCount = subscriberStats.get(subscriberId) || 0;
subscriberStats.set(subscriberId, currentCount + 1);
// Remove from processing lock
processingLock.delete(message.body);
}
// Ack all messages in the batch that were actually processed
if (messageIds.length > 0) {
await mq.batchAck(queueName, messageIds);
}
},
);
workers.push(worker);
}
// Send messages in smaller batches to allow better distribution
const messageBodies = [];
for (let i = 0; i < totalMessages; i++) {
messageBodies.push(`message-${i}`);
}
// Send in chunks to allow workers to process and become available
const chunkSize = 50;
for (let i = 0; i < messageBodies.length; i += chunkSize) {
const chunk = messageBodies.slice(i, i + chunkSize);
await mq.batchSend(queueName, chunk);
// Give a small delay to allow processing
await time.tickAsync(10);
}
// Wait for all messages to be processed
const maxWaitTime = 10000; // 10 seconds max wait
let waitTime = 0;
while (processedMessages.size < totalMessages && waitTime < maxWaitTime) {
await time.tickAsync(100);
waitTime += 100;
}
// Verify all messages were processed
assertEquals(
processedMessages.size,
totalMessages,
`Expected ${totalMessages} messages processed, got ${processedMessages.size}`,
);
// Wait a bit more to ensure all acks are processed
await time.tickAsync(100);
// Verify queue is empty
const stats = await mq.getQueueStats(queueName);
assertEquals(stats.queued, 0, "Queue should be empty");
assertEquals(stats.processing, 0, "No messages should be processing");
// Verify work distribution among subscribers
let totalProcessedAcrossSubscribers = 0;
for (const [subscriberId, count] of subscriberStats.entries()) {
totalProcessedAcrossSubscribers += count;
console.log(`Subscriber ${subscriberId} processed ${count} messages`);
}
assertEquals(
totalProcessedAcrossSubscribers,
totalMessages,
"Total processed should match sent messages",
);
// Verify at least one subscriber processed messages (relaxed requirement since MQ may favor one worker)
const activeSubscribers =
Array.from(subscriberStats.values()).filter((count) => count > 0).length;
assertEquals(
activeSubscribers >= 1,
true,
"At least one subscriber should have processed messages",
);
// Stop all workers
workers.forEach((worker) => worker.stop());
// Verify no queue waiters remain
assertEquals(mq.queueWaiters.size, 0);
} finally {
await db.close();
time.restore();
}
});
+423
View File
@@ -0,0 +1,423 @@
import type { DataStore } from "./datastore.ts";
import { parseExpressionString } from "../space_lua/parse.ts";
import { LuaEnv } from "../space_lua/runtime.ts";
import type {
KV,
KvKey,
MQMessage,
MQStats,
MQSubscribeOptions,
} from "../../plug-api/types/datastore.ts";
import { race, sleep } from "@silverbulletmd/silverbullet/lib/async";
export type ProcessingMessage = MQMessage & {
ts: number;
};
const queuedPrefix = ["mq", "queued"];
const processingPrefix = ["mq", "processing"];
const dlqPrefix = ["mq", "dlq"];
export class QueueWorker {
stopping = false;
private stopReject?: (e: any) => void;
constructor(
private mq: DataStoreMQ,
readonly queue: string,
readonly options: MQSubscribeOptions,
private callback: (messages: MQMessage[]) => Promise<void> | void,
) {
}
/**
* This is the main loop of the worker, whenever it exits the loop it means the worker has stopped
*/
async run() {
try {
while (true) {
if (this.stopping) {
break;
}
// Poll for messages
const messages = await this.mq.poll(
this.queue,
this.options.batchSize || 1,
);
if (messages.length > 0) {
// We have messages, process them, then immediately loop to poll again
await this.callback(messages);
} else {
// No messages, wait to be woken up or a timeout
try {
await race([
// Wait to be woken up explicitly
new Promise<void>((resolve, reject) => {
this.stopReject = reject;
this.mq.queueWorker(this.queue, resolve, reject);
}),
// Or a poll interval timeout
sleep(this.options.pollInterval || 1000).then(() => {
// Remove self from waiters
this.mq.removeQueuedWorker(this.queue, this.stopReject!);
}),
]);
} catch (e: any) {
// Only scenario we should end up here is stop being called
console.info(e.message);
break;
}
}
}
} catch (e) {
console.error("Error in queue worker", e);
}
}
stop() {
this.stopping = true;
if (this.stopReject) {
// Worker was in a waiting state, reject the promise to wake it up and remove from waiters
this.mq.removeQueuedWorker(this.queue, this.stopReject);
this.stopReject(new Error("Queue worker stopped"));
}
}
}
export class DataStoreMQ {
// Internal sequencer for messages, only really necessary when batch sending tons of messages within a millisecond
seq = 0;
queueWaiters = new Map<
string,
({ resolve: () => void; reject: (e: any) => void })[]
>();
constructor(
private ds: DataStore,
) {
}
/// Worker management
public queueWorker(
queue: string,
resolve: () => void,
reject: (e: any) => void,
) {
let waiters = this.queueWaiters.get(queue);
if (!waiters) {
waiters = [];
this.queueWaiters.set(queue, waiters);
}
// console.log("[mq]", "Queuing a worker for queue", queue);
waiters.push({ resolve, reject });
}
/**
* Wakes up a single worker waiting on the given queue, if any
* @param queue
*/
wakeupWorker(queue: string) {
const waiters = this.queueWaiters.get(queue);
if (waiters && waiters.length > 0) {
// console.log("[mq]", "Waking up a worker for queue", queue);
const { resolve } = waiters.shift()!;
resolve();
if (waiters.length === 0) {
// Clean up empty arrays
this.queueWaiters.delete(queue);
}
}
}
removeQueuedWorker(queue: string, reject: (e: any) => void) {
const waiters = this.queueWaiters.get(queue);
if (waiters) {
const index = waiters.findIndex((w) => w.reject === reject);
if (index !== -1) {
waiters.splice(index, 1);
}
if (waiters.length === 0) {
// Let's not keep empty arrays around
this.queueWaiters.delete(queue);
}
}
}
/**
* Sends a batch of messages to a queue.
* @param queue the name of the queue
* @param bodies the bodies of the messages to send
* @returns
*/
async batchSend(queue: string, bodies: any[]): Promise<void> {
if (bodies.length === 0) {
return;
}
const messages: KV<MQMessage>[] = bodies.map((body) => {
const id = `${Date.now()}-${String(++this.seq).padStart(6, "0")}`;
const key = [...queuedPrefix, queue, id];
return {
key,
value: { id, queue, body },
};
});
await this.ds.batchSet(messages);
this.wakeupWorker(queue);
}
send(queue: string, body: any): Promise<void> {
return this.batchSend(queue, [body]);
}
async poll(queue: string, maxItems: number): Promise<MQMessage[]> {
// Note: this is not happening in a transactional way, so we may get duplicate message delivery
// Retrieve a batch of messages
const messages = await this.ds.luaQuery<MQMessage>(
[...queuedPrefix, queue],
{
limit: maxItems,
},
);
if (messages.length === 0) {
return [];
}
// Put them in the processing queue
await this.ds.batchSet(
messages.map((m) => ({
key: [...processingPrefix, queue, m.id],
value: {
...m,
ts: Date.now(),
},
})),
);
// Delete them from the queued queue
await this.ds.batchDelete(
messages.map((m) => [...queuedPrefix, queue, m.id]),
);
// Return them
return messages;
}
/**
* @param queue
* @param batchSize
* @param callback
* @returns a function to be called to unsubscribe
*/
subscribe(
queue: string,
options: MQSubscribeOptions,
callback: (messages: MQMessage[]) => Promise<void> | void,
): QueueWorker {
const worker = new QueueWorker(this, queue, options, callback);
// Start the worker asynchronously
worker.run();
return worker;
}
ack(queue: string, id: string) {
return this.batchAck(queue, [id]);
}
async batchAck(queue: string, ids: string[]) {
if (ids.length === 0) {
return;
}
await this.ds.batchDelete(
ids.map((id) => [...processingPrefix, queue, id]),
);
}
async requeueTimeouts(
timeout: number,
maxRetries?: number,
disableDLQ?: boolean,
) {
const now = Date.now();
const env = new LuaEnv();
env.setLocal("ts", now - timeout);
const messages = await this.ds.luaQuery<ProcessingMessage>(
processingPrefix,
{
objectVariable: "m",
where: parseExpressionString("m.ts < ts"),
},
env,
);
if (messages.length === 0) {
return;
}
await this.ds.batchDelete(
messages.map((m) => [...processingPrefix, m.queue, m.id]),
);
const newMessages: KV<ProcessingMessage>[] = [];
for (const m of messages) {
const retries = (m.retries || 0) + 1;
if (maxRetries && retries > maxRetries) {
if (disableDLQ) {
console.warn(
"[mq]",
"Message exceeded max retries, flushing message",
m,
);
} else {
console.warn(
"[mq]",
"Message exceeded max retries, moving to DLQ",
m,
);
newMessages.push({
key: [...dlqPrefix, m.queue, m.id],
value: {
queue: m.queue,
id: m.id,
body: m.body,
ts: Date.now(),
retries,
},
});
}
} else {
console.info("[mq]", "Message ack timed out, requeueing", m);
newMessages.push({
key: [...queuedPrefix, m.queue, m.id],
value: {
...m,
retries,
},
});
}
}
await this.ds.batchSet(newMessages);
}
async fetchDLQMessages(): Promise<ProcessingMessage[]> {
return (await this.ds.luaQuery<ProcessingMessage>(dlqPrefix, {}));
}
async fetchProcessingMessages(): Promise<ProcessingMessage[]> {
return (await this.ds.luaQuery<ProcessingMessage>(processingPrefix, {}));
}
async flushDLQ(): Promise<void> {
const ids: KvKey[] = [];
for (const item of await this.ds.luaQuery<MQMessage>(dlqPrefix, {})) {
ids.push([...dlqPrefix, item.queue, item.id]);
}
await this.ds.batchDelete(ids);
}
/**
* Flushes a queue, including all queued, processing and DLQ messages
* @param queue
*/
async flushQueue(queue: string): Promise<void> {
const ids: KvKey[] = [];
for (
const item of await this.ds.luaQuery<MQMessage>(
[...queuedPrefix, queue],
{},
)
) {
ids.push([...queuedPrefix, queue, item.id]);
}
for (
const item of await this.ds.luaQuery<ProcessingMessage>([
...processingPrefix,
queue,
], {})
) {
ids.push([...processingPrefix, queue, item.id]);
}
for (
const item of await this.ds.luaQuery<ProcessingMessage>([
...dlqPrefix,
queue,
], {})
) {
ids.push([...dlqPrefix, queue, item.id]);
}
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;
return {
queued,
processing,
dlq,
};
}
async awaitEmptyQueue(queue: string): Promise<void> {
while (true) {
const stats = await this.getQueueStats(queue);
if (stats.queued === 0 && stats.processing === 0) {
break;
}
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;
}
}
+51
View File
@@ -0,0 +1,51 @@
import type { KvPrimitives, KvQueryOptions } from "./kv_primitives.ts";
import type { KV, KvKey } from "../../plug-api/types/datastore.ts";
/**
* Turns any KvPrimitives into a KvPrimitives that automatically prefixes all keys (and removes them again when reading)
*/
export class PrefixedKvPrimitives implements KvPrimitives {
constructor(private wrapped: KvPrimitives, private prefix: KvKey) {
}
clear(): Promise<void> {
return this.wrapped.clear();
}
batchGet(keys: KvKey[]): Promise<any[]> {
return this.wrapped.batchGet(keys.map((key) => this.applyPrefix(key)));
}
batchSet(entries: KV[]): Promise<void> {
return this.wrapped.batchSet(
entries.map(({ key, value }) => ({ key: this.applyPrefix(key), value })),
);
}
batchDelete(keys: KvKey[]): Promise<void> {
return this.wrapped.batchDelete(keys.map((key) => this.applyPrefix(key)));
}
async *query(options: KvQueryOptions): AsyncIterableIterator<KV> {
for await (
const result of this.wrapped.query({
prefix: this.applyPrefix(options.prefix),
})
) {
yield { key: this.stripPrefix(result.key), value: result.value };
}
}
close(): void {
this.wrapped.close();
}
private applyPrefix(key?: KvKey): KvKey {
return [...this.prefix, ...(key ? key : [])];
}
private stripPrefix(key: KvKey): KvKey {
return key.slice(this.prefix.length);
}
}