feat(FN-3905): use merger session model for post-merge workflow steps

The merge implements post-merge model resolution, using the merger's own session model for executing post-merge prompt workflow steps instead of falling back to executor model defaults. Tests cover the resolution hierarchy, and documentation clarifies the precedence order for downstream consumers.

Fusion-Task-Id: FN-3905
This commit is contained in:
Fusion
2026-05-11 09:19:57 -07:00
committed by gsxdsm
parent 758073c799
commit 3c2f1bdff0
13 changed files with 379 additions and 46 deletions

View File

@@ -5,6 +5,7 @@ import { join } from "node:path";
import { afterEach, beforeEach, describe, expect, it } from "vitest";
import { ChatStore, Database } from "@fusion/core";
import { request } from "../test-request.js";
import { RoomReplyGenerationError } from "../chat.js";
class MockStore {
constructor(private readonly rootDir: string, private readonly db: Database) {}
@@ -190,6 +191,62 @@ describe("Chat Room API Routes", () => {
expect(emptyContent.status).toBe(400);
});
it("surfaces room responder failures instead of returning silent success", async () => {
const { createServer } = await import("../server.js");
const appWithFailingRoomReplies = createServer(store as any, {
chatStore,
chatManager: {
sendRoomMessage: async () => {
throw new Error("Failed to generate room replies for room room-1: agent-a: Room responder returned an empty reply");
},
} as any,
});
const createRoomRes = await request(appWithFailingRoomReplies, "POST", "/api/chat/rooms", JSON.stringify({ name: "Product" }), {
"content-type": "application/json",
});
const roomId = (createRoomRes.body as any).room.id as string;
const postRes = await request(
appWithFailingRoomReplies,
"POST",
`/api/chat/rooms/${roomId}/messages`,
JSON.stringify({ content: "hello world" }),
{ "content-type": "application/json" },
);
expect(postRes.status).toBe(500);
expect(JSON.stringify(postRes.body)).toContain("Failed to generate room replies");
});
it("surfaces deterministic room-reply generation failures", async () => {
const { createServer } = await import("../server.js");
const failingApp = createServer(store as any, {
chatStore,
chatManager: {
sendRoomMessage: async (_roomId: string) => {
throw new RoomReplyGenerationError("No active room responders available", _roomId);
},
} as any,
});
const createRoomRes = await request(failingApp, "POST", "/api/chat/rooms", JSON.stringify({ name: "Product" }), {
"content-type": "application/json",
});
const roomId = (createRoomRes.body as any).room.id as string;
const postRes = await request(
failingApp,
"POST",
`/api/chat/rooms/${roomId}/messages`,
JSON.stringify({ content: "hello world" }),
{ "content-type": "application/json" },
);
expect(postRes.status).toBe(502);
expect((postRes.body as any).error).toContain("No active room responders available");
});
it("deletes messages and supports pagination", async () => {
const room = chatStore.createRoom({ name: "QA" });
const m1 = chatStore.addRoomMessage(room.id, { role: "user", content: "one" });

View File

@@ -190,6 +190,10 @@ describe("Chat HTTP + SSE routes — rooms (FN-3805..FN-3811 contract)", () => {
const persisted = chatStore.getRoomMessage(messageId);
expect(persisted?.content).toBe("hello @agent_room");
const assistantMessages = chatStore.getRoomMessages(roomId).filter((entry) => entry.role === "assistant");
expect(assistantMessages).toHaveLength(1);
expect(assistantMessages[0]).toMatchObject({ senderAgentId: "agent-room" });
const invalidSender = await request(
appWithRoomReplies,
"POST",

View File

@@ -1,5 +1,5 @@
import { beforeEach, describe, expect, it, vi } from "vitest";
import { ChatManager, __setCreateResolvedAgentSession, __resetChatState } from "../chat.js";
import { ChatManager, RoomReplyGenerationError, __setCreateResolvedAgentSession, __resetChatState } from "../chat.js";
const mockChatStore = {
listRoomMembers: vi.fn(),
@@ -89,6 +89,20 @@ describe("Chat orchestration — rooms (FN-3805..FN-3811 contract)", () => {
expect(assistantWrite).toMatchObject({ role: "assistant", senderAgentId: "agent-a", content: "Room reply" });
});
it("fails deterministically when no member responder can be resolved", async () => {
mockChatStore.listRoomMembers.mockReturnValue([
{ roomId: "room-1", agentId: "agent-a", role: "member", addedAt: "2026-01-01" },
]);
mockAgentStore.listAgents.mockRejectedValue(new Error("list failed"));
mockAgentStore.getAgent.mockResolvedValue(null);
const manager = new ChatManager(mockChatStore as any, "/tmp", mockAgentStore as any);
await expect(manager.sendRoomMessage("room-1", "hello")).rejects.toBeInstanceOf(RoomReplyGenerationError);
expect(mockChatStore.addRoomMessage).toHaveBeenCalledTimes(1);
expect(mockChatStore.addRoomMessage.mock.calls[0]?.[1]).toMatchObject({ role: "user", content: "hello" });
});
it("falls back to room-member getAgent lookup when listAgents fails", async () => {
mockChatStore.listRoomMembers.mockReturnValue([
{ roomId: "room-1", agentId: "agent-a", role: "member", addedAt: "2026-01-01" },
@@ -235,5 +249,46 @@ describe("Chat orchestration — rooms (FN-3805..FN-3811 contract)", () => {
expect(prompt).toContain("history-item-28");
expect(prompt).not.toContain("history-item-0");
});
it("throws surfaced error when room has members but no resolvable responders", async () => {
mockChatStore.listRoomMembers.mockReturnValue([
{ roomId: "room-1", agentId: "agent-a", role: "member", addedAt: "2026-01-01" },
]);
mockAgentStore.listAgents.mockResolvedValue([]);
mockAgentStore.getAgent.mockResolvedValue(null);
const manager = new ChatManager(mockChatStore as any, "/tmp", mockAgentStore as any);
await expect(manager.sendRoomMessage("room-1", "hello")).rejects.toThrow(
"No active room responders available for room room-1",
);
});
it("throws surfaced error when all room responders fail to reply", async () => {
mockChatStore.listRoomMembers.mockReturnValue([
{ roomId: "room-1", agentId: "agent-a", role: "member", addedAt: "2026-01-01" },
]);
mockAgentStore.listAgents.mockResolvedValue([{ id: "agent-a", name: "Alpha", role: "executor" }]);
mockAgentStore.getAgent.mockResolvedValue({ id: "agent-a", name: "Alpha", role: "executor" });
__setCreateResolvedAgentSession(async () => ({
session: {
prompt: vi.fn().mockResolvedValue(undefined),
dispose: vi.fn(),
state: {
messages: [{ role: "assistant", content: " " }],
},
},
} as any));
const manager = new ChatManager(mockChatStore as any, "/tmp", mockAgentStore as any);
await expect(manager.sendRoomMessage("room-1", "hello @Alpha")).rejects.toThrow(
/Failed to generate room replies for room room-1:/,
);
const assistantWrites = mockChatStore.addRoomMessage.mock.calls
.map((call: any[]) => call[1])
.filter((input: any) => input.role === "assistant" && input.senderAgentId);
expect(assistantWrites).toHaveLength(0);
});
});
});

View File

@@ -45,6 +45,13 @@ vi.mock("@fusion/engine", () => ({
dispose: vi.fn(),
},
})),
createResolvedAgentSession: vi.fn(async () => ({
session: {
state: { messages: [] as Array<{ role: string; content: string }> },
prompt: vi.fn(),
dispose: vi.fn(),
},
})),
AgentReflectionService: class {
async generateReflection(): Promise<never> { throw new Error("Reflection service unavailable"); }
async buildReflectionContext(): Promise<never> { throw new Error("Reflection service unavailable"); }

View File

@@ -491,6 +491,16 @@ export function getRateLimitResetTime(ip: string): Date | null {
* Manages AI agent chat sessions.
* Creates sessions, sends messages, and streams AI responses via SSE.
*/
export class RoomReplyGenerationError extends Error {
readonly roomId: string;
constructor(message: string, roomId: string) {
super(message);
this.name = "RoomReplyGenerationError";
this.roomId = roomId;
}
}
export class ChatManager {
private agentStoreReady?: Promise<void>;
private generationCounter = 0;
@@ -885,6 +895,7 @@ export class ChatManager {
...(Array.isArray(attachments) ? { attachments } : {}),
});
const roomMembers = this.chatStore.listRoomMembers(roomId);
const responders = [...responderPlan.direct, ...responderPlan.ambient];
if (responders.length === 0) {
if (responderPlan.nonMemberMentions.length > 0) {
@@ -897,29 +908,51 @@ export class ChatManager {
content: `I couldn't route ${labels} because they are not members of this room.`,
});
}
if (roomMembers.length > 0) {
throw new RoomReplyGenerationError(`No active room responders available for room ${roomId}`, roomId);
}
return { userMessage, responders: [] };
}
for (const responder of responders) {
const response = await this.generateRoomResponderReply({
roomId,
roomName: room.name,
content: trimmedContent,
latestUserMessageId: userMessage.id,
mentions,
responder,
modelProvider,
modelId,
});
const successfulResponderIds: string[] = [];
const responderFailures: string[] = [];
this.chatStore.addRoomMessage(roomId, {
role: "assistant",
content: response.content,
thinkingOutput: response.thinkingOutput,
metadata: response.metadata,
senderAgentId: responder.id,
mentions: mentions.map((mention) => mention.agentId),
});
for (const responder of responders) {
try {
const response = await this.generateRoomResponderReply({
roomId,
roomName: room.name,
content: trimmedContent,
latestUserMessageId: userMessage.id,
mentions,
responder,
modelProvider,
modelId,
});
this.chatStore.addRoomMessage(roomId, {
role: "assistant",
content: response.content,
thinkingOutput: response.thinkingOutput,
metadata: response.metadata,
senderAgentId: responder.id,
mentions: mentions.map((mention) => mention.agentId),
});
successfulResponderIds.push(responder.id);
} catch (error) {
const reason = error instanceof Error ? error.message : String(error);
diagnostics.error(`Room responder ${responder.id} failed in room ${roomId}: ${reason}`);
responderFailures.push(`${responder.id}: ${reason}`);
}
}
if (successfulResponderIds.length === 0) {
throw new RoomReplyGenerationError(
`Failed to generate room replies for room ${roomId}: ${responderFailures.join("; ")}`,
roomId,
);
}
if (responderPlan.nonMemberMentions.length > 0) {
@@ -935,7 +968,7 @@ export class ChatManager {
return {
userMessage,
responders: responders.map((responder) => responder.id),
responders: successfulResponderIds,
};
}
@@ -1011,8 +1044,18 @@ export class ChatManager {
.join("");
}
const stateError = (resolvedSession.session.state as { errorMessage?: string } | undefined)?.errorMessage;
if (stateError?.trim()) {
throw new Error(stateError.trim());
}
const finalContent = content.trim();
if (!finalContent) {
throw new Error("Room responder returned an empty reply");
}
return {
content: content.trim() || "(no response)",
content: finalContent,
thinkingOutput: null,
metadata: {
roomId: input.roomId,

View File

@@ -29,6 +29,7 @@ vi.mock("@fusion/core", async () => {
vi.mock("@fusion/engine", () => ({
createFnAgent: vi.fn(async () => ({ session: { state: { messages: [] }, prompt: vi.fn(), dispose: vi.fn() } })),
createResolvedAgentSession: vi.fn(async () => ({ session: { state: { messages: [] }, prompt: vi.fn(), dispose: vi.fn() } })),
promptWithFallback: vi.fn(),
}));

View File

@@ -1,4 +1,5 @@
import type { ChatAttachment, ChatRoomCreateInput, ChatRoomStatus, ChatRoomUpdateInput } from "@fusion/core";
import { RoomReplyGenerationError } from "../chat.js";
import { ApiError, badRequest, internalError, notFound } from "../api-error.js";
import { rateLimit, RATE_LIMITS } from "../rate-limit.js";
import type { ApiRoutesContext } from "./types.js";
@@ -265,6 +266,9 @@ export function registerChatRoomRoutes(ctx: ApiRoutesContext): void {
res.status(201).json({ message: result.userMessage });
} catch (err: unknown) {
if (err instanceof ApiError) throw err;
if (err instanceof RoomReplyGenerationError) {
throw new ApiError(502, err.message);
}
rethrowAsApiError(err, "Failed to create chat room message");
}
});

View File

@@ -5196,6 +5196,159 @@ describe("aiMergeTask — post-merge workflow steps", () => {
expect(store.moveTask).toHaveBeenCalledWith("FN-050", "done");
});
it("uses assigned agent runtime model for post-merge prompt step when workflow step has no override", async () => {
const store = createMockStore();
(store as any).getWorkflowStep = vi.fn().mockResolvedValue({
id: "WS-001",
name: "Post-merge Notify",
description: "Send notifications after merge",
prompt: "Check merged code.",
phase: "post-merge",
mode: "prompt",
enabled: true,
createdAt: new Date().toISOString(),
updatedAt: new Date().toISOString(),
});
const baseTask = {
id: "FN-050",
title: "Test task",
description: "Test",
column: "in-review",
dependencies: [],
worktree: "/tmp/root/.worktrees/KB-050",
steps: [],
currentStep: 0,
log: [],
assignedAgentId: "agent-001",
enabledWorkflowSteps: ["WS-001"],
createdAt: new Date().toISOString(),
updatedAt: new Date().toISOString(),
};
store.getTask = vi.fn().mockResolvedValue({ ...baseTask, prompt: "# test" });
await aiMergeTask(store, "/tmp/root", "FN-050", {
agentStore: {
listAgents: vi.fn().mockResolvedValue([]),
getAgent: vi.fn().mockResolvedValue({
id: "agent-001",
runtimeConfig: {
model: "anthropic/claude-3-5-sonnet-20241022",
},
}),
} as any,
});
const postMergeAgentCall = mockedCreateFnAgent.mock.calls.find(
(c: any) => c[0]?.systemPrompt?.includes("post-merge workflow step agent"),
);
expect(postMergeAgentCall?.[0]?.defaultProvider).toBe("anthropic");
expect(postMergeAgentCall?.[0]?.defaultModelId).toBe("claude-3-5-sonnet-20241022");
});
it("uses workflow-step model override over assigned agent runtime model", async () => {
const store = createMockStore();
(store as any).getWorkflowStep = vi.fn().mockResolvedValue({
id: "WS-001",
name: "Post-merge Notify",
description: "Send notifications after merge",
prompt: "Check merged code.",
phase: "post-merge",
mode: "prompt",
modelProvider: "openai",
modelId: "gpt-4.1-mini",
enabled: true,
createdAt: new Date().toISOString(),
updatedAt: new Date().toISOString(),
});
const baseTask = {
id: "FN-050",
title: "Test task",
description: "Test",
column: "in-review",
dependencies: [],
worktree: "/tmp/root/.worktrees/KB-050",
steps: [],
currentStep: 0,
log: [],
assignedAgentId: "agent-001",
enabledWorkflowSteps: ["WS-001"],
createdAt: new Date().toISOString(),
updatedAt: new Date().toISOString(),
};
store.getTask = vi.fn().mockResolvedValue({ ...baseTask, prompt: "# test" });
await aiMergeTask(store, "/tmp/root", "FN-050", {
agentStore: {
listAgents: vi.fn().mockResolvedValue([]),
getAgent: vi.fn().mockResolvedValue({
id: "agent-001",
runtimeConfig: {
model: "anthropic/claude-3-5-sonnet-20241022",
},
}),
} as any,
});
const postMergeAgentCall = mockedCreateFnAgent.mock.calls.find(
(c: any) => c[0]?.systemPrompt?.includes("post-merge workflow step agent"),
);
expect(postMergeAgentCall?.[0]?.defaultProvider).toBe("openai");
expect(postMergeAgentCall?.[0]?.defaultModelId).toBe("gpt-4.1-mini");
const modelLogCall = (store.logEntry as ReturnType<typeof vi.fn>).mock.calls.find(
(call: any) => String(call[1]).includes("Workflow step 'Post-merge Notify' using model:"),
);
expect(modelLogCall?.[1]).toContain("(workflow step override)");
});
it("falls back to project default override model when no workflow-step or assigned-agent model is set", async () => {
const store = createMockStore();
(store as any).getWorkflowStep = vi.fn().mockResolvedValue({
id: "WS-001",
name: "Post-merge Notify",
description: "Send notifications after merge",
prompt: "Check merged code.",
phase: "post-merge",
mode: "prompt",
enabled: true,
createdAt: new Date().toISOString(),
updatedAt: new Date().toISOString(),
});
const baseTask = {
id: "FN-050",
title: "Test task",
description: "Test",
column: "in-review",
dependencies: [],
worktree: "/tmp/root/.worktrees/KB-050",
steps: [],
currentStep: 0,
log: [],
enabledWorkflowSteps: ["WS-001"],
createdAt: new Date().toISOString(),
updatedAt: new Date().toISOString(),
};
store.getTask = vi.fn().mockResolvedValue({ ...baseTask, prompt: "# test" });
(store.getSettings as ReturnType<typeof vi.fn>).mockResolvedValue({
...DEFAULT_SETTINGS,
defaultProviderOverride: "openai",
defaultModelIdOverride: "gpt-4o-mini",
defaultProvider: "anthropic",
defaultModelId: "claude-3-5-haiku-latest",
});
await aiMergeTask(store, "/tmp/root", "FN-050");
const postMergeAgentCall = mockedCreateFnAgent.mock.calls.find(
(c: any) => c[0]?.systemPrompt?.includes("post-merge workflow step agent"),
);
expect(postMergeAgentCall?.[0]?.defaultProvider).toBe("openai");
expect(postMergeAgentCall?.[0]?.defaultModelId).toBe("gpt-4o-mini");
});
it("does not run pre-merge workflow steps in merger", async () => {
const store = createMockStore();
(store as any).getWorkflowStep = vi.fn().mockResolvedValue({

View File

@@ -36,7 +36,6 @@ import {
getTaskMergeBlocker,
normalizeMergeConflictStrategy,
resolveTaskMergeTarget,
resolveProjectDefaultModel,
resolveTitleSummarizerSettingsModel,
resolveAgentPrompt,
summarizeCommitBody,
@@ -7741,28 +7740,6 @@ If issues are found that need attention, describe them clearly and include concr
});
try {
const defaultModel = resolveProjectDefaultModel(settings);
const stepProvider = workflowStep.modelProvider || defaultModel.provider;
const stepModelId = workflowStep.modelId || defaultModel.modelId;
const useOverride = !!(workflowStep.modelProvider && workflowStep.modelId);
// Post-merge step agents inherit merger instructions
let postMergeInstructions = "";
if (mergeOptions.agentStore) {
try {
const agents = await mergeOptions.agentStore.listAgents({ role: "merger" });
for (const agent of agents) {
if (agent.instructionsText || agent.instructionsPath) {
postMergeInstructions = await resolveAgentInstructions(agent, rootDir);
break;
}
}
} catch {
// Graceful fallback
}
}
const postMergeSystemPrompt = buildSystemPromptWithInstructions(systemPrompt, postMergeInstructions);
// Build skill selection context for post-merge session
let postMergeSkillContext = undefined;
let taskForSkillContext: Awaited<ReturnType<typeof store.getTask>> | null = null;
@@ -7788,6 +7765,28 @@ If issues are found that need attention, describe them clearly and include concr
const assignedAgent = assignedAgentId && agentStoreWithGetAgent
? await agentStoreWithGetAgent.getAgent(assignedAgentId).catch(() => null)
: null;
const mergerSessionModel = resolveMergerSessionModel(settings, assignedAgent?.runtimeConfig);
const stepProvider = workflowStep.modelProvider || mergerSessionModel.provider;
const stepModelId = workflowStep.modelId || mergerSessionModel.modelId;
const useOverride = !!(workflowStep.modelProvider && workflowStep.modelId);
// Post-merge step agents inherit merger instructions
let postMergeInstructions = "";
if (mergeOptions.agentStore) {
try {
const agents = await mergeOptions.agentStore.listAgents({ role: "merger" });
for (const agent of agents) {
if (agent.instructionsText || agent.instructionsPath) {
postMergeInstructions = await resolveAgentInstructions(agent, rootDir);
break;
}
}
} catch {
// Graceful fallback
}
}
const postMergeSystemPrompt = buildSystemPromptWithInstructions(systemPrompt, postMergeInstructions);
const mergerRuntimeHint = extractRuntimeHint(assignedAgent?.runtimeConfig);
const { session } = await createResolvedAgentSession({
sessionPurpose: "merger",