FN-8870: add structural approval and report mail
Add typed report and approval metadata to the mailbox contract. - Define and validate structural mail kinds, report sections, and approval references. - Let agents send validated reports while reserving approval mail for engine emission. - Emit idempotent, fail-soft approval notifications and document the contract. Files changed: .changeset/fn-8870-structural-mail-contract.md | 7 ++ docs/agents.md | 2 + docs/architecture.md | 4 + .../message-metadata-structural-mail.test.ts | 25 ++++ packages/core/src/index.gate.ts | 2 +- packages/core/src/index.ts | 2 +- packages/core/src/types.ts | 30 +++++ packages/core/src/types/messaging/messages.ts | 21 ++++ ...gent-tools-send-message-structural-mail.test.ts | 34 +++++ .../src/__tests__/approval-mail-emission.test.ts | 137 +++++++++++++++++++++ packages/engine/src/agent-heartbeat.ts | 3 + packages/engine/src/agent-tools.ts | 28 ++++- packages/engine/src/agents/approval-mail.ts | 40 ++++++ packages/engine/src/executor.ts | 3 + 14 files changed, 334 insertions(+), 4 deletions(-) Fusion-Task-Id: FN-8870 Fusion-Task-Lineage: 7d30b2f9-c004-4975-86b1-9600e2c8d76f Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
7
.changeset/fn-8870-structural-mail-contract.md
Normal file
7
.changeset/fn-8870-structural-mail-contract.md
Normal file
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@runfusion/fusion": minor
|
||||
---
|
||||
|
||||
summary: Add structural reports and approval items to agent mail.
|
||||
category: feature
|
||||
dev: Adds mailKind/report/approvalRequestId metadata, fn_send_message report params, and approval-mail:<approvalRequestId> idempotency.
|
||||
@@ -946,6 +946,8 @@ Messaging is available in dashboard mailbox UI and CLI. In dashboard Mailbox →
|
||||
|
||||
Agent-backed dashboard chat sessions (including plugin-runtime agents such as Hermes/OpenClaw/Paperclip) also expose mailbox tools (`fn_send_message`, `fn_read_messages`) when a `MessageStore` is wired for that project. Model-only chats without an attached agent do not expose these tools.
|
||||
|
||||
Mail has an optional structural metadata contract: `mailKind` distinguishes ordinary messages, reports, and approvals; reports carry a small serializable title/section writeup, while `approvalRequestId` is a reference to live approval state rather than a copied snapshot. `fn_send_message` can send reports with `mail_kind: "report"` and `report`; use Chat for quick back-and-forth. Approval items are engine-emitted only.
|
||||
|
||||
### Dashboard Chat workspace tools
|
||||
|
||||
Dashboard Chat, Chat Room responders, and task-detail Planner Chat run at the interactive project checkout with coding workspace tools: `read`, `write`, `edit`, `bash`, `grep`, `find`, and `ls`. Use them for user-directed file changes and shell investigation. When a durable agent is bound, its permanent-agent permission policy still governs file writes/deletes and command execution; unbound model Chat has no durable-principal policy gate. Chat must keep the checkout branch sticky: inspect Git freely, but do not use `git checkout` or `git switch` unless the operator explicitly requests it.
|
||||
|
||||
@@ -6,6 +6,10 @@ This document describes the actual architecture of Fusion as implemented in this
|
||||
|
||||
---
|
||||
|
||||
## Structural approval mail
|
||||
|
||||
Action gates emit a best-effort, idempotent system-mail item for each pending approval request. The item uses `mailKind: "approval"` and references its live `approvalRequestId`; `sendMessageOnce` keys it as `approval-mail:<approvalRequestId>`, so repeated deduped gate attempts do not spam the inbox and a mailbox failure cannot interrupt the safety pause.
|
||||
|
||||
## Terminal task wedge notifications
|
||||
|
||||
Actionable terminal task updates are classified into bounded reasons such as a named merge gate, retry exhaustion, or a completion blocker. The PostgreSQL-backed task row persists an active/resolved episode with an opaque identity, so `NotificationService` sends one `task-wedged` provider event and one dashboard system-mailbox message per active reason even across restarts. Repeated observations remain quiet until an authoritative non-wedge task update resolves the episode; changed and resolved-then-reentered reasons notify again.
|
||||
|
||||
@@ -0,0 +1,25 @@
|
||||
import { describe, expect, it } from "vitest";
|
||||
import { validateMessageMetadata } from "../types.js";
|
||||
|
||||
describe("validateMessageMetadata structural mail", () => {
|
||||
it("accepts absent, legacy, report, approval, and mail-kind-only metadata", () => {
|
||||
expect(() => validateMessageMetadata(undefined)).not.toThrow();
|
||||
expect(() => validateMessageMetadata({ replyTo: { messageId: "m-1" }, wakeRecipient: true, nativeStructures: [{ kind: "mission", id: "M-1" }], kind: "task-proposal", proposedTask: { title: "Title", description: "Description" }, proposalIdempotencyKey: "proposal-1" })).not.toThrow();
|
||||
expect(() => validateMessageMetadata({ mailKind: "report", report: { title: "Brief", sections: [{ heading: "Summary", body: "Ready" }] } })).not.toThrow();
|
||||
expect(() => validateMessageMetadata({ mailKind: "approval", approvalRequestId: "approval-1" })).not.toThrow();
|
||||
expect(() => validateMessageMetadata({ mailKind: "message" })).not.toThrow();
|
||||
});
|
||||
|
||||
it.each([
|
||||
[{ mailKind: "other" }, "metadata.mailKind"],
|
||||
[{ report: "bad" }, "metadata.report must be an object"],
|
||||
[{ report: { title: " ", sections: [{ heading: "H", body: "B" }] } }, "metadata.report.title"],
|
||||
[{ report: { title: "T", sections: "bad" } }, "metadata.report.sections must be an array"],
|
||||
[{ report: { title: "T", sections: [] } }, "metadata.report.sections must not be empty"],
|
||||
[{ report: { title: "T", sections: [{ body: "B" }] } }, "metadata.report.sections[].heading"],
|
||||
[{ report: { title: "T", sections: [{ heading: "H", body: " " }] } }, "metadata.report.sections[].body"],
|
||||
[{ approvalRequestId: " " }, "metadata.approvalRequestId"],
|
||||
])("rejects malformed structural mail %#", (metadata, message) => {
|
||||
expect(() => validateMessageMetadata(metadata as never)).toThrow(message);
|
||||
});
|
||||
});
|
||||
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
@@ -1404,6 +1404,9 @@ import type {
|
||||
EphemeralTaskCreationPolicy,
|
||||
ProposedTaskMetadata,
|
||||
NativeStructureEmbed,
|
||||
MailKind,
|
||||
MailReportSection,
|
||||
MailReport,
|
||||
MessageMetadata,
|
||||
Message,
|
||||
MessageCreateInput,
|
||||
@@ -1416,6 +1419,9 @@ export type {
|
||||
EphemeralTaskCreationPolicy,
|
||||
ProposedTaskMetadata,
|
||||
NativeStructureEmbed,
|
||||
MailKind,
|
||||
MailReportSection,
|
||||
MailReport,
|
||||
MessageMetadata,
|
||||
Message,
|
||||
MessageCreateInput,
|
||||
@@ -1448,6 +1454,30 @@ export function validateMessageMetadata(metadata: MessageMetadata | undefined):
|
||||
plugin-owned read adapter at render time, so attachment metadata remains a ref rather than a
|
||||
duplicated persistence snapshot; labels are optional attach-time fallbacks.
|
||||
*/
|
||||
/*
|
||||
FNXC:StructuralMail 2026-08-09-07:16:
|
||||
Report sections are required: an empty report is a quick message wearing a report label and
|
||||
would render as an empty shell for the structural-mail consumer.
|
||||
*/
|
||||
if (metadata.mailKind !== undefined && !["message", "report", "approval"].includes(metadata.mailKind)) {
|
||||
throw new Error("metadata.mailKind is invalid");
|
||||
}
|
||||
if (metadata.report !== undefined) {
|
||||
const report = metadata.report;
|
||||
if (typeof report !== "object" || report === null || Array.isArray(report)) throw new Error("metadata.report must be an object");
|
||||
if (typeof report.title !== "string" || !report.title.trim()) throw new Error("metadata.report.title must be a non-empty string");
|
||||
if (!Array.isArray(report.sections)) throw new Error("metadata.report.sections must be an array");
|
||||
if (report.sections.length === 0) throw new Error("metadata.report.sections must not be empty");
|
||||
for (const section of report.sections) {
|
||||
if (typeof section !== "object" || section === null || Array.isArray(section)) throw new Error("metadata.report.sections entries must be objects");
|
||||
if (typeof section.heading !== "string" || !section.heading.trim()) throw new Error("metadata.report.sections[].heading must be a non-empty string");
|
||||
if (typeof section.body !== "string" || !section.body.trim()) throw new Error("metadata.report.sections[].body must be a non-empty string");
|
||||
}
|
||||
}
|
||||
if (metadata.approvalRequestId !== undefined && (typeof metadata.approvalRequestId !== "string" || !metadata.approvalRequestId.trim())) {
|
||||
throw new Error("metadata.approvalRequestId must be a non-empty string");
|
||||
}
|
||||
|
||||
if (metadata.nativeStructures !== undefined) {
|
||||
if (!Array.isArray(metadata.nativeStructures)) {
|
||||
throw new Error("metadata.nativeStructures must be an array");
|
||||
|
||||
@@ -72,6 +72,18 @@ export interface ProposedTaskMetadata {
|
||||
*/
|
||||
export type NativeStructureEmbed = NativeStructureRef & { label?: string };
|
||||
|
||||
export type MailKind = "message" | "report" | "approval";
|
||||
|
||||
export interface MailReportSection {
|
||||
heading: string;
|
||||
body: string;
|
||||
}
|
||||
|
||||
export interface MailReport {
|
||||
title: string;
|
||||
sections: MailReportSection[];
|
||||
}
|
||||
|
||||
export interface MessageMetadata extends Record<string, unknown> {
|
||||
/** Optional link to the original message when this message is a reply. */
|
||||
replyTo?: MessageReplyReference;
|
||||
@@ -110,6 +122,15 @@ export interface MessageMetadata extends Record<string, unknown> {
|
||||
* unavailable target a human-readable fallback after lazy preview resolution.
|
||||
*/
|
||||
nativeStructures?: NativeStructureEmbed[];
|
||||
/*
|
||||
FNXC:StructuralMail 2026-08-09-07:16:
|
||||
Absent mailKind preserves every legacy row's quick-message semantics. Reports are a small
|
||||
serializable writeup alongside (not instead of) nativeStructures; approvalRequestId is a
|
||||
reference only, never an approval snapshot, so live state resolves from ApprovalRequestStore.
|
||||
*/
|
||||
mailKind?: MailKind;
|
||||
report?: MailReport;
|
||||
approvalRequestId?: string;
|
||||
}
|
||||
|
||||
/** Message record stored in the system */
|
||||
|
||||
@@ -0,0 +1,34 @@
|
||||
import { describe, expect, it, vi } from "vitest";
|
||||
import { createSendMessageTool } from "../agent-tools.js";
|
||||
|
||||
const execute = async (params: Record<string, unknown>, parent?: Record<string, unknown>) => {
|
||||
const sendMessage = vi.fn(async () => ({ id: "msg-1" }));
|
||||
const tool = createSendMessageTool({ sendMessage, getMessage: vi.fn(async () => parent ?? null) } as never, "agent-a");
|
||||
return { result: await tool.execute("1", params as never), sendMessage };
|
||||
};
|
||||
const text = (result: any) => result.content[0].text as string;
|
||||
|
||||
describe("fn_send_message structural mail", () => {
|
||||
it("persists report metadata and merges a reply reference", async () => {
|
||||
const report = { title: "Brief", sections: [{ heading: "Summary", body: "Done" }] };
|
||||
const { result, sendMessage } = await execute({ to_id: "dashboard", type: "agent-to-user", content: "See report", mail_kind: "report", report });
|
||||
expect(text(result)).toContain("Message sent");
|
||||
expect(sendMessage).toHaveBeenCalledWith(expect.objectContaining({ metadata: { mailKind: "report", report } }));
|
||||
const reply = await execute({ content: "See report", reply_to_message_id: "parent", mail_kind: "report", report }, { fromId: "dashboard", fromType: "user", toId: "agent-a", toType: "agent" });
|
||||
expect(reply.sendMessage).toHaveBeenCalledWith(expect.objectContaining({ metadata: { replyTo: { messageId: "parent" }, mailKind: "report", report } }));
|
||||
});
|
||||
it("keeps plain sends metadata-free", async () => {
|
||||
const { sendMessage } = await execute({ to_id: "agent-b", content: "hello" });
|
||||
expect(sendMessage).toHaveBeenCalledWith(expect.not.objectContaining({ metadata: expect.anything() }));
|
||||
});
|
||||
it.each([
|
||||
{ to_id: "agent-b", content: "x", mail_kind: "report" },
|
||||
{ to_id: "agent-b", content: "x", report: { title: " ", sections: [{ heading: "H", body: "B" }] } },
|
||||
{ to_id: "agent-b", content: "x", report: { title: "T", sections: [] } },
|
||||
{ to_id: "agent-b", content: "x", report: { title: "T", sections: [{ heading: " ", body: "B" }] } },
|
||||
{ to_id: "agent-b", content: "x", mail_kind: "approval" },
|
||||
])("rejects malformed structural input", async (params) => {
|
||||
const { result, sendMessage } = await execute(params);
|
||||
expect(text(result)).toMatch(/^ERROR:/); expect(sendMessage).not.toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
137
packages/engine/src/__tests__/approval-mail-emission.test.ts
Normal file
137
packages/engine/src/__tests__/approval-mail-emission.test.ts
Normal file
@@ -0,0 +1,137 @@
|
||||
import "./executor-test-helpers.js";
|
||||
import { describe, expect, it, vi } from "vitest";
|
||||
import { DASHBOARD_USER_ID } from "@fusion/core";
|
||||
import { TaskExecutor } from "../executor.js";
|
||||
import { HeartbeatMonitor } from "../agent-heartbeat.js";
|
||||
import { emitApprovalMail } from "../agents/approval-mail.js";
|
||||
|
||||
function createMessageStore(options: { reject?: boolean } = {}) {
|
||||
const persisted = new Map<string, unknown>();
|
||||
const sendMessageOnce = vi.fn(async (message: unknown, key: string) => {
|
||||
if (options.reject) throw new Error("mail store unavailable");
|
||||
const inserted = !persisted.has(key);
|
||||
if (inserted) persisted.set(key, message);
|
||||
return { message: { id: key }, inserted };
|
||||
});
|
||||
return { sendMessageOnce, persisted };
|
||||
}
|
||||
|
||||
function createExecutorStore() {
|
||||
const listeners = new Map<string, Set<(...args: unknown[]) => void>>();
|
||||
return {
|
||||
on: vi.fn((event: string, listener: (...args: unknown[]) => void) => {
|
||||
const handlers = listeners.get(event) ?? new Set();
|
||||
handlers.add(listener);
|
||||
listeners.set(event, handlers);
|
||||
}),
|
||||
off: vi.fn(),
|
||||
getSettings: vi.fn().mockResolvedValue({ globalPause: false, enginePaused: false }),
|
||||
listTasks: vi.fn().mockResolvedValue([]),
|
||||
pauseTask: vi.fn().mockResolvedValue(undefined),
|
||||
logEntry: vi.fn().mockResolvedValue(undefined),
|
||||
};
|
||||
}
|
||||
|
||||
const actionDecision = {
|
||||
disposition: "require-approval",
|
||||
category: "command_execution",
|
||||
toolName: "bash",
|
||||
operation: "shell command",
|
||||
summary: "bash: shell command",
|
||||
resourceType: "command",
|
||||
approvalDedupeKey: "approval-dedupe-key",
|
||||
metadata: {},
|
||||
};
|
||||
|
||||
const heartbeatAgent = { id: "agent-1", name: "Ada", permissionPolicy: undefined };
|
||||
|
||||
/**
|
||||
* FNXC:StructuralMail 2026-08-09-07:16:
|
||||
* Approval mail is produced by four independently dispatched gate closures. Exercise those
|
||||
* closures, rather than only scanning their source, so future refactors cannot silently drop a
|
||||
* mailbox write from one executor or heartbeat gate path.
|
||||
*/
|
||||
describe("approval-mail emission", () => {
|
||||
it("writes one reference-only approval mailbox item with a deterministic key", async () => {
|
||||
const messageStore = createMessageStore();
|
||||
await emitApprovalMail({ messageStore: messageStore as never, approvalRequestId: "approval-1", toolName: "fn_bash", taskId: "FN-1", agentName: "Ada" });
|
||||
|
||||
expect(messageStore.sendMessageOnce).toHaveBeenCalledWith(
|
||||
expect.objectContaining({
|
||||
toId: DASHBOARD_USER_ID,
|
||||
metadata: { mailKind: "approval", approvalRequestId: "approval-1", taskId: "FN-1" },
|
||||
}),
|
||||
"approval-mail:approval-1",
|
||||
);
|
||||
const message = messageStore.persisted.get("approval-mail:approval-1") as { metadata: Record<string, unknown> };
|
||||
expect(message.metadata).not.toHaveProperty("status");
|
||||
expect(message.metadata).not.toHaveProperty("decision");
|
||||
});
|
||||
|
||||
it("keeps a deduped re-raised approval to one durable mailbox row", async () => {
|
||||
const messageStore = createMessageStore();
|
||||
const input = { messageStore: messageStore as never, approvalRequestId: "approval-repeat", toolName: "bash" };
|
||||
|
||||
await emitApprovalMail(input);
|
||||
await emitApprovalMail(input);
|
||||
|
||||
expect(messageStore.sendMessageOnce).toHaveBeenCalledTimes(2);
|
||||
expect(messageStore.persisted).toHaveLength(1);
|
||||
});
|
||||
|
||||
it("swallows mail failures and supports no store", async () => {
|
||||
await expect(emitApprovalMail({ messageStore: createMessageStore({ reject: true }) as never, approvalRequestId: "a", toolName: "tool" })).resolves.toBeUndefined();
|
||||
await expect(emitApprovalMail({ approvalRequestId: "a", toolName: "tool" })).resolves.toBeUndefined();
|
||||
});
|
||||
|
||||
it("emits approval mail from both executor gate closures without interrupting their pauses", async () => {
|
||||
const messageStore = createMessageStore();
|
||||
const store = createExecutorStore();
|
||||
const executor = new TaskExecutor(store as never, "/tmp/test", { messageStore: messageStore as never });
|
||||
vi.spyOn(executor, "awaitAbortInFlightTaskWork").mockResolvedValue(undefined);
|
||||
|
||||
const actionContext = (executor as any).buildActionGateContext("FN-action", null, undefined);
|
||||
const permanentContext = (executor as any).buildPermanentAgentGatingContext("FN-permanent", null, undefined);
|
||||
await actionContext.pauseForApproval({ approvalRequestId: "executor-action", decision: actionDecision });
|
||||
await permanentContext.pauseForApproval({ approvalRequestId: "executor-permanent", toolName: "git_push" });
|
||||
await new Promise((resolve) => setImmediate(resolve));
|
||||
|
||||
expect(store.pauseTask).toHaveBeenCalledTimes(2);
|
||||
expect(messageStore.persisted).toHaveLength(2);
|
||||
expect([...messageStore.persisted.values()]).toEqual(expect.arrayContaining([
|
||||
expect.objectContaining({ toId: DASHBOARD_USER_ID, metadata: expect.objectContaining({ mailKind: "approval", approvalRequestId: "executor-action" }) }),
|
||||
expect.objectContaining({ toId: DASHBOARD_USER_ID, metadata: expect.objectContaining({ mailKind: "approval", approvalRequestId: "executor-permanent" }) }),
|
||||
]));
|
||||
});
|
||||
|
||||
it("emits approval mail from both heartbeat gate closures and tolerates mailbox failure", async () => {
|
||||
const messageStore = createMessageStore({ reject: true });
|
||||
const taskStore = { pauseTask: vi.fn().mockResolvedValue(undefined), logEntry: vi.fn().mockResolvedValue(undefined) };
|
||||
const agentStore = {
|
||||
updateAgentState: vi.fn().mockResolvedValue(undefined),
|
||||
updateAgent: vi.fn().mockResolvedValue(undefined),
|
||||
};
|
||||
const monitor = new HeartbeatMonitor({ store: agentStore as never, taskStore: taskStore as never, messageStore: messageStore as never, rootDir: "/tmp/test" });
|
||||
|
||||
const actionContext = (monitor as any).buildActionGateContext(heartbeatAgent, "FN-heartbeat-action", "run-1");
|
||||
const permanentContext = (monitor as any).buildPermanentAgentGatingContext(heartbeatAgent, "FN-heartbeat-permanent", "run-1");
|
||||
await expect(actionContext.pauseForApproval({ approvalRequestId: "heartbeat-action", decision: actionDecision })).resolves.toBeUndefined();
|
||||
await expect(permanentContext.pauseForApproval({ approvalRequestId: "heartbeat-permanent", toolName: "git_push" })).resolves.toBeUndefined();
|
||||
await new Promise((resolve) => setImmediate(resolve));
|
||||
|
||||
expect(taskStore.pauseTask).toHaveBeenCalledTimes(2);
|
||||
expect(agentStore.updateAgentState).toHaveBeenCalledTimes(2);
|
||||
expect(messageStore.sendMessageOnce).toHaveBeenCalledWith(expect.objectContaining({ metadata: expect.objectContaining({ mailKind: "approval", approvalRequestId: "heartbeat-action" }) }), "approval-mail:heartbeat-action");
|
||||
expect(messageStore.sendMessageOnce).toHaveBeenCalledWith(expect.objectContaining({ metadata: expect.objectContaining({ mailKind: "approval", approvalRequestId: "heartbeat-permanent" }) }), "approval-mail:heartbeat-permanent");
|
||||
});
|
||||
|
||||
it("keeps all gate paths usable without a message store", async () => {
|
||||
const taskStore = { pauseTask: vi.fn().mockResolvedValue(undefined), logEntry: vi.fn().mockResolvedValue(undefined) };
|
||||
const agentStore = { updateAgentState: vi.fn().mockResolvedValue(undefined), updateAgent: vi.fn().mockResolvedValue(undefined) };
|
||||
const monitor = new HeartbeatMonitor({ store: agentStore as never, taskStore: taskStore as never, rootDir: "/tmp/test" });
|
||||
|
||||
await expect((monitor as any).buildActionGateContext(heartbeatAgent, "FN-no-mail", "run-1").pauseForApproval({ approvalRequestId: "no-mail", decision: actionDecision })).resolves.toBeUndefined();
|
||||
expect(taskStore.pauseTask).toHaveBeenCalledOnce();
|
||||
expect(agentStore.updateAgentState).toHaveBeenCalledOnce();
|
||||
});
|
||||
});
|
||||
@@ -44,6 +44,7 @@ import { Type, type Static } from "@earendil-works/pi-ai";
|
||||
import { createHash } from "node:crypto";
|
||||
import { createTaskCreateTool, createTaskLogToolWithContext, createTaskLogsReadTool, createTaskDocumentWriteTool, createTaskDocumentReadTool, createTaskReadTools, createArtifactRegisterTool, createArtifactListTool, createArtifactViewTool, createListAgentsTool, createDelegateTaskTool, createTaskAssignTool, createGetAgentConfigTool, createUpdateAgentConfigTool, createAgentCreateTool, createAgentDeleteTool, createSendMessageTool, createReadMessagesTool, createPostRoomMessageTool, createMemoryTools, createGoalRetrievalTools, createMissionTools, createIdeationTools, createReadEvaluationsTool, createUpdateIdentityTool, createReflectOnPerformanceTool, createWebFetchTool, createWorkflowListTool, createWorkflowGetTool, createWorkflowValidateTool, createWorkflowSelectTool, createTaskPromoteTool, createWorkflowCreateTool, createWorkflowUpdateTool, createWorkflowDeleteTool, createWorkflowSettingsTool, createTraitListTool, createAskQuestionTool, createResearchTools, readAgentMemoryWorkspaceLongTerm, taskCreateParams } from "./agent-tools.js";
|
||||
import { AgentLogger } from "./agents/agent-logger.js";
|
||||
import { emitApprovalMail } from "./agents/approval-mail.js";
|
||||
import {
|
||||
resolveAgentInstructionsWithRatings,
|
||||
buildPluginPromptSection,
|
||||
@@ -1099,6 +1100,7 @@ export class HeartbeatMonitor {
|
||||
`Approval required for ${decision.toolName}. Request ${approvalRequestId} created; task and agent paused awaiting decision.`,
|
||||
);
|
||||
}
|
||||
void emitApprovalMail({ messageStore: this.messageStore, approvalRequestId, toolName: decision.toolName, taskId, agentId: agent.id, agentName: agent.name });
|
||||
await this.store.updateAgentState(agent.id, "paused");
|
||||
await this.store.updateAgent(agent.id, { pauseReason: "awaiting-approval" });
|
||||
},
|
||||
@@ -1171,6 +1173,7 @@ export class HeartbeatMonitor {
|
||||
}
|
||||
await this.store.updateAgentState(agent.id, "paused");
|
||||
await this.store.updateAgent(agent.id, { pauseReason: "awaiting-approval" });
|
||||
void emitApprovalMail({ messageStore: this.messageStore, approvalRequestId, toolName, taskId, agentId: agent.id, agentName: agent.name });
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
@@ -560,6 +560,13 @@ export const sendMessageParams = Type.Object({
|
||||
reply_to_message_id: Type.Optional(
|
||||
Type.String({ description: "Optional ID of the message you are replying to. Parent-based recipient inference is allowed only when that message was addressed to you." }),
|
||||
),
|
||||
mail_kind: Type.Optional(Type.Union([
|
||||
Type.Literal("message"), Type.Literal("report"), Type.Literal("approval"),
|
||||
], { description: "Structural mail kind. Use report for a composed writeup; approval is engine-managed." })),
|
||||
report: Type.Optional(Type.Object({
|
||||
title: Type.String({ description: "Report title" }),
|
||||
sections: Type.Array(Type.Object({ heading: Type.String(), body: Type.String() }), { description: "Non-empty report sections" }),
|
||||
}, { description: "Structured report payload for mail" })),
|
||||
});
|
||||
|
||||
export const readMessagesParams = Type.Object({
|
||||
@@ -5648,7 +5655,7 @@ export function createSendMessageTool(
|
||||
"Send a message to another agent or user. The recipient will be woken if they have " +
|
||||
"`messageResponseMode: 'immediate'` configured. When replying, include `reply_to_message_id`; omit " +
|
||||
"`to_id` to reply to that message's sender only when the parent was addressed to you. Otherwise provide " +
|
||||
"the exact recipient ID and appropriate type explicitly.",
|
||||
"the exact recipient ID and appropriate type explicitly. Use mail for structured reports and approval items; use chat for quick back-and-forth.",
|
||||
parameters: sendMessageParams,
|
||||
// eslint-disable-next-line @typescript-eslint/no-explicit-any
|
||||
execute: async (_id: string, params: Static<typeof sendMessageParams>, _signal?: any, _onUpdate?: any, _ctx?: any) => {
|
||||
@@ -5667,6 +5674,17 @@ export function createSendMessageTool(
|
||||
}
|
||||
|
||||
try {
|
||||
if (params.mail_kind === "approval") {
|
||||
return { content: [{ type: "text" as const, text: "ERROR: approval mail is emitted by the engine only" }], details: {} };
|
||||
}
|
||||
if (params.mail_kind === "report" && !params.report) {
|
||||
return { content: [{ type: "text" as const, text: "ERROR: mail_kind report requires report" }], details: {} };
|
||||
}
|
||||
if (params.report) {
|
||||
if (!params.report.title.trim()) return { content: [{ type: "text" as const, text: "ERROR: report.title must be a non-empty string" }], details: {} };
|
||||
if (params.report.sections.length === 0) return { content: [{ type: "text" as const, text: "ERROR: report.sections must not be empty" }], details: {} };
|
||||
if (params.report.sections.some((section) => !section.heading.trim() || !section.body.trim())) return { content: [{ type: "text" as const, text: "ERROR: report sections require non-empty heading and body" }], details: {} };
|
||||
}
|
||||
const replyToMessageId = params.reply_to_message_id?.trim();
|
||||
|
||||
if (params.reply_to_message_id !== undefined && !replyToMessageId) {
|
||||
@@ -5753,7 +5771,13 @@ export function createSendMessageTool(
|
||||
toType: recipient.type,
|
||||
content,
|
||||
type: messageType,
|
||||
...(replyToMessageId ? { metadata: { replyTo: { messageId: replyToMessageId } } } : {}),
|
||||
...((replyToMessageId || params.mail_kind || params.report) ? {
|
||||
metadata: {
|
||||
...(replyToMessageId ? { replyTo: { messageId: replyToMessageId } } : {}),
|
||||
...(params.mail_kind ? { mailKind: params.mail_kind } : {}),
|
||||
...(params.report ? { report: params.report } : {}),
|
||||
},
|
||||
} : {}),
|
||||
}),
|
||||
correlation: { kind: "direct", fromAgentId, toId: recipient.id },
|
||||
}, options?.autoRecovery ?? { mode: "deterministic-only", maxRetries: 3 }, async () => {
|
||||
|
||||
40
packages/engine/src/agents/approval-mail.ts
Normal file
40
packages/engine/src/agents/approval-mail.ts
Normal file
@@ -0,0 +1,40 @@
|
||||
import { DASHBOARD_USER_ID, type MessageCreateInput, type MessageStore } from "@fusion/core";
|
||||
import { createLogger } from "../logger.js";
|
||||
|
||||
const approvalMailLog = createLogger("approval-mail");
|
||||
|
||||
/**
|
||||
* FNXC:StructuralMail 2026-08-09-07:16:
|
||||
* Pending-approval lookup deliberately re-invokes its pause callbacks, so mailbox delivery is
|
||||
* idempotent by approvalRequestId. The write is fail-soft: mailbox persistence must never break
|
||||
* the approval pause that protects the gated action.
|
||||
*/
|
||||
export async function emitApprovalMail(input: {
|
||||
messageStore?: Pick<MessageStore, "sendMessageOnce">;
|
||||
approvalRequestId: string;
|
||||
toolName: string;
|
||||
taskId?: string;
|
||||
agentId?: string;
|
||||
agentName?: string;
|
||||
}): Promise<void> {
|
||||
try {
|
||||
if (!input.messageStore?.sendMessageOnce) return;
|
||||
const requester = input.agentName ?? input.agentId ?? "An agent";
|
||||
const message: MessageCreateInput = {
|
||||
fromId: "system",
|
||||
fromType: "system",
|
||||
toId: DASHBOARD_USER_ID,
|
||||
toType: "user",
|
||||
type: "system",
|
||||
content: `**Approval required for ${input.toolName}**\n\n${requester} requested approval before this action can continue.`,
|
||||
metadata: {
|
||||
mailKind: "approval",
|
||||
approvalRequestId: input.approvalRequestId,
|
||||
...(input.taskId ? { taskId: input.taskId } : {}),
|
||||
},
|
||||
};
|
||||
await input.messageStore.sendMessageOnce(message, `approval-mail:${input.approvalRequestId}`);
|
||||
} catch (error) {
|
||||
approvalMailLog.warn(`approval mailbox message failed: ${error instanceof Error ? error.message : String(error)}`);
|
||||
}
|
||||
}
|
||||
@@ -204,6 +204,7 @@ import {
|
||||
import { BranchAttributionError, filterFilesToOwnTaskCommits } from "./execution/branch-attribution.js";
|
||||
import { resolveIntegrationBranch } from "./merge/integration-branch.js";
|
||||
import { AgentLogger } from "./agents/agent-logger.js";
|
||||
import { emitApprovalMail } from "./agents/approval-mail.js";
|
||||
import { createLogger, executorLog, reviewerLog, formatError } from "./logger.js";
|
||||
import { TokenCapDetector } from "./errors/token-cap-detector.js";
|
||||
import { isUsageLimitError, checkSessionError, type UsageLimitPauser } from "./errors/usage-limit-detector.js";
|
||||
@@ -2937,6 +2938,7 @@ export class TaskExecutor {
|
||||
executorLog.warn(`${taskId}: failed to suspend in-flight session while awaiting approval: ${error instanceof Error ? error.message : String(error)}`);
|
||||
});
|
||||
}
|
||||
void emitApprovalMail({ messageStore: this.options.messageStore, approvalRequestId, toolName: decision.toolName, taskId, agentId: actorId, agentName: actorName });
|
||||
if (agent && this.options.agentStore) {
|
||||
await this.options.agentStore.updateAgentState(agent.id, "paused");
|
||||
await this.options.agentStore.updateAgent(agent.id, { pauseReason: "awaiting-approval" });
|
||||
@@ -3028,6 +3030,7 @@ export class TaskExecutor {
|
||||
this.approvalSuspended.delete(taskId);
|
||||
throw error;
|
||||
}
|
||||
void emitApprovalMail({ messageStore: this.options.messageStore, approvalRequestId, toolName, taskId, agentId: actorId, agentName: actorName });
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user