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

@@ -0,0 +1,5 @@
---
"@runfusion/fusion": patch
---
Post-merge prompt workflow-step agent sessions now honor the assigned agent runtime model (`runtimeConfig.model`) when the workflow step does not provide its own model override, matching the rest of the merger session model resolution path.

View File

@@ -280,8 +280,9 @@ Intentional exclusions from shared snapshots:
- The Rooms tab in `packages/dashboard/app/components/ChatView.tsx` is wired through `useChatRooms` (`packages/dashboard/app/hooks/useChatRooms.ts`).
- `useChatRooms` owns room list fetch/sort, active-room selection, member+message hydration, room creation/deletion, and room message sends.
- The hook subscribes to `/api/events` and consumes `chat:room:created`, `chat:room:updated`, `chat:room:deleted`, `chat:room:member:added`, `chat:room:member:removed`, `chat:room:message:added`, `chat:room:message:updated`, and `chat:room:message:deleted` to keep UI state in sync.
- Room messages persist through `POST /api/chat/rooms/:id/messages`; the route persists the user message first, then calls `ChatManager.sendRoomMessage(...)` to orchestrate room-member responders and persist assistant room replies with `chatStore.addRoomMessage(...)`.
- Room messages persist through `POST /api/chat/rooms/:id/messages`; the route persists the user message first, then calls `ChatManager.sendRoomMessage(...)` to orchestrate room-member responders and persist assistant room replies with `chatStore.addRoomMessage(...)` (including `senderAgentId` for each responder).
- `sendRoomMessage(...)` uses existing room-member + mention resolution rules: mentioned members are direct responders, non-mentioned members are ambient responders (capped by `ROOM_AMBIENT_MAX_RESPONDERS`), and non-member mentions are handled explicitly by the manager instead of silently disappearing.
- Room-reply generation is now non-silent on failure: if a room has members but no active responders can be resolved, or all responder generations fail/return empty output, `sendRoomMessage(...)` throws `RoomReplyGenerationError` and the route surfaces HTTP 502 instead of returning a silent user-only success.
- UI does not optimistically insert room messages; `useChatRooms.sendRoomMessage()` re-fetches authoritative room messages immediately after `POST /api/chat/rooms/:id/messages` so persisted assistant replies remain visible across SSE timing gaps, then continues applying `chat:room:message:*` events for live updates.
- Mention UI in rooms keeps direct-chat behavior unchanged while adding room affordances:
- `AgentMentionPopup` receives room membership context and shows members first with a `status-dot` member indicator (`aria-label="Room member"`).

View File

@@ -114,6 +114,8 @@ Chat Rooms are project-scoped group conversations for multiple agents. They are
- After a successful room send, the room composer is cleared (matching direct-chat composer behavior) so stale text is not left in the input.
- On mobile, room threads use the same keyboard-aware thread anchoring as direct chat, keeping the composer pinned above the soft keyboard while typing.
- The dashboard backend now orchestrates room responders on that POST: mentioned members are routed as direct responders, additional ambient members may reply (up to the room ambient responder cap), and each assistant reply is persisted with `senderAgentId` via `chatStore.addRoomMessage(...)`.
- If room replies cannot be generated (for example no resolvable responders or all responders fail), the POST fails with an API error (HTTP 502) instead of silently returning only the user message.
- If room responders cannot be resolved or all room-reply generations fail, the POST now returns an error instead of silently succeeding with only the user message, so failures are surfaced deterministically.
- Room responder prompt construction now includes a bounded recent room transcript (role, sender label, timestamp, content) plus an explicit latest-user-message marker, so replies stay thread-aware without unbounded prompt growth.
- The UI still avoids optimistic room echo; after `POST /api/chat/rooms/:id/messages`, it immediately re-fetches authoritative room messages to surface persisted user/assistant replies even if SSE delivery is delayed, and it continues to apply `chat:room:message:*` SSE updates for live fan-out.
- Relationship summary: direct Chat runs one target (agent or model) per session; rooms are shared threads with multiple agent members and now use the same message contract as direct Chat; Quick Chat stays a floating single-target panel and does not host rooms.

View File

@@ -530,6 +530,8 @@ When heartbeat has both (1) and (2-5), the runtime model is used as primary and
3. Global `defaultProvider` + `defaultModelId`
4. Automatic provider/model resolution
For post-merge prompt workflow steps, explicit step-level `modelProvider` + `modelId` overrides take precedence over the merger lane above.
### Title summarization model
Used for task title auto-summarization and (when enabled) AI merge commit summaries.

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",