fix(team-mode): skip pending mailbox reinjection
This commit is contained in:
@@ -1,8 +1,8 @@
|
|||||||
/// <reference types="bun-types" />
|
/// <reference types="bun-types" />
|
||||||
|
|
||||||
import { afterEach, describe, expect, mock, test } from "bun:test"
|
import { afterEach, describe, expect, mock, test } from "bun:test"
|
||||||
import { readdir } from "node:fs/promises"
|
|
||||||
import { randomUUID } from "node:crypto"
|
import { randomUUID } from "node:crypto"
|
||||||
|
import { readdir } from "node:fs/promises"
|
||||||
import { tmpdir } from "node:os"
|
import { tmpdir } from "node:os"
|
||||||
import path from "node:path"
|
import path from "node:path"
|
||||||
|
|
||||||
@@ -145,7 +145,7 @@ describe("pollAndBuildInjection", () => {
|
|||||||
expect(inboxEntries).not.toContain("processed")
|
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
|
// given
|
||||||
const { teamRunId, config } = await setupRuntime(["m1"])
|
const { teamRunId, config } = await setupRuntime(["m1"])
|
||||||
const messageId = randomUUID()
|
const messageId = randomUUID()
|
||||||
@@ -160,12 +160,56 @@ describe("pollAndBuildInjection", () => {
|
|||||||
}, teamRunId, config, { isLead: true, activeMembers: ["m1"] })
|
}, teamRunId, config, { isLead: true, activeMembers: ["m1"] })
|
||||||
|
|
||||||
// when
|
// when
|
||||||
await pollAndBuildInjection("session-1", "m1", teamRunId, config, "turn-A")
|
const firstInjection = await pollAndBuildInjection("session-1", "m1", teamRunId, config, "turn-A")
|
||||||
await pollAndBuildInjection("session-1", "m1", teamRunId, config, "turn-B")
|
const secondInjection = await pollAndBuildInjection("session-1", "m1", teamRunId, config, "turn-B")
|
||||||
const runtimeState = await loadRuntimeState(teamRunId, config)
|
const runtimeState = await loadRuntimeState(teamRunId, config)
|
||||||
const member = runtimeState.members.find((entry) => entry.name === "m1")
|
const member = runtimeState.members.find((entry) => entry.name === "m1")
|
||||||
|
|
||||||
// then
|
// then
|
||||||
|
expect(firstInjection.injected).toBe(true)
|
||||||
|
expect(secondInjection).toEqual({
|
||||||
|
injected: false,
|
||||||
|
messageIds: [],
|
||||||
|
reason: "pending ack",
|
||||||
|
})
|
||||||
expect(member?.pendingInjectedMessageIds).toEqual([messageId])
|
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])
|
||||||
|
})
|
||||||
})
|
})
|
||||||
|
|||||||
@@ -1,5 +1,5 @@
|
|||||||
import type { TeamModeConfig } from "../../../config/schema/team-mode"
|
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 type { Message } from "../types"
|
||||||
import { listUnreadMessages } from "./inbox"
|
import { listUnreadMessages } from "./inbox"
|
||||||
|
|
||||||
@@ -58,8 +58,14 @@ export async function pollAndBuildInjection(
|
|||||||
return { injected: false, messageIds: [], reason: "already injected this turn" }
|
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 (unreadMessages.length === 0) {
|
||||||
|
if (pendingMessageIds.size > 0) {
|
||||||
|
return { injected: false, messageIds: [], reason: "pending ack" }
|
||||||
|
}
|
||||||
|
|
||||||
return { injected: false, messageIds: [], reason: "no unread" }
|
return { injected: false, messageIds: [], reason: "no unread" }
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -5,13 +5,13 @@ import { tmpdir } from "node:os"
|
|||||||
import path from "node:path"
|
import path from "node:path"
|
||||||
|
|
||||||
import { TeamModeConfigSchema } from "../../config/schema/team-mode"
|
import { TeamModeConfigSchema } from "../../config/schema/team-mode"
|
||||||
|
import { sendMessage } from "../../features/team-mode/team-mailbox/send"
|
||||||
import {
|
import {
|
||||||
clearTeamSessionRegistry,
|
clearTeamSessionRegistry,
|
||||||
registerTeamSession,
|
registerTeamSession,
|
||||||
} from "../../features/team-mode/team-session-registry"
|
} 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 { saveRuntimeState } from "../../features/team-mode/team-state-store/store"
|
||||||
|
import type { RuntimeState } from "../../features/team-mode/types"
|
||||||
import { createTeamMailboxInjector } from "./hook"
|
import { createTeamMailboxInjector } from "./hook"
|
||||||
|
|
||||||
function createRuntimeState(sessionID: string, teamRunId = randomUUID()): RuntimeState {
|
function createRuntimeState(sessionID: string, teamRunId = randomUUID()): RuntimeState {
|
||||||
@@ -182,6 +182,48 @@ describe("createTeamMailboxInjector", () => {
|
|||||||
expect(secondOutput.messages).toEqual(originalSecondMessages)
|
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 () => {
|
it("injects mailbox messages during the spawn race when the registry has the fresh member session but disk state is stale", async () => {
|
||||||
// given
|
// given
|
||||||
const baseDir = await createTemporaryBaseDir()
|
const baseDir = await createTemporaryBaseDir()
|
||||||
|
|||||||
Reference in New Issue
Block a user