feat(team-mode): add team mailbox inbox with tests
This commit is contained in:
@@ -0,0 +1,63 @@
|
|||||||
|
/// <reference types="bun-types" />
|
||||||
|
|
||||||
|
import { describe, expect, mock, test } from "bun:test"
|
||||||
|
import { mkdir, mkdtemp, writeFile } from "node:fs/promises"
|
||||||
|
import { randomUUID } from "node:crypto"
|
||||||
|
import { tmpdir } from "node:os"
|
||||||
|
import path from "node:path"
|
||||||
|
|
||||||
|
const logCalls: Array<[string, unknown?]> = []
|
||||||
|
|
||||||
|
mock.module("../../../shared/logger", () => ({
|
||||||
|
log: (message: string, data?: unknown) => {
|
||||||
|
logCalls.push([message, data])
|
||||||
|
},
|
||||||
|
}))
|
||||||
|
|
||||||
|
const { listUnreadMessages } = await import("./inbox")
|
||||||
|
const { TeamModeConfigSchema } = await import("../../../config/schema/team-mode")
|
||||||
|
const { getInboxDir, resolveBaseDir } = await import("../team-registry/paths")
|
||||||
|
|
||||||
|
async function createBaseDirectory(): Promise<string> {
|
||||||
|
return await mkdtemp(path.join(tmpdir(), "team-mailbox-inbox-"))
|
||||||
|
}
|
||||||
|
|
||||||
|
describe("listUnreadMessages", () => {
|
||||||
|
test("returns FIFO messages while skipping malformed, processed, and dot files", async () => {
|
||||||
|
// given
|
||||||
|
const config = TeamModeConfigSchema.parse({ base_dir: await createBaseDirectory() })
|
||||||
|
const teamRunId = randomUUID()
|
||||||
|
const inboxDir = getInboxDir(resolveBaseDir(config), teamRunId, "m1")
|
||||||
|
await mkdir(path.join(inboxDir, "processed"), { recursive: true })
|
||||||
|
|
||||||
|
await writeFile(path.join(inboxDir, "later.json"), JSON.stringify({
|
||||||
|
version: 1,
|
||||||
|
messageId: randomUUID(),
|
||||||
|
from: "m2",
|
||||||
|
to: "m1",
|
||||||
|
kind: "message",
|
||||||
|
body: "later",
|
||||||
|
timestamp: 200,
|
||||||
|
}))
|
||||||
|
await writeFile(path.join(inboxDir, "earlier.json"), JSON.stringify({
|
||||||
|
version: 1,
|
||||||
|
messageId: randomUUID(),
|
||||||
|
from: "m3",
|
||||||
|
to: "m1",
|
||||||
|
kind: "message",
|
||||||
|
body: "earlier",
|
||||||
|
timestamp: 100,
|
||||||
|
}))
|
||||||
|
await writeFile(path.join(inboxDir, "bad.json"), "{not-json")
|
||||||
|
await writeFile(path.join(inboxDir, ".hidden.json"), "{}")
|
||||||
|
await writeFile(path.join(inboxDir, "processed", "done.json"), "{}")
|
||||||
|
|
||||||
|
// when
|
||||||
|
const unreadMessages = await listUnreadMessages(teamRunId, "m1", config)
|
||||||
|
|
||||||
|
// then
|
||||||
|
expect(unreadMessages.map((message) => message.body)).toEqual(["earlier", "later"])
|
||||||
|
expect(logCalls).toHaveLength(1)
|
||||||
|
expect(logCalls[0]?.[0]).toContain("skipped unreadable message")
|
||||||
|
})
|
||||||
|
})
|
||||||
@@ -0,0 +1,76 @@
|
|||||||
|
import type { Dirent } from "node:fs"
|
||||||
|
import { readdir, readFile } from "node:fs/promises"
|
||||||
|
import path from "node:path"
|
||||||
|
|
||||||
|
import type { TeamModeConfig } from "../../../config/schema/team-mode"
|
||||||
|
import { log } from "../../../shared/logger"
|
||||||
|
import { getInboxDir, resolveBaseDir } from "../team-registry/paths"
|
||||||
|
import { MessageSchema } from "../types"
|
||||||
|
import type { Message } from "../types"
|
||||||
|
|
||||||
|
function isInboxMessageFile(entry: Dirent): boolean {
|
||||||
|
return entry.isFile() && entry.name.endsWith(".json") && !entry.name.startsWith(".")
|
||||||
|
}
|
||||||
|
|
||||||
|
function isMissingDirectoryError(error: unknown): error is NodeJS.ErrnoException {
|
||||||
|
return error instanceof Error && "code" in error && error.code === "ENOENT"
|
||||||
|
}
|
||||||
|
|
||||||
|
async function readInboxMessage(
|
||||||
|
inboxDir: string,
|
||||||
|
fileName: string,
|
||||||
|
memberName: string,
|
||||||
|
teamRunId: string,
|
||||||
|
): Promise<Message | null> {
|
||||||
|
const filePath = path.join(inboxDir, fileName)
|
||||||
|
const messageContext = { memberName, teamRunId, fileName }
|
||||||
|
|
||||||
|
try {
|
||||||
|
const fileContent = await readFile(filePath, "utf8")
|
||||||
|
const parsedMessage = MessageSchema.safeParse(JSON.parse(fileContent))
|
||||||
|
if (!parsedMessage.success) {
|
||||||
|
log("team mailbox skipped malformed message", {
|
||||||
|
event: "team-mailbox-malformed-message",
|
||||||
|
...messageContext,
|
||||||
|
issues: parsedMessage.error.issues,
|
||||||
|
})
|
||||||
|
return null
|
||||||
|
}
|
||||||
|
|
||||||
|
return parsedMessage.data
|
||||||
|
} catch (error) {
|
||||||
|
log("team mailbox skipped unreadable message", {
|
||||||
|
event: "team-mailbox-unreadable-message",
|
||||||
|
...messageContext,
|
||||||
|
error: error instanceof Error ? error.message : String(error),
|
||||||
|
})
|
||||||
|
return null
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
export async function listUnreadMessages(
|
||||||
|
teamRunId: string,
|
||||||
|
memberName: string,
|
||||||
|
config: TeamModeConfig,
|
||||||
|
): Promise<Message[]> {
|
||||||
|
const inboxDir = getInboxDir(resolveBaseDir(config), teamRunId, memberName)
|
||||||
|
|
||||||
|
try {
|
||||||
|
const directoryEntries = await readdir(inboxDir, { withFileTypes: true })
|
||||||
|
const unreadMessages = await Promise.all(
|
||||||
|
directoryEntries
|
||||||
|
.filter(isInboxMessageFile)
|
||||||
|
.map((entry) => readInboxMessage(inboxDir, entry.name, memberName, teamRunId)),
|
||||||
|
)
|
||||||
|
|
||||||
|
return unreadMessages
|
||||||
|
.filter((message): message is Message => message !== null)
|
||||||
|
.sort((leftMessage, rightMessage) => leftMessage.timestamp - rightMessage.timestamp)
|
||||||
|
} catch (error) {
|
||||||
|
if (isMissingDirectoryError(error)) {
|
||||||
|
return []
|
||||||
|
}
|
||||||
|
|
||||||
|
throw error
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user