feat(FN-3588): add immediate wake controls for agent inbox and message API

The merge delivers four major bodies of work. The dominant theme is FN-3588, which adds an "immediate wake" override to the message inbox API, exposes it in the MailboxModal UI, and updates `agent-heartbeat.ts` and `executor.ts` to honor the override alongside timer/signal triggers. FN-3705 gates re

Fusion-Task-Id: FN-3588
This commit is contained in:
Fusion
2026-05-07 10:41:11 -07:00
committed by gsxdsm
parent af82119a75
commit 0efba0890f
10 changed files with 358 additions and 26 deletions

View File

@@ -7613,6 +7613,7 @@ export interface SendMessageInput {
content: string;
type: MessageType;
metadata?: MessageMetadata;
wakeImmediately?: boolean;
}
/** Fetch inbox messages for the current user. */

View File

@@ -889,7 +889,7 @@
align-items: center;
gap: var(--space-sm);
font-size: 0.85rem;
color: var(--text-secondary, var(--text-muted));
color: var(--text-muted);
cursor: pointer;
}
@@ -897,6 +897,10 @@
margin: 0;
}
.message-composer-wake-label:has(input:disabled) {
cursor: default;
}
.message-composer-wake-hint {
margin-left: var(--space-xs);
font-size: 0.75rem;

View File

@@ -1,4 +1,4 @@
import { useState, useCallback } from "react";
import { useState, useCallback, useMemo } from "react";
import { X, Send, Loader2, Bot, AlertCircle } from "lucide-react";
import type { ParticipantType, MessageType } from "@fusion/core";
import { getErrorMessage } from "@fusion/core";
@@ -47,6 +47,11 @@ export function MessageComposer({
const [isSending, setIsSending] = useState(false);
const [error, setError] = useState<string | null>(null);
const selectedAgent = useMemo(() => agents.find((agent) => agent.id === toId), [agents, toId]);
const recipientIsAgent = toType === "agent";
const recipientAlwaysImmediate = recipientIsAgent && selectedAgent?.runtimeConfig?.messageResponseMode === "immediate";
const wakeImmediately = recipientIsAgent && (wakeRecipient || recipientAlwaysImmediate);
const isValid = toId.trim() !== "" && content.trim().length > 0 && content.length <= MAX_CONTENT_LENGTH;
const handleSend = useCallback(async () => {
@@ -57,14 +62,11 @@ export function MessageComposer({
try {
const messageType: MessageType = toType === "agent" ? "user-to-agent" : "system";
const includeWake = wakeRecipient && toType === "agent";
const metadata =
replyContext || includeWake
? {
...(replyContext ? { replyTo: { messageId: replyContext.messageId } } : {}),
...(includeWake ? { wakeRecipient: true } : {}),
}
replyContext
? { replyTo: { messageId: replyContext.messageId } }
: undefined;
const sendWakeImmediately = wakeImmediately;
await sendMessage(
{
toId: toId.trim(),
@@ -72,6 +74,7 @@ export function MessageComposer({
content: content.trim(),
type: messageType,
...(metadata ? { metadata } : {}),
...(sendWakeImmediately ? { wakeImmediately: true } : {}),
},
projectId,
);
@@ -83,7 +86,7 @@ export function MessageComposer({
} finally {
setIsSending(false);
}
}, [isValid, isSending, toId, toType, content, wakeRecipient, replyContext, projectId, onSend, addToast]);
}, [isValid, isSending, toId, toType, content, wakeImmediately, replyContext, projectId, onSend, addToast]);
const handleAgentSelect = useCallback((agentId: string) => {
setToId(agentId);
@@ -174,19 +177,22 @@ export function MessageComposer({
</div>
{/* Wake recipient toggle (agents only) */}
{toType === "agent" && (
{recipientIsAgent && (
<div className="message-composer-field message-composer-field--wake">
<label className="message-composer-wake-label">
<input
type="checkbox"
checked={wakeRecipient}
checked={wakeImmediately}
disabled={recipientAlwaysImmediate}
onChange={(e) => setWakeRecipient(e.target.checked)}
data-testid="message-composer-wake"
/>
<span>
Wake recipient immediately
<span className="message-composer-wake-hint">
(overrides their messageResponseMode)
Wake agent immediately
<span className="message-composer-wake-hint" data-testid="message-composer-wake-hint">
{recipientAlwaysImmediate
? "(agent is already set to immediate response mode)"
: "(one-off override for this message only)"}
</span>
</span>
</label>

View File

@@ -197,7 +197,7 @@ describe("MessageComposer", () => {
});
});
it("forwards wakeRecipient metadata when the wake checkbox is ticked for an agent recipient", async () => {
it("forwards wakeImmediately when the wake checkbox is ticked for an agent recipient", async () => {
render(<MessageComposer {...defaultProps} agents={mockAgents} />);
fireEvent.change(screen.getByTestId("message-composer-recipient"), {
target: { value: "agent-001" },
@@ -211,14 +211,14 @@ describe("MessageComposer", () => {
await waitFor(() => {
expect(mockSendMessage).toHaveBeenCalledWith(
expect.objectContaining({
metadata: { wakeRecipient: true },
wakeImmediately: true,
}),
undefined,
);
});
});
it("merges wakeRecipient with replyTo metadata when replying", async () => {
it("sends wakeImmediately alongside replyTo metadata when replying", async () => {
render(
<MessageComposer
{...defaultProps}
@@ -236,9 +236,9 @@ describe("MessageComposer", () => {
await waitFor(() => {
expect(mockSendMessage).toHaveBeenCalledWith(
expect.objectContaining({
wakeImmediately: true,
metadata: {
replyTo: { messageId: "msg-orig" },
wakeRecipient: true,
},
}),
undefined,
@@ -246,7 +246,7 @@ describe("MessageComposer", () => {
});
});
it("omits wakeRecipient metadata when the checkbox is left unchecked", async () => {
it("omits wakeImmediately when the checkbox is left unchecked", async () => {
render(<MessageComposer {...defaultProps} agents={mockAgents} />);
fireEvent.change(screen.getByTestId("message-composer-recipient"), {
target: { value: "agent-001" },
@@ -258,10 +258,40 @@ describe("MessageComposer", () => {
await waitFor(() => {
const callArgs = mockSendMessage.mock.calls[0][0];
expect(callArgs.metadata).toBeUndefined();
expect(callArgs.wakeImmediately).toBeUndefined();
});
});
it("locks wake checkbox as checked when selected agent is already immediate mode", () => {
const immediateAgents: Agent[] = [
{
...mockAgents[0],
runtimeConfig: { messageResponseMode: "immediate" },
},
];
render(<MessageComposer {...defaultProps} agents={immediateAgents} />);
fireEvent.change(screen.getByTestId("message-composer-recipient"), {
target: { value: "agent-001" },
});
const wakeCheckbox = screen.getByTestId("message-composer-wake") as HTMLInputElement;
expect(wakeCheckbox.checked).toBe(true);
expect(wakeCheckbox.disabled).toBe(true);
expect(screen.getByTestId("message-composer-wake-hint").textContent).toContain("already set to immediate response mode");
});
it("hides wake checkbox for non-agent recipients", () => {
render(
<MessageComposer
{...defaultProps}
agents={mockAgents}
recipient={{ id: "dashboard", type: "user" }}
/>,
);
expect(screen.queryByTestId("message-composer-wake")).toBeNull();
});
it("passes projectId to sendMessage", async () => {
render(<MessageComposer {...defaultProps} agents={mockAgents} projectId="proj-1" />);
fireEvent.change(screen.getByTestId("message-composer-recipient"), {

View File

@@ -3104,6 +3104,198 @@ describe("Messaging Routes", () => {
expect(res.body.id).toBe("msg-runtime-1");
});
it("triggers executeHeartbeat when wakeImmediately is true for agent recipients", async () => {
const executeHeartbeat = vi.fn().mockResolvedValue({ id: "run-1" });
const wakeApp = express();
wakeApp.use(express.json());
wakeApp.use("/api", createApiRoutes(store, {
heartbeatMonitor: { executeHeartbeat, rootDir } as any,
}));
const res = await REQUEST(
wakeApp,
"POST",
"/api/messages",
JSON.stringify({
toId: "agent-wake-1",
toType: "agent",
content: "wake now",
type: "user-to-agent",
wakeImmediately: true,
}),
{ "Content-Type": "application/json" },
);
expect(res.status).toBe(201);
expect(executeHeartbeat).toHaveBeenCalledWith({
agentId: "agent-wake-1",
source: "on_demand",
triggerDetail: "wake-on-message",
});
});
it("does not trigger executeHeartbeat when wakeImmediately is omitted/false or recipient is not an agent", async () => {
const executeHeartbeat = vi.fn().mockResolvedValue({ id: "run-1" });
const wakeApp = express();
wakeApp.use(express.json());
wakeApp.use("/api", createApiRoutes(store, {
heartbeatMonitor: { executeHeartbeat, rootDir } as any,
}));
const noWake = await REQUEST(
wakeApp,
"POST",
"/api/messages",
JSON.stringify({
toId: "agent-no-wake",
toType: "agent",
content: "normal message",
type: "user-to-agent",
}),
{ "Content-Type": "application/json" },
);
const userRecipient = await REQUEST(
wakeApp,
"POST",
"/api/messages",
JSON.stringify({
toId: "dashboard",
toType: "user",
content: "user message",
type: "system",
wakeImmediately: true,
}),
{ "Content-Type": "application/json" },
);
expect(noWake.status).toBe(201);
expect(userRecipient.status).toBe(201);
expect(executeHeartbeat).not.toHaveBeenCalled();
});
it("uses project-scoped heartbeat monitor resolution when default monitor belongs to another root", async () => {
const defaultExecuteHeartbeat = vi.fn().mockResolvedValue({ id: "run-default" });
const projectExecuteHeartbeat = vi.fn().mockResolvedValue({ id: "run-project" });
const engineManager = {
getAllEngines: vi.fn().mockReturnValue(new Map([
[
"project-1",
{
getWorkingDirectory: () => rootDir,
getHeartbeatMonitor: () => ({ executeHeartbeat: projectExecuteHeartbeat }),
},
],
])),
};
const wakeApp = express();
wakeApp.use(express.json());
wakeApp.use("/api", createApiRoutes(store, {
heartbeatMonitor: { executeHeartbeat: defaultExecuteHeartbeat, rootDir: join(rootDir, "other-project") } as any,
engineManager: engineManager as any,
}));
const res = await REQUEST(
wakeApp,
"POST",
"/api/messages",
JSON.stringify({
toId: "agent-project-scope",
toType: "agent",
content: "wake in scoped project",
type: "user-to-agent",
wakeImmediately: true,
}),
{ "Content-Type": "application/json" },
);
expect(res.status).toBe(201);
expect(defaultExecuteHeartbeat).not.toHaveBeenCalled();
expect(projectExecuteHeartbeat).toHaveBeenCalledWith({
agentId: "agent-project-scope",
source: "on_demand",
triggerDetail: "wake-on-message",
});
});
it("returns created message even when wakeImmediately execution throws", async () => {
const executeHeartbeat = vi.fn().mockRejectedValue(new Error("wake failed"));
const wakeApp = express();
wakeApp.use(express.json());
wakeApp.use("/api", createApiRoutes(store, {
heartbeatMonitor: { executeHeartbeat, rootDir } as any,
}));
const res = await REQUEST(
wakeApp,
"POST",
"/api/messages",
JSON.stringify({
toId: "agent-wake-failure",
toType: "agent",
content: "wake best effort",
type: "user-to-agent",
wakeImmediately: true,
}),
{ "Content-Type": "application/json" },
);
expect(res.status).toBe(201);
expect(res.body.toId).toBe("agent-wake-failure");
expect(executeHeartbeat).toHaveBeenCalledTimes(1);
});
it("supports metadata.wakeRecipient as an immediate-wake request", async () => {
const executeHeartbeat = vi.fn().mockResolvedValue({ id: "run-1" });
const wakeApp = express();
wakeApp.use(express.json());
wakeApp.use("/api", createApiRoutes(store, {
heartbeatMonitor: { executeHeartbeat, rootDir } as any,
}));
const res = await REQUEST(
wakeApp,
"POST",
"/api/messages",
JSON.stringify({
toId: "agent-wake-meta",
toType: "agent",
content: "wake via metadata",
type: "user-to-agent",
metadata: { wakeRecipient: true },
}),
{ "Content-Type": "application/json" },
);
expect(res.status).toBe(201);
expect(executeHeartbeat).toHaveBeenCalledWith({
agentId: "agent-wake-meta",
source: "on_demand",
triggerDetail: "wake-on-message",
});
});
it("gracefully no-ops wakeImmediately when no heartbeat monitor is configured", async () => {
const res = await REQUEST(
app,
"POST",
"/api/messages",
JSON.stringify({
toId: "agent-no-monitor",
toType: "agent",
content: "wake request without monitor",
type: "user-to-agent",
wakeImmediately: true,
}),
{ "Content-Type": "application/json" },
);
expect(res.status).toBe(201);
expect(res.body.toId).toBe("agent-no-monitor");
});
it("GET /api/messages/inbox returns dashboard inbox messages", async () => {
const inboxMessage = messageStore.sendMessage({
fromId: "agent-1",
@@ -3231,6 +3423,16 @@ describe("Messaging Routes", () => {
},
message: "metadata.replyTo.messageId must be a non-empty string",
},
{
body: {
toId: "agent-1",
toType: "agent",
content: "x",
type: "user-to-agent",
wakeImmediately: "yes",
},
message: "wakeImmediately must be a boolean",
},
];
for (const testCase of invalidCases) {

View File

@@ -1,11 +1,12 @@
import type { Request } from "express";
import { resolve } from "node:path";
import { DASHBOARD_USER_ID, MessageStore, type MessageType, type ParticipantType, validateMessageMetadata } from "@fusion/core";
import { ApiError, badRequest, notFound } from "../api-error.js";
import { getTerminalService } from "../terminal-service.js";
import type { ApiRoutesContext } from "./types.js";
export function registerMessagingScriptRoutes(ctx: ApiRoutesContext): void {
const { router, options, getProjectContext, rethrowAsApiError } = ctx;
const { router, options, getProjectContext, rethrowAsApiError, runtimeLogger } = ctx;
// ── Scripts API ──────────────────────────────────────────────────────────
@@ -190,6 +191,44 @@ export function registerMessagingScriptRoutes(ctx: ApiRoutesContext): void {
const VALID_MESSAGE_TYPES: MessageType[] = ["agent-to-agent", "agent-to-user", "user-to-agent", "system"];
const VALID_PARTICIPANT_TYPES: ParticipantType[] = ["agent", "user", "system"];
type HeartbeatMonitorHandle = NonNullable<NonNullable<ApiRoutesContext["options"]>["heartbeatMonitor"]>;
const heartbeatMonitor = options?.heartbeatMonitor;
function isHeartbeatMonitorForProject(scopedStore: import("@fusion/core").TaskStore): boolean {
if (!heartbeatMonitor?.rootDir) return true;
try {
const monitorRoot = resolve(heartbeatMonitor.rootDir);
const storeRoot = resolve(scopedStore.getRootDir());
return monitorRoot === storeRoot;
} catch {
return true;
}
}
function resolveHeartbeatMonitor(scopedStore: import("@fusion/core").TaskStore): HeartbeatMonitorHandle | undefined {
const engineManager = options?.engineManager;
if (!engineManager) return undefined;
try {
const storeRoot = resolve(scopedStore.getRootDir());
for (const engine of engineManager.getAllEngines().values()) {
if (resolve(engine.getWorkingDirectory()) === storeRoot) {
const monitor = engine.getHeartbeatMonitor();
if (!monitor) {
return undefined;
}
return {
rootDir: engine.getWorkingDirectory(),
startRun: monitor.startRun.bind(monitor),
executeHeartbeat: monitor.executeHeartbeat.bind(monitor),
stopRun: monitor.stopRun.bind(monitor),
};
}
}
} catch {
// no-op: fallback handled by caller
}
return undefined;
}
router.get("/messages/inbox", async (req, res) => {
try {
@@ -258,7 +297,7 @@ export function registerMessagingScriptRoutes(ctx: ApiRoutesContext): void {
router.post("/messages", async (req, res) => {
try {
const { toId, toType, content, type, metadata } = req.body;
const { toId, toType, content, type, metadata, wakeImmediately } = req.body;
if (!toId || typeof toId !== "string") {
throw badRequest("toId is required");
@@ -276,6 +315,9 @@ export function registerMessagingScriptRoutes(ctx: ApiRoutesContext): void {
if (metadata !== undefined && (typeof metadata !== "object" || metadata === null || Array.isArray(metadata))) {
throw badRequest("metadata must be an object");
}
if (wakeImmediately !== undefined && typeof wakeImmediately !== "boolean") {
throw badRequest("wakeImmediately must be a boolean");
}
try {
validateMessageMetadata(metadata);
@@ -293,6 +335,28 @@ export function registerMessagingScriptRoutes(ctx: ApiRoutesContext): void {
type,
metadata,
});
const shouldWakeImmediately = toType === "agent" && (wakeImmediately === true || metadata?.wakeRecipient === true);
if (shouldWakeImmediately) {
try {
const { store: scopedStore } = await getProjectContext(req);
const resolvedMonitor =
isHeartbeatMonitorForProject(scopedStore)
? heartbeatMonitor
: resolveHeartbeatMonitor(scopedStore);
if (resolvedMonitor) {
await resolvedMonitor.executeHeartbeat({
agentId: toId,
source: "on_demand",
triggerDetail: "wake-on-message",
});
}
} catch (wakeErr) {
runtimeLogger.warn(`POST /api/messages wakeImmediately best-effort wake failed: ${wakeErr instanceof Error ? wakeErr.message : String(wakeErr)}`);
}
}
res.status(201).json(message);
} catch (err: unknown) {
if (err instanceof ApiError) {