From e8de8b79a8fcaacb71284abebf491db9ae68c5d8 Mon Sep 17 00:00:00 2001 From: YeonGyu-Kim Date: Fri, 15 May 2026 23:15:46 +0900 Subject: [PATCH] fix(team-mode): skip pending mailbox reinjection --- .../team-mode/team-mailbox/poll.test.ts | 52 +++++++++++++++++-- src/features/team-mode/team-mailbox/poll.ts | 10 +++- src/hooks/team-mailbox-injector/hook.test.ts | 46 +++++++++++++++- 3 files changed, 100 insertions(+), 8 deletions(-) diff --git a/src/features/team-mode/team-mailbox/poll.test.ts b/src/features/team-mode/team-mailbox/poll.test.ts index efed5415c..424c60eaa 100644 --- a/src/features/team-mode/team-mailbox/poll.test.ts +++ b/src/features/team-mode/team-mailbox/poll.test.ts @@ -1,8 +1,8 @@ /// import { afterEach, describe, expect, mock, test } from "bun:test" -import { readdir } from "node:fs/promises" import { randomUUID } from "node:crypto" +import { readdir } from "node:fs/promises" import { tmpdir } from "node:os" import path from "node:path" @@ -145,7 +145,7 @@ describe("pollAndBuildInjection", () => { expect(inboxEntries).not.toContain("processed") }) - test("deduplicates pendingInjectedMessageIds when the same unread message surfaces across turns", async () => { + test("does not re-inject a pending message on a later turn", async () => { // given const { teamRunId, config } = await setupRuntime(["m1"]) const messageId = randomUUID() @@ -160,12 +160,56 @@ describe("pollAndBuildInjection", () => { }, teamRunId, config, { isLead: true, activeMembers: ["m1"] }) // when - await pollAndBuildInjection("session-1", "m1", teamRunId, config, "turn-A") - await pollAndBuildInjection("session-1", "m1", teamRunId, config, "turn-B") + const firstInjection = await pollAndBuildInjection("session-1", "m1", teamRunId, config, "turn-A") + const secondInjection = await pollAndBuildInjection("session-1", "m1", teamRunId, config, "turn-B") const runtimeState = await loadRuntimeState(teamRunId, config) const member = runtimeState.members.find((entry) => entry.name === "m1") // then + expect(firstInjection.injected).toBe(true) + expect(secondInjection).toEqual({ + injected: false, + messageIds: [], + reason: "pending ack", + }) expect(member?.pendingInjectedMessageIds).toEqual([messageId]) }) + + test("injects only new unread messages when older unread messages are pending ack", async () => { + // given + const { teamRunId, config } = await setupRuntime(["m1"]) + const pendingMessageId = randomUUID() + const newMessageId = randomUUID() + await sendMessage({ + version: 1, + messageId: pendingMessageId, + from: "lead", + to: "m1", + kind: "message", + body: "already injected", + timestamp: 100, + }, teamRunId, config, { isLead: true, activeMembers: ["m1"] }) + await pollAndBuildInjection("session-1", "m1", teamRunId, config, "turn-A") + await sendMessage({ + version: 1, + messageId: newMessageId, + from: "lead", + to: "m1", + kind: "message", + body: "fresh message", + timestamp: 200, + }, teamRunId, config, { isLead: true, activeMembers: ["m1"] }) + + // when + const result = await pollAndBuildInjection("session-1", "m1", teamRunId, config, "turn-B") + const runtimeState = await loadRuntimeState(teamRunId, config) + const member = runtimeState.members.find((entry) => entry.name === "m1") + + // then + expect(result.injected).toBe(true) + expect(result.messageIds).toEqual([newMessageId]) + expect(result.content).toContain("fresh message") + expect(result.content).not.toContain("already injected") + expect(member?.pendingInjectedMessageIds).toEqual([pendingMessageId, newMessageId]) + }) }) diff --git a/src/features/team-mode/team-mailbox/poll.ts b/src/features/team-mode/team-mailbox/poll.ts index 853ba0de5..2e8dcd8f4 100644 --- a/src/features/team-mode/team-mailbox/poll.ts +++ b/src/features/team-mode/team-mailbox/poll.ts @@ -1,5 +1,5 @@ import type { TeamModeConfig } from "../../../config/schema/team-mode" -import { transitionRuntimeState, loadRuntimeState } from "../team-state-store/store" +import { loadRuntimeState, transitionRuntimeState } from "../team-state-store/store" import type { Message } from "../types" import { listUnreadMessages } from "./inbox" @@ -58,8 +58,14 @@ export async function pollAndBuildInjection( return { injected: false, messageIds: [], reason: "already injected this turn" } } - const unreadMessages = await listUnreadMessages(teamRunId, memberName, config) + const pendingMessageIds = new Set(runtimeMember.pendingInjectedMessageIds) + const unreadMessages = (await listUnreadMessages(teamRunId, memberName, config)) + .filter((message) => !pendingMessageIds.has(message.messageId)) if (unreadMessages.length === 0) { + if (pendingMessageIds.size > 0) { + return { injected: false, messageIds: [], reason: "pending ack" } + } + return { injected: false, messageIds: [], reason: "no unread" } } diff --git a/src/hooks/team-mailbox-injector/hook.test.ts b/src/hooks/team-mailbox-injector/hook.test.ts index b99250f29..513f43dbd 100644 --- a/src/hooks/team-mailbox-injector/hook.test.ts +++ b/src/hooks/team-mailbox-injector/hook.test.ts @@ -5,13 +5,13 @@ import { tmpdir } from "node:os" import path from "node:path" import { TeamModeConfigSchema } from "../../config/schema/team-mode" +import { sendMessage } from "../../features/team-mode/team-mailbox/send" import { clearTeamSessionRegistry, registerTeamSession, } from "../../features/team-mode/team-session-registry" -import { sendMessage } from "../../features/team-mode/team-mailbox/send" -import type { RuntimeState } from "../../features/team-mode/types" import { saveRuntimeState } from "../../features/team-mode/team-state-store/store" +import type { RuntimeState } from "../../features/team-mode/types" import { createTeamMailboxInjector } from "./hook" function createRuntimeState(sessionID: string, teamRunId = randomUUID()): RuntimeState { @@ -182,6 +182,48 @@ describe("createTeamMailboxInjector", () => { expect(secondOutput.messages).toEqual(originalSecondMessages) }) + it("does not re-inject pending mailbox messages on a later turn marker", async () => { + // given + const baseDir = await createTemporaryBaseDir() + temporaryDirectories.push(baseDir) + const hook = createHook(baseDir) + const runtimeState = createRuntimeState("session-member") + await seedRuntimeState(baseDir, runtimeState) + await sendMessage({ + version: 1, + messageId: randomUUID(), + from: "lead", + to: "member-a", + kind: "message", + body: "hello", + timestamp: 1, + }, runtimeState.teamRunId, TeamModeConfigSchema.parse({ base_dir: baseDir, enabled: true }), { isLead: true, activeMembers: ["lead", "member-a"] }) + const firstOutput = createOutput("session-member") + const secondOutput = createOutput("session-member") + secondOutput.messages.unshift({ + info: { + role: "assistant", + sessionID: "session-member", + }, + parts: [{ type: "text", text: "assistant turn" }], + }) + const originalSecondMessages = structuredClone(secondOutput.messages) + + // when + await hook["experimental.chat.messages.transform"]?.( + { sessionID: "session-member" }, + firstOutput, + ) + await hook["experimental.chat.messages.transform"]?.( + { sessionID: "session-member" }, + secondOutput, + ) + + // then + expect(firstOutput.messages).toHaveLength(2) + expect(secondOutput.messages).toEqual(originalSecondMessages) + }) + it("injects mailbox messages during the spawn race when the registry has the fresh member session but disk state is stale", async () => { // given const baseDir = await createTemporaryBaseDir()