refactor(FN-1633): migrate MessageStore from filesystem to SQLite backend

- Replace filesystem-based message storage with SQLite backend
- Add MessageStore class using better-sqlite3 with WAL mode
- Update message.ts CLI command to use new MessageStore API
- Update dashboard routes and engine runtime for SQLite integration
- Update all related tests for new storage implementation
This commit is contained in:
Fusion
2026-04-15 07:02:38 -07:00
committed by gsxdsm
parent 940494272d
commit cd8dfa452f
7 changed files with 504 additions and 491 deletions

View File

@@ -2,20 +2,24 @@ import { describe, it, expect, beforeEach, afterEach, vi } from "vitest";
import { mkdtempSync, rmSync } from "node:fs";
import { join } from "node:path";
import { tmpdir } from "node:os";
import { Database } from "./db.js";
import { MessageStore } from "./message-store.js";
import type { Message, Mailbox } from "./types.js";
describe("MessageStore", () => {
let store: MessageStore;
let db: Database;
let tempDir: string;
beforeEach(async () => {
beforeEach(() => {
tempDir = mkdtempSync(join(tmpdir(), "kb-msg-test-"));
store = new MessageStore({ rootDir: tempDir });
await store.init();
db = new Database(tempDir);
db.init();
store = new MessageStore(db);
});
afterEach(() => {
db.close();
try {
rmSync(tempDir, { recursive: true, force: true });
} catch {
@@ -23,22 +27,9 @@ describe("MessageStore", () => {
}
});
describe("init()", () => {
it("creates messages directory and index file", async () => {
const { existsSync } = await import("node:fs");
expect(existsSync(join(tempDir, "messages"))).toBe(true);
expect(existsSync(join(tempDir, "messages", "index.json"))).toBe(true);
});
it("is idempotent — calling init twice does not throw", async () => {
await store.init();
await store.init();
});
});
describe("sendMessage() and getMessage()", () => {
it("creates and retrieves a message", async () => {
const message = await store.sendMessage({
it("creates and retrieves a message", () => {
const message = store.sendMessage({
fromId: "user-1",
fromType: "user",
toId: "agent-1",
@@ -59,12 +50,12 @@ describe("MessageStore", () => {
expect(message.createdAt).toBeTruthy();
expect(message.updatedAt).toBeTruthy();
const retrieved = await store.getMessage(message.id);
const retrieved = store.getMessage(message.id);
expect(retrieved).toEqual(message);
});
it("auto-fills sender as system when not provided", async () => {
const message = await store.sendMessage({
it("auto-fills sender as system when not provided", () => {
const message = store.sendMessage({
toId: "user-1",
toType: "user",
content: "System notification",
@@ -75,8 +66,8 @@ describe("MessageStore", () => {
expect(message.fromType).toBe("system");
});
it("stores metadata when provided", async () => {
const message = await store.sendMessage({
it("stores metadata when provided", () => {
const message = store.sendMessage({
fromId: "agent-1",
fromType: "agent",
toId: "user-1",
@@ -89,19 +80,18 @@ describe("MessageStore", () => {
expect(message.metadata).toEqual({ taskId: "FN-001", priority: "high" });
});
it("returns null for non-existent message", async () => {
const result = await store.getMessage("msg-nonexistent");
it("returns null for non-existent message", () => {
const result = store.getMessage("msg-nonexistent");
expect(result).toBeNull();
});
});
describe("message-to-agent hook", () => {
it("does not call the hook for non-agent recipients", async () => {
it("does not call the hook for non-agent recipients", () => {
const hook = vi.fn();
const hookedStore = new MessageStore({ rootDir: tempDir, onMessageToAgent: hook });
await hookedStore.init();
const hookedStore = new MessageStore(db, { onMessageToAgent: hook });
await hookedStore.sendMessage({
hookedStore.sendMessage({
fromId: "agent-1",
fromType: "agent",
toId: "user-1",
@@ -113,12 +103,11 @@ describe("MessageStore", () => {
expect(hook).not.toHaveBeenCalled();
});
it("calls the hook when a message is sent to an agent", async () => {
it("calls the hook when a message is sent to an agent", () => {
const hook = vi.fn();
const hookedStore = new MessageStore({ rootDir: tempDir, onMessageToAgent: hook });
await hookedStore.init();
const hookedStore = new MessageStore(db, { onMessageToAgent: hook });
const message = await hookedStore.sendMessage({
const message = hookedStore.sendMessage({
fromId: "user-1",
fromType: "user",
toId: "agent-1",
@@ -131,8 +120,8 @@ describe("MessageStore", () => {
expect(hook).toHaveBeenCalledWith(message);
});
it("does nothing when no hook is configured", async () => {
await expect(
it("does nothing when no hook is configured", () => {
expect(() => {
store.sendMessage({
fromId: "user-1",
fromType: "user",
@@ -140,17 +129,16 @@ describe("MessageStore", () => {
toType: "agent",
content: "No hook configured",
type: "user-to-agent",
}),
).resolves.toMatchObject({ toId: "agent-1", toType: "agent" });
});
}).not.toThrow();
});
it("setMessageToAgentHook updates the hook used for subsequent messages", async () => {
it("setMessageToAgentHook updates the hook used for subsequent messages", () => {
const firstHook = vi.fn();
const secondHook = vi.fn();
const hookedStore = new MessageStore({ rootDir: tempDir, onMessageToAgent: firstHook });
await hookedStore.init();
const hookedStore = new MessageStore(db, { onMessageToAgent: firstHook });
await hookedStore.sendMessage({
hookedStore.sendMessage({
fromId: "user-1",
fromType: "user",
toId: "agent-1",
@@ -161,7 +149,7 @@ describe("MessageStore", () => {
hookedStore.setMessageToAgentHook(secondHook);
await hookedStore.sendMessage({
hookedStore.sendMessage({
fromId: "user-1",
fromType: "user",
toId: "agent-1",
@@ -176,8 +164,8 @@ describe("MessageStore", () => {
});
describe("getInbox()", () => {
it("returns inbox messages for a participant", async () => {
await store.sendMessage({
it("returns inbox messages for a participant", () => {
store.sendMessage({
fromId: "agent-1",
fromType: "agent",
toId: "user-1",
@@ -186,7 +174,7 @@ describe("MessageStore", () => {
type: "agent-to-user",
});
await store.sendMessage({
store.sendMessage({
fromId: "agent-2",
fromType: "agent",
toId: "user-1",
@@ -195,20 +183,20 @@ describe("MessageStore", () => {
type: "agent-to-user",
});
const inbox = await store.getInbox("user-1", "user");
const inbox = store.getInbox("user-1", "user");
expect(inbox).toHaveLength(2);
// Newest first
expect(inbox[0].content).toBe("Message 2");
expect(inbox[1].content).toBe("Message 1");
});
it("returns empty array for participant with no messages", async () => {
const inbox = await store.getInbox("user-99", "user");
it("returns empty array for participant with no messages", () => {
const inbox = store.getInbox("user-99", "user");
expect(inbox).toEqual([]);
});
it("filters by read status", async () => {
const msg1 = await store.sendMessage({
it("filters by read status", () => {
const msg1 = store.sendMessage({
fromId: "agent-1",
fromType: "agent",
toId: "user-1",
@@ -217,7 +205,7 @@ describe("MessageStore", () => {
type: "agent-to-user",
});
const msg2 = await store.sendMessage({
const msg2 = store.sendMessage({
fromId: "agent-1",
fromType: "agent",
toId: "user-1",
@@ -226,20 +214,20 @@ describe("MessageStore", () => {
type: "agent-to-user",
});
await store.markAsRead(msg2.id);
store.markAsRead(msg2.id);
const unreadOnly = await store.getInbox("user-1", "user", { read: false });
const unreadOnly = store.getInbox("user-1", "user", { read: false });
expect(unreadOnly).toHaveLength(1);
expect(unreadOnly[0].id).toBe(msg1.id);
const readOnly = await store.getInbox("user-1", "user", { read: true });
const readOnly = store.getInbox("user-1", "user", { read: true });
expect(readOnly).toHaveLength(1);
expect(readOnly[0].id).toBe(msg2.id);
});
it("applies pagination (limit/offset)", async () => {
it("applies pagination (limit/offset)", () => {
for (let i = 0; i < 5; i++) {
await store.sendMessage({
store.sendMessage({
fromId: "agent-1",
fromType: "agent",
toId: "user-1",
@@ -249,18 +237,18 @@ describe("MessageStore", () => {
});
}
const page1 = await store.getInbox("user-1", "user", { limit: 2, offset: 0 });
const page1 = store.getInbox("user-1", "user", { limit: 2, offset: 0 });
expect(page1).toHaveLength(2);
const page2 = await store.getInbox("user-1", "user", { limit: 2, offset: 2 });
const page2 = store.getInbox("user-1", "user", { limit: 2, offset: 2 });
expect(page2).toHaveLength(2);
// No overlap
expect(page1[0].id).not.toBe(page2[0].id);
});
it("filters by message type", async () => {
await store.sendMessage({
it("filters by message type", () => {
store.sendMessage({
fromId: "agent-1",
fromType: "agent",
toId: "user-1",
@@ -269,7 +257,7 @@ describe("MessageStore", () => {
type: "agent-to-user",
});
await store.sendMessage({
store.sendMessage({
fromId: "system",
fromType: "system",
toId: "user-1",
@@ -278,19 +266,19 @@ describe("MessageStore", () => {
type: "system",
});
const agentOnly = await store.getInbox("user-1", "user", { type: "agent-to-user" });
const agentOnly = store.getInbox("user-1", "user", { type: "agent-to-user" });
expect(agentOnly).toHaveLength(1);
expect(agentOnly[0].type).toBe("agent-to-user");
const systemOnly = await store.getInbox("user-1", "user", { type: "system" });
const systemOnly = store.getInbox("user-1", "user", { type: "system" });
expect(systemOnly).toHaveLength(1);
expect(systemOnly[0].type).toBe("system");
});
});
describe("getOutbox()", () => {
it("returns sent messages for a participant", async () => {
await store.sendMessage({
it("returns sent messages for a participant", () => {
store.sendMessage({
fromId: "user-1",
fromType: "user",
toId: "agent-1",
@@ -299,7 +287,7 @@ describe("MessageStore", () => {
type: "user-to-agent",
});
await store.sendMessage({
store.sendMessage({
fromId: "user-1",
fromType: "user",
toId: "agent-2",
@@ -308,21 +296,21 @@ describe("MessageStore", () => {
type: "user-to-agent",
});
const outbox = await store.getOutbox("user-1", "user");
const outbox = store.getOutbox("user-1", "user");
expect(outbox).toHaveLength(2);
expect(outbox[0].content).toBe("Outgoing 2");
expect(outbox[1].content).toBe("Outgoing 1");
});
it("returns empty array when no messages sent", async () => {
const outbox = await store.getOutbox("user-99", "user");
it("returns empty array when no messages sent", () => {
const outbox = store.getOutbox("user-99", "user");
expect(outbox).toEqual([]);
});
});
describe("markAsRead()", () => {
it("marks a message as read", async () => {
const message = await store.sendMessage({
it("marks a message as read", () => {
const message = store.sendMessage({
fromId: "agent-1",
fromType: "agent",
toId: "user-1",
@@ -333,15 +321,15 @@ describe("MessageStore", () => {
expect(message.read).toBe(false);
const updated = await store.markAsRead(message.id);
const updated = store.markAsRead(message.id);
expect(updated.read).toBe(true);
const retrieved = await store.getMessage(message.id);
const retrieved = store.getMessage(message.id);
expect(retrieved!.read).toBe(true);
});
it("is idempotent for already-read messages", async () => {
const message = await store.sendMessage({
it("is idempotent for already-read messages", () => {
const message = store.sendMessage({
fromId: "agent-1",
fromType: "agent",
toId: "user-1",
@@ -350,19 +338,19 @@ describe("MessageStore", () => {
type: "agent-to-user",
});
await store.markAsRead(message.id);
const updated = await store.markAsRead(message.id);
store.markAsRead(message.id);
const updated = store.markAsRead(message.id);
expect(updated.read).toBe(true);
});
it("throws for non-existent message", async () => {
await expect(store.markAsRead("msg-nonexistent")).rejects.toThrow("not found");
it("throws for non-existent message", () => {
expect(() => store.markAsRead("msg-nonexistent")).toThrow("not found");
});
});
describe("markAllAsRead()", () => {
it("marks all unread messages as read", async () => {
await store.sendMessage({
it("marks all unread messages as read", () => {
store.sendMessage({
fromId: "agent-1",
fromType: "agent",
toId: "user-1",
@@ -371,7 +359,7 @@ describe("MessageStore", () => {
type: "agent-to-user",
});
await store.sendMessage({
store.sendMessage({
fromId: "agent-2",
fromType: "agent",
toId: "user-1",
@@ -380,22 +368,22 @@ describe("MessageStore", () => {
type: "agent-to-user",
});
const count = await store.markAllAsRead("user-1", "user");
const count = store.markAllAsRead("user-1", "user");
expect(count).toBe(2);
const inbox = await store.getInbox("user-1", "user");
const inbox = store.getInbox("user-1", "user");
expect(inbox.every((m) => m.read)).toBe(true);
});
it("returns 0 when no unread messages", async () => {
const count = await store.markAllAsRead("user-99", "user");
it("returns 0 when no unread messages", () => {
const count = store.markAllAsRead("user-99", "user");
expect(count).toBe(0);
});
});
describe("deleteMessage()", () => {
it("deletes a message", async () => {
const message = await store.sendMessage({
it("deletes a message", () => {
const message = store.sendMessage({
fromId: "user-1",
fromType: "user",
toId: "agent-1",
@@ -404,14 +392,14 @@ describe("MessageStore", () => {
type: "user-to-agent",
});
await store.deleteMessage(message.id);
store.deleteMessage(message.id);
const retrieved = await store.getMessage(message.id);
const retrieved = store.getMessage(message.id);
expect(retrieved).toBeNull();
});
it("removes message from inbox index", async () => {
const message = await store.sendMessage({
it("removes message from inbox", () => {
const message = store.sendMessage({
fromId: "agent-1",
fromType: "agent",
toId: "user-1",
@@ -420,14 +408,14 @@ describe("MessageStore", () => {
type: "agent-to-user",
});
await store.deleteMessage(message.id);
store.deleteMessage(message.id);
const inbox = await store.getInbox("user-1", "user");
const inbox = store.getInbox("user-1", "user");
expect(inbox).toHaveLength(0);
});
it("removes message from outbox index", async () => {
const message = await store.sendMessage({
it("removes message from outbox", () => {
const message = store.sendMessage({
fromId: "user-1",
fromType: "user",
toId: "agent-1",
@@ -436,21 +424,21 @@ describe("MessageStore", () => {
type: "user-to-agent",
});
await store.deleteMessage(message.id);
store.deleteMessage(message.id);
const outbox = await store.getOutbox("user-1", "user");
const outbox = store.getOutbox("user-1", "user");
expect(outbox).toHaveLength(0);
});
it("throws for non-existent message", async () => {
await expect(store.deleteMessage("msg-nonexistent")).rejects.toThrow("not found");
it("throws for non-existent message", () => {
expect(() => store.deleteMessage("msg-nonexistent")).toThrow("not found");
});
});
describe("getConversation()", () => {
it("returns all messages between two participants", async () => {
it("returns all messages between two participants", () => {
// user-1 sends to agent-1
await store.sendMessage({
store.sendMessage({
fromId: "user-1",
fromType: "user",
toId: "agent-1",
@@ -460,7 +448,7 @@ describe("MessageStore", () => {
});
// agent-1 replies to user-1
await store.sendMessage({
store.sendMessage({
fromId: "agent-1",
fromType: "agent",
toId: "user-1",
@@ -470,7 +458,7 @@ describe("MessageStore", () => {
});
// Unrelated message
await store.sendMessage({
store.sendMessage({
fromId: "agent-2",
fromType: "agent",
toId: "user-1",
@@ -479,7 +467,7 @@ describe("MessageStore", () => {
type: "agent-to-user",
});
const conversation = await store.getConversation(
const conversation = store.getConversation(
{ id: "user-1", type: "user" },
{ id: "agent-1", type: "agent" },
);
@@ -490,8 +478,8 @@ describe("MessageStore", () => {
expect(conversation[1].content).toBe("Hi there");
});
it("returns empty array when no conversation exists", async () => {
const conversation = await store.getConversation(
it("returns empty array when no conversation exists", () => {
const conversation = store.getConversation(
{ id: "user-1", type: "user" },
{ id: "agent-99", type: "agent" },
);
@@ -500,8 +488,8 @@ describe("MessageStore", () => {
});
describe("getMailbox()", () => {
it("returns mailbox summary with unread count", async () => {
await store.sendMessage({
it("returns mailbox summary with unread count", () => {
store.sendMessage({
fromId: "agent-1",
fromType: "agent",
toId: "user-1",
@@ -510,7 +498,7 @@ describe("MessageStore", () => {
type: "agent-to-user",
});
await store.sendMessage({
store.sendMessage({
fromId: "agent-1",
fromType: "agent",
toId: "user-1",
@@ -519,7 +507,7 @@ describe("MessageStore", () => {
type: "agent-to-user",
});
const mailbox = await store.getMailbox("user-1", "user");
const mailbox = store.getMailbox("user-1", "user");
expect(mailbox.ownerId).toBe("user-1");
expect(mailbox.ownerType).toBe("user");
@@ -528,14 +516,14 @@ describe("MessageStore", () => {
expect(mailbox.lastMessage!.content).toBe("Unread 2");
});
it("returns 0 unread when no messages", async () => {
const mailbox = await store.getMailbox("user-99", "user");
it("returns 0 unread when no messages", () => {
const mailbox = store.getMailbox("user-99", "user");
expect(mailbox.unreadCount).toBe(0);
expect(mailbox.lastMessage).toBeUndefined();
});
it("counts only unread messages", async () => {
const msg1 = await store.sendMessage({
it("counts only unread messages", () => {
const msg1 = store.sendMessage({
fromId: "agent-1",
fromType: "agent",
toId: "user-1",
@@ -544,7 +532,7 @@ describe("MessageStore", () => {
type: "agent-to-user",
});
await store.sendMessage({
store.sendMessage({
fromId: "agent-1",
fromType: "agent",
toId: "user-1",
@@ -553,19 +541,19 @@ describe("MessageStore", () => {
type: "agent-to-user",
});
await store.markAsRead(msg1.id);
store.markAsRead(msg1.id);
const mailbox = await store.getMailbox("user-1", "user");
const mailbox = store.getMailbox("user-1", "user");
expect(mailbox.unreadCount).toBe(1);
});
});
describe("events", () => {
it("emits message:sent event on send", async () => {
it("emits message:sent event on send", () => {
const events: Message[] = [];
store.on("message:sent", (msg) => events.push(msg));
await store.sendMessage({
store.sendMessage({
fromId: "user-1",
fromType: "user",
toId: "agent-1",
@@ -578,11 +566,11 @@ describe("MessageStore", () => {
expect(events[0].content).toBe("Hello");
});
it("emits message:received event on send", async () => {
it("emits message:received event on send", () => {
const events: Message[] = [];
store.on("message:received", (msg) => events.push(msg));
await store.sendMessage({
store.sendMessage({
fromId: "user-1",
fromType: "user",
toId: "agent-1",
@@ -594,11 +582,11 @@ describe("MessageStore", () => {
expect(events).toHaveLength(1);
});
it("emits message:read event on mark as read", async () => {
it("emits message:read event on mark as read", () => {
const events: Message[] = [];
store.on("message:read", (msg) => events.push(msg));
const message = await store.sendMessage({
const message = store.sendMessage({
fromId: "agent-1",
fromType: "agent",
toId: "user-1",
@@ -607,17 +595,17 @@ describe("MessageStore", () => {
type: "agent-to-user",
});
await store.markAsRead(message.id);
store.markAsRead(message.id);
expect(events).toHaveLength(1);
expect(events[0].read).toBe(true);
});
it("emits message:deleted event on delete", async () => {
it("emits message:deleted event on delete", () => {
const events: string[] = [];
store.on("message:deleted", (id) => events.push(id));
const message = await store.sendMessage({
const message = store.sendMessage({
fromId: "user-1",
fromType: "user",
toId: "agent-1",
@@ -626,7 +614,7 @@ describe("MessageStore", () => {
type: "user-to-agent",
});
await store.deleteMessage(message.id);
store.deleteMessage(message.id);
expect(events).toHaveLength(1);
expect(events[0]).toBe(message.id);

View File

@@ -1,19 +1,19 @@
/**
* MessageStore - Filesystem-based persistence for the messaging system.
* MessageStore - SQLite-based persistence for the messaging system.
*
* Messages are stored at `.fusion/messages/{messageId}.json` with their metadata.
* An index file at `.fusion/messages/index.json` provides efficient mailbox lookups.
* Messages are stored in the `messages` table with indexed lookups
* for inbox/outbox/conversation queries.
*
* File Structure:
* - messages/{messageId}.json: Individual message data
* - messages/index.json: Owner-to-message index for inbox/outbox queries
* Follows the same patterns as ChatStore:
* - EventEmitter for change notifications
* - SQLite for structured data storage (synchronous)
* - JSON columns for optional metadata
*/
import { mkdir, readFile, writeFile, readdir, unlink, rename } from "node:fs/promises";
import { existsSync } from "node:fs";
import { join } from "node:path";
import { randomUUID } from "node:crypto";
import { EventEmitter } from "node:events";
import { randomUUID } from "node:crypto";
import type { Database } from "./db.js";
import { fromJson, toJsonNullable } from "./db.js";
import type {
Message,
MessageCreateInput,
@@ -23,67 +23,109 @@ import type {
ParticipantType,
} from "./types.js";
// ── Event Types ─────────────────────────────────────────────────────
/** Events emitted by MessageStore */
export interface MessageStoreEvents {
/** Emitted when a new message is created and sent */
"message:sent": (message: Message) => void;
"message:sent": [message: Message];
/** Emitted when a message is received by a participant */
"message:received": (message: Message) => void;
"message:received": [message: Message];
/** Emitted when a message is marked as read */
"message:read": (message: Message) => void;
"message:read": [message: Message];
/** Emitted when a message is deleted */
"message:deleted": (messageId: string) => void;
"message:deleted": [messageId: string];
}
// ── Options Types ────────────────────────────────────────────────────
/** Options for MessageStore constructor */
export interface MessageStoreOptions {
/** Root directory for fn data (default: .fusion) */
rootDir?: string;
/** Optional hook invoked when a message is addressed to an agent */
onMessageToAgent?: (message: Message) => void;
}
/** Index structure for mailbox lookups */
interface MessageIndex {
/** Map of "type:id" -> { inbox: [msgId, ...], outbox: [msgId, ...] } */
byOwner: Record<string, { inbox: string[]; outbox: string[] }>;
}
// ── MessageStore Class ───────────────────────────────────────────────
/**
* MessageStore manages messages between agents, users, and the system.
* Uses filesystem-based persistence following the AgentStore pattern.
* Uses SQLite for persistent storage with efficient indexed queries.
*/
export class MessageStore extends EventEmitter {
private rootDir: string;
private messagesDir: string;
private indexPath: string;
export class MessageStore extends EventEmitter<MessageStoreEvents> {
private onMessageToAgent?: (message: Message) => void;
constructor(options: MessageStoreOptions = {}) {
// Prepared statements for frequently-run queries
private stmtInsert!: ReturnType<Database["prepare"]>;
private stmtGetById!: ReturnType<Database["prepare"]>;
private stmtUpdateRead!: ReturnType<Database["prepare"]>;
private stmtDelete!: ReturnType<Database["prepare"]>;
private stmtCountUnread!: ReturnType<Database["prepare"]>;
private stmtGetLastMessage!: ReturnType<Database["prepare"]>;
constructor(
private db: Database,
options?: MessageStoreOptions,
) {
super();
this.rootDir = options.rootDir ?? ".fusion";
this.messagesDir = join(this.rootDir, "messages");
this.indexPath = join(this.messagesDir, "index.json");
this.onMessageToAgent = options.onMessageToAgent;
this.setMaxListeners(100);
this.onMessageToAgent = options?.onMessageToAgent;
// Prepare frequently-run statements
this.stmtInsert = this.db.prepare(`
INSERT INTO messages (id, fromId, fromType, toId, toType, content, type, read, metadata, createdAt, updatedAt)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
`);
this.stmtGetById = this.db.prepare(`
SELECT * FROM messages WHERE id = ?
`);
this.stmtUpdateRead = this.db.prepare(`
UPDATE messages SET read = 1, updatedAt = ? WHERE id = ?
`);
this.stmtDelete = this.db.prepare(`
DELETE FROM messages WHERE id = ?
`);
this.stmtCountUnread = this.db.prepare(`
SELECT COUNT(*) as count FROM messages WHERE toId = ? AND toType = ? AND read = 0
`);
this.stmtGetLastMessage = this.db.prepare(`
SELECT * FROM messages WHERE toId = ? AND toType = ? ORDER BY createdAt DESC, rowid DESC LIMIT 1
`);
}
// ── Row-to-Object Converters ───────────────────────────────────────
/**
* Initialize the store by creating necessary directories and index file.
* Should be called before other operations.
* Convert a database row to a Message object.
*/
async init(): Promise<void> {
await mkdir(this.messagesDir, { recursive: true });
if (!existsSync(this.indexPath)) {
await this.writeIndex({ byOwner: {} });
}
private rowToMessage(row: any): Message {
return {
id: row.id,
fromId: row.fromId,
fromType: row.fromType as ParticipantType,
toId: row.toId,
toType: row.toType as ParticipantType,
content: row.content,
type: row.type as MessageType,
read: row.read === 1,
metadata: fromJson<Record<string, unknown>>(row.metadata),
createdAt: row.createdAt,
updatedAt: row.updatedAt,
};
}
// ── Public API ────────────────────────────────────────────────────
/**
* Create and store a new message.
* @param input - Message creation parameters
* @returns The created message
*/
async sendMessage(input: MessageCreateInput): Promise<Message> {
sendMessage(input: MessageCreateInput): Message {
const now = new Date().toISOString();
const messageId = `msg-${randomUUID().slice(0, 8)}`;
@@ -104,12 +146,21 @@ export class MessageStore extends EventEmitter {
updatedAt: now,
};
// Write message file
await this.writeMessageFile(message);
// Update index
await this.addToIndex(message);
this.stmtInsert.run(
message.id,
message.fromId,
message.fromType,
message.toId,
message.toType,
message.content,
message.type,
message.read ? 1 : 0,
toJsonNullable(message.metadata),
message.createdAt,
message.updatedAt,
);
this.db.bumpLastModified();
this.emit("message:sent", message);
this.emit("message:received", message);
@@ -125,17 +176,10 @@ export class MessageStore extends EventEmitter {
* @param id - The message ID
* @returns The message, or null if not found
*/
async getMessage(id: string): Promise<Message | null> {
try {
const path = join(this.messagesDir, `${id}.json`);
const content = await readFile(path, "utf-8");
return JSON.parse(content) as Message;
} catch (err) {
if ((err as NodeJS.ErrnoException).code === "ENOENT") {
return null;
}
throw err;
}
getMessage(id: string): Message | null {
const row = this.stmtGetById.get(id);
if (!row) return null;
return this.rowToMessage(row);
}
/**
@@ -145,17 +189,36 @@ export class MessageStore extends EventEmitter {
* @param filter - Optional filter criteria
* @returns Array of messages (newest first)
*/
async getInbox(
getInbox(
ownerId: string,
ownerType: ParticipantType,
filter?: MessageFilter,
): Promise<Message[]> {
const index = await this.readIndex();
const key = `${ownerType}:${ownerId}`;
const inboxIds = index.byOwner[key]?.inbox ?? [];
): Message[] {
const whereClauses: string[] = ["toId = ?", "toType = ?"];
const params: (string | number)[] = [ownerId, ownerType];
const messages = await this.loadMessagesByIds(inboxIds);
return this.applyFilter(messages, filter);
if (filter?.type) {
whereClauses.push("type = ?");
params.push(filter.type);
}
if (filter?.read !== undefined) {
whereClauses.push("read = ?");
params.push(filter.read ? 1 : 0);
}
const whereSql = whereClauses.join(" AND ");
const limit = filter?.limit ?? 100;
const offset = filter?.offset ?? 0;
const rows = this.db.prepare(`
SELECT * FROM messages
WHERE ${whereSql}
ORDER BY createdAt DESC, rowid DESC
LIMIT ? OFFSET ?
`).all(...params, limit, offset);
return (rows as any[]).map((row) => this.rowToMessage(row));
}
/**
@@ -165,17 +228,36 @@ export class MessageStore extends EventEmitter {
* @param filter - Optional filter criteria
* @returns Array of messages (newest first)
*/
async getOutbox(
getOutbox(
ownerId: string,
ownerType: ParticipantType,
filter?: MessageFilter,
): Promise<Message[]> {
const index = await this.readIndex();
const key = `${ownerType}:${ownerId}`;
const outboxIds = index.byOwner[key]?.outbox ?? [];
): Message[] {
const whereClauses: string[] = ["fromId = ?", "fromType = ?"];
const params: (string | number)[] = [ownerId, ownerType];
const messages = await this.loadMessagesByIds(outboxIds);
return this.applyFilter(messages, filter);
if (filter?.type) {
whereClauses.push("type = ?");
params.push(filter.type);
}
if (filter?.read !== undefined) {
whereClauses.push("read = ?");
params.push(filter.read ? 1 : 0);
}
const whereSql = whereClauses.join(" AND ");
const limit = filter?.limit ?? 100;
const offset = filter?.offset ?? 0;
const rows = this.db.prepare(`
SELECT * FROM messages
WHERE ${whereSql}
ORDER BY createdAt DESC, rowid DESC
LIMIT ? OFFSET ?
`).all(...params, limit, offset);
return (rows as any[]).map((row) => this.rowToMessage(row));
}
/**
@@ -184,24 +266,22 @@ export class MessageStore extends EventEmitter {
* @returns The updated message
* @throws Error if message not found
*/
async markAsRead(messageId: string): Promise<Message> {
const message = await this.getMessage(messageId);
if (!message) {
markAsRead(messageId: string): Message {
// First check if the message exists
const existing = this.getMessage(messageId);
if (!existing) {
throw new Error(`Message ${messageId} not found`);
}
if (message.read) return message;
if (existing.read) return existing;
const updated: Message = {
...message,
read: true,
updatedAt: new Date().toISOString(),
};
const now = new Date().toISOString();
this.stmtUpdateRead.run(now, messageId);
this.db.bumpLastModified();
await this.writeMessageFile(updated);
this.emit("message:read", updated);
return updated;
const updated = this.getMessage(messageId);
this.emit("message:read", updated!);
return updated!;
}
/**
@@ -210,19 +290,23 @@ export class MessageStore extends EventEmitter {
* @param ownerType - The participant type
* @returns Number of messages marked as read
*/
async markAllAsRead(
markAllAsRead(
ownerId: string,
ownerType: ParticipantType,
): Promise<number> {
const inbox = await this.getInbox(ownerId, ownerType);
const unread = inbox.filter((m) => !m.read);
): number {
const now = new Date().toISOString();
// Get count of unread messages before updating
const unreadRow = this.db.prepare(`
SELECT COUNT(*) as count FROM messages WHERE toId = ? AND toType = ? AND read = 0
`).get(ownerId, ownerType) as { count: number } | undefined;
const count = unreadRow?.count ?? 0;
let count = 0;
for (const message of unread) {
await this.markAsRead(message.id);
count++;
}
// Mark all as read
this.db.prepare(`
UPDATE messages SET read = 1, updatedAt = ? WHERE toId = ? AND toType = ? AND read = 0
`).run(now, ownerId, ownerType);
this.db.bumpLastModified();
return count;
}
@@ -231,19 +315,15 @@ export class MessageStore extends EventEmitter {
* @param id - The message ID
* @throws Error if message not found
*/
async deleteMessage(id: string): Promise<void> {
const message = await this.getMessage(id);
if (!message) {
deleteMessage(id: string): void {
// First check if the message exists
const existing = this.getMessage(id);
if (!existing) {
throw new Error(`Message ${id} not found`);
}
// Remove message file
const path = join(this.messagesDir, `${id}.json`);
await unlink(path);
// Remove from index
await this.removeFromIndex(message);
this.stmtDelete.run(id);
this.db.bumpLastModified();
this.emit("message:deleted", id);
}
@@ -253,28 +333,28 @@ export class MessageStore extends EventEmitter {
* @param participantB - Second participant
* @returns Array of messages (oldest first for conversation ordering)
*/
async getConversation(
getConversation(
participantA: { id: string; type: ParticipantType },
participantB: { id: string; type: ParticipantType },
): Promise<Message[]> {
const index = await this.readIndex();
const keyA = `${participantA.type}:${participantA.id}`;
const keyB = `${participantB.type}:${participantB.id}`;
): Message[] {
// Find messages where either participant is sender or receiver
// This captures all messages between the two participants
const rows = this.db.prepare(`
SELECT * FROM messages
WHERE (
(fromId = ? AND fromType = ? AND toId = ? AND toType = ?)
OR
(fromId = ? AND fromType = ? AND toId = ? AND toType = ?)
)
ORDER BY createdAt ASC
`).all(
participantA.id, participantA.type,
participantB.id, participantB.type,
participantB.id, participantB.type,
participantA.id, participantA.type,
);
const aInbox = index.byOwner[keyA]?.inbox ?? [];
const aOutbox = index.byOwner[keyA]?.outbox ?? [];
const allA = new Set([...aInbox, ...aOutbox]);
const bInbox = index.byOwner[keyB]?.inbox ?? [];
const bOutbox = index.byOwner[keyB]?.outbox ?? [];
const allB = new Set([...bInbox, ...bOutbox]);
// Find intersection: messages both participants have
const conversationIds = [...allA].filter((id) => allB.has(id));
const messages = await this.loadMessagesByIds(conversationIds);
// Conversation order: oldest first
return [...messages].reverse();
return (rows as any[]).map((row) => this.rowToMessage(row));
}
/**
@@ -283,13 +363,15 @@ export class MessageStore extends EventEmitter {
* @param ownerType - The participant type
* @returns Mailbox summary with unread count and last message
*/
async getMailbox(
getMailbox(
ownerId: string,
ownerType: ParticipantType,
): Promise<Mailbox> {
const inbox = await this.getInbox(ownerId, ownerType);
const unreadCount = inbox.filter((m) => !m.read).length;
const lastMessage = inbox.length > 0 ? inbox[0] : undefined;
): Mailbox {
const unreadRow = this.stmtCountUnread.get(ownerId, ownerType) as { count: number } | undefined;
const unreadCount = unreadRow?.count ?? 0;
const lastRow = this.stmtGetLastMessage.get(ownerId, ownerType);
const lastMessage = lastRow ? this.rowToMessage(lastRow) : undefined;
return {
ownerId,
@@ -305,100 +387,4 @@ export class MessageStore extends EventEmitter {
setMessageToAgentHook(hook: (message: Message) => void): void {
this.onMessageToAgent = hook;
}
// ─────────────────────────────────────────────────────────────────────────
// Private helpers
// ─────────────────────────────────────────────────────────────────────────
private async writeMessageFile(message: Message): Promise<void> {
const path = join(this.messagesDir, `${message.id}.json`);
const tempPath = `${path}.tmp.${Date.now()}`;
await writeFile(tempPath, JSON.stringify(message, null, 2));
await rename(tempPath, path);
}
private async readIndex(): Promise<MessageIndex> {
try {
const content = await readFile(this.indexPath, "utf-8");
return JSON.parse(content) as MessageIndex;
} catch {
return { byOwner: {} };
}
}
private async writeIndex(index: MessageIndex): Promise<void> {
const tempPath = `${this.indexPath}.tmp.${Date.now()}`;
await writeFile(tempPath, JSON.stringify(index, null, 2));
await rename(tempPath, this.indexPath);
}
private async addToIndex(message: Message): Promise<void> {
const index = await this.readIndex();
// Add to recipient's inbox
const toKey = `${message.toType}:${message.toId}`;
if (!index.byOwner[toKey]) {
index.byOwner[toKey] = { inbox: [], outbox: [] };
}
index.byOwner[toKey].inbox.unshift(message.id);
// Add to sender's outbox
const fromKey = `${message.fromType}:${message.fromId}`;
if (!index.byOwner[fromKey]) {
index.byOwner[fromKey] = { inbox: [], outbox: [] };
}
index.byOwner[fromKey].outbox.unshift(message.id);
await this.writeIndex(index);
}
private async removeFromIndex(message: Message): Promise<void> {
const index = await this.readIndex();
// Remove from recipient's inbox
const toKey = `${message.toType}:${message.toId}`;
if (index.byOwner[toKey]) {
index.byOwner[toKey].inbox = index.byOwner[toKey].inbox.filter((id) => id !== message.id);
index.byOwner[toKey].outbox = index.byOwner[toKey].outbox.filter((id) => id !== message.id);
}
// Remove from sender's outbox
const fromKey = `${message.fromType}:${message.fromId}`;
if (index.byOwner[fromKey]) {
index.byOwner[fromKey].inbox = index.byOwner[fromKey].inbox.filter((id) => id !== message.id);
index.byOwner[fromKey].outbox = index.byOwner[fromKey].outbox.filter((id) => id !== message.id);
}
await this.writeIndex(index);
}
private async loadMessagesByIds(ids: string[]): Promise<Message[]> {
const messages: Message[] = [];
for (const id of ids) {
const message = await this.getMessage(id);
if (message) {
messages.push(message);
}
}
return messages;
}
private applyFilter(messages: Message[], filter?: MessageFilter): Message[] {
let result = messages;
if (filter?.type) {
result = result.filter((m) => m.type === filter.type);
}
if (filter?.read !== undefined) {
result = result.filter((m) => m.read === filter.read);
}
// Apply pagination
const offset = filter?.offset ?? 0;
const limit = filter?.limit ?? result.length;
result = result.slice(offset, offset + limit);
return result;
}
}