From 07368b4685126368fe02874588709d8e2479cd9b Mon Sep 17 00:00:00 2001 From: YeonGyu-Kim Date: Tue, 28 Apr 2026 10:46:22 +0900 Subject: [PATCH] feat(team-mode): add team mailbox send with tests --- .../team-mode/team-mailbox/send.test.ts | 189 ++++++++++++++++++ src/features/team-mode/team-mailbox/send.ts | 166 +++++++++++++++ 2 files changed, 355 insertions(+) create mode 100644 src/features/team-mode/team-mailbox/send.test.ts create mode 100644 src/features/team-mode/team-mailbox/send.ts diff --git a/src/features/team-mode/team-mailbox/send.test.ts b/src/features/team-mode/team-mailbox/send.test.ts new file mode 100644 index 000000000..7f556be7a --- /dev/null +++ b/src/features/team-mode/team-mailbox/send.test.ts @@ -0,0 +1,189 @@ +/// + +import { describe, expect, test } from "bun:test" +import { mkdir, mkdtemp, readdir, readFile, writeFile } from "node:fs/promises" +import { randomUUID } from "node:crypto" +import { tmpdir } from "node:os" +import path from "node:path" + +import { TeamModeConfigSchema } from "../../../config/schema/team-mode" +import { getInboxDir, resolveBaseDir } from "../team-registry/paths" +import { MessageSchema } from "../types" +import { + BroadcastNotPermittedError, + DuplicateMessageIdError, + PayloadTooLargeError, + RecipientBackpressureError, + sendMessage, +} from "./send" + +async function createBaseDirectory(): Promise { + return await mkdtemp(path.join(tmpdir(), "team-mailbox-send-")) +} + +function createConfig(baseDir: string) { + return TeamModeConfigSchema.parse({ base_dir: baseDir }) +} + +function createMessage(overrides?: Partial[0]>) { + return MessageSchema.parse({ + version: 1, + messageId: randomUUID(), + from: "lead", + to: "m1", + kind: "message", + body: "hello", + timestamp: Date.now(), + ...overrides, + }) +} + +describe("sendMessage", () => { + test("writes distinct files for concurrent writers targeting the same recipient", async () => { + // given + const baseDir = await createBaseDirectory() + const config = createConfig(baseDir) + const teamRunId = randomUUID() + const messages = Array.from({ length: 4 }, (_, index) => createMessage({ + from: `m${index + 1}`, + body: `message-${index + 1}`, + timestamp: 100 + index, + })) + + // when + await Promise.all(messages.map(async (message) => { + await sendMessage(message, teamRunId, config, { isLead: false, activeMembers: ["m1"] }) + })) + + // then + const inboxDir = getInboxDir(resolveBaseDir(config), teamRunId, "m1") + const fileNames = (await readdir(inboxDir)).filter((entry) => entry.endsWith(".json")) + expect(fileNames).toHaveLength(4) + + const parsedMessages = await Promise.all(fileNames.map(async (fileName) => { + const fileContent = await readFile(path.join(inboxDir, fileName), "utf8") + return MessageSchema.parse(JSON.parse(fileContent)) + })) + expect(new Set(parsedMessages.map((message) => message.messageId)).size).toBe(4) + }) + + test("rejects payloads larger than 32 KB", async () => { + // given + const config = createConfig(await createBaseDirectory()) + const message = createMessage({ body: "가".repeat(20_000) }) + + // when + const result = sendMessage(message, randomUUID(), config, { isLead: false, activeMembers: ["m1"] }) + + // then + try { + await result + throw new Error("expected sendMessage to reject") + } catch (error) { + expect(error).toBeInstanceOf(PayloadTooLargeError) + } + }) + + test("rejects sends when recipient unread bytes exceed the backpressure limit", async () => { + // given + const baseDir = await createBaseDirectory() + const config = createConfig(baseDir) + const teamRunId = randomUUID() + const inboxDir = getInboxDir(resolveBaseDir(config), teamRunId, "m1") + await mkdir(inboxDir, { recursive: true }) + await writeFile(path.join(inboxDir, "full.json"), "x".repeat(config.recipient_unread_max_bytes + 1), { flag: "w" }) + + // when + const result = sendMessage(createMessage(), teamRunId, config, { isLead: false, activeMembers: ["m1"] }) + + // then + try { + await result + throw new Error("expected sendMessage to reject") + } catch (error) { + expect(error).toBeInstanceOf(RecipientBackpressureError) + } + }) + + test("counts in-flight .delivering-* reservations toward recipient backpressure", async () => { + // given + const baseDir = await createBaseDirectory() + const config = createConfig(baseDir) + const teamRunId = randomUUID() + const inboxDir = getInboxDir(resolveBaseDir(config), teamRunId, "m1") + await mkdir(inboxDir, { recursive: true }) + const pendingMessageId = randomUUID() + await writeFile( + path.join(inboxDir, `.delivering-${pendingMessageId}.json`), + "x".repeat(config.recipient_unread_max_bytes + 1), + { flag: "w" }, + ) + + // when + const result = sendMessage(createMessage(), teamRunId, config, { isLead: false, activeMembers: ["m1"] }) + + // then + try { + await result + throw new Error("expected sendMessage to reject") + } catch (error) { + expect(error).toBeInstanceOf(RecipientBackpressureError) + } + }) + + test("rejects duplicate message ids for the same recipient", async () => { + // given + const baseDir = await createBaseDirectory() + const config = createConfig(baseDir) + const teamRunId = randomUUID() + const message = createMessage() + await sendMessage(message, teamRunId, config, { isLead: false, activeMembers: ["m1"] }) + + // when + const result = sendMessage(message, teamRunId, config, { isLead: false, activeMembers: ["m1"] }) + + // then + try { + await result + throw new Error("expected sendMessage to reject") + } catch (error) { + expect(error).toBeInstanceOf(DuplicateMessageIdError) + } + }) + + test("gates broadcasts to leads and fans out to each active member", async () => { + // given + const baseDir = await createBaseDirectory() + const config = createConfig(baseDir) + const teamRunId = randomUUID() + const broadcastMessage = createMessage({ to: "*" }) + + // when + const rejectedSend = sendMessage(broadcastMessage, teamRunId, config, { + isLead: false, + activeMembers: ["m1", "m2"], + }) + const deliveredSend = sendMessage(broadcastMessage, teamRunId, config, { + isLead: true, + activeMembers: ["m1", "m2"], + }) + + // then + try { + await rejectedSend + throw new Error("expected sendMessage to reject") + } catch (error) { + expect(error).toBeInstanceOf(BroadcastNotPermittedError) + } + + expect(await deliveredSend).toEqual({ + messageId: broadcastMessage.messageId, + deliveredTo: ["m1", "m2"], + }) + + const memberOneFiles = await readdir(getInboxDir(resolveBaseDir(config), teamRunId, "m1")) + const memberTwoFiles = await readdir(getInboxDir(resolveBaseDir(config), teamRunId, "m2")) + expect(memberOneFiles.filter((entry) => entry.endsWith(".json"))).toHaveLength(1) + expect(memberTwoFiles.filter((entry) => entry.endsWith(".json"))).toHaveLength(1) + }) +}) diff --git a/src/features/team-mode/team-mailbox/send.ts b/src/features/team-mode/team-mailbox/send.ts new file mode 100644 index 000000000..a5d055672 --- /dev/null +++ b/src/features/team-mode/team-mailbox/send.ts @@ -0,0 +1,166 @@ +import { Buffer } from "node:buffer" +import { mkdir, readdir, stat } from "node:fs/promises" +import path from "node:path" + +import type { TeamModeConfig } from "../../../config/schema/team-mode" +import { getInboxDir, resolveBaseDir } from "../team-registry/paths" +import { loadRuntimeState } from "../team-state-store/store" +import { atomicWrite, withLock } from "../team-state-store/locks" +import type { Message } from "../types" + +type SendContext = { + isLead: boolean + activeMembers: string[] + reservedRecipients?: ReadonlySet +} + +export class BroadcastNotPermittedError extends Error { + constructor(message = "broadcast requires lead role") { + super(message) + this.name = "BroadcastNotPermittedError" + } +} + +export class PayloadTooLargeError extends Error { + constructor(message = "payload exceeds 32 KB") { + super(message) + this.name = "PayloadTooLargeError" + } +} + +export class RecipientBackpressureError extends Error { + constructor(message = "recipient inbox full (backpressure)") { + super(message) + this.name = "RecipientBackpressureError" + } +} + +export class DuplicateMessageIdError extends Error { + constructor(message = "duplicate message id") { + super(message) + this.name = "DuplicateMessageIdError" + } +} + +export class TeamDeletingError extends Error { + constructor(message = "team is deleting") { + super(message) + this.name = "TeamDeletingError" + } +} + +function isMissingPathError(error: unknown): boolean { + return typeof error === "object" + && error !== null + && "code" in error + && error.code === "ENOENT" +} + +async function assertTeamAcceptsMessages(teamRunId: string, config: TeamModeConfig): Promise { + try { + const runtimeState = await loadRuntimeState(teamRunId, config) + if (runtimeState.status === "deleting" || runtimeState.status === "deleted") { + throw new TeamDeletingError() + } + } catch (error) { + if (isMissingPathError(error)) { + return + } + + throw error + } +} + +function resolveRecipients(message: Message, context: SendContext): string[] { + if (message.to !== "*") { + return [message.to] + } + + return [...new Set(context.activeMembers)] +} + +async function getUnreadSizeBytes(inboxDir: string): Promise { + try { + const directoryEntries = await readdir(inboxDir, { withFileTypes: true }) + const unreadEntries = directoryEntries.filter((entry) => { + if (!entry.isFile() || !entry.name.endsWith(".json")) return false + if (entry.name.startsWith(".delivering-")) return true + return !entry.name.startsWith(".") + }) + + const sizes = await Promise.all(unreadEntries.map(async (entry) => { + const fileStats = await stat(path.join(inboxDir, entry.name)) + return fileStats.size + })) + + return sizes.reduce((totalBytes, fileSize) => totalBytes + fileSize, 0) + } catch (error) { + if (isMissingPathError(error)) { + return 0 + } + + throw error + } +} + +async function fileExists(filePath: string): Promise { + try { + await stat(filePath) + return true + } catch (error) { + if (isMissingPathError(error)) { + return false + } + + throw error + } +} + +export async function sendMessage( + message: Message, + teamRunId: string, + config: TeamModeConfig, + context: SendContext, +): Promise<{ messageId: string; deliveredTo: string[] }> { + const serializedMessage = `${JSON.stringify(message, null, 2)}\n` + const serializedMessageBytes = Buffer.byteLength(serializedMessage, "utf8") + const payloadBytes = Buffer.byteLength(message.body, "utf8") + if (payloadBytes > config.message_payload_max_bytes) { + throw new PayloadTooLargeError() + } + + await assertTeamAcceptsMessages(teamRunId, config) + + if (message.to === "*" && !context.isLead) { + throw new BroadcastNotPermittedError() + } + + const baseDir = resolveBaseDir(config) + const deliveredTo: string[] = [] + const reservedRecipients = context.reservedRecipients ?? new Set() + + for (const recipient of resolveRecipients(message, context)) { + const inboxDir = getInboxDir(baseDir, teamRunId, recipient) + await mkdir(inboxDir, { recursive: true, mode: 0o700 }) + + await withLock(`${inboxDir}.lock`, async () => { + const unreadSizeBytes = await getUnreadSizeBytes(inboxDir) + const nextUnreadSizeBytes = unreadSizeBytes + serializedMessageBytes + if (nextUnreadSizeBytes > config.recipient_unread_max_bytes) { + throw new RecipientBackpressureError() + } + + const unreservedPath = path.join(inboxDir, `${message.messageId}.json`) + const reservedPath = path.join(inboxDir, `.delivering-${message.messageId}.json`) + if (await fileExists(unreservedPath) || await fileExists(reservedPath)) { + throw new DuplicateMessageIdError() + } + + const targetPath = reservedRecipients.has(recipient) ? reservedPath : unreservedPath + await atomicWrite(targetPath, serializedMessage) + deliveredTo.push(recipient) + }, { ownerTag: `team-mailbox:${recipient}` }) + } + + return { messageId: message.messageId, deliveredTo } +}