feat(FN-5425): add coordination notice system for multi-agent room task fil

Adds a `room-coordination.ts` helper module for multi-agent coordination messaging, wires coordination notices into the agent heartbeat path, and ships tests plus documentation for the feature. A changeset prepares the `@runfusion/fusion` package for release.

Fusion-Task-Id: FN-5425
This commit is contained in:
Fusion (runfusion.ai)
2026-05-21 00:10:32 -07:00
committed by gsxdsm
parent 5c15031416
commit 6fdfb1ad5c
9 changed files with 752 additions and 7 deletions

View File

@@ -0,0 +1,12 @@
---
"@runfusion/fusion": patch
---
Add a Room Coordination Notices prompt-injected advisory that fires when a
user posts an explicit "file a task" / "create a task" request into a chat
room with multiple agent members. Each agent is instructed to either post a
one-line claim before calling fn_task_create (claim branch) or to defer and
acknowledge a peer's prior claim/announcement (defer-suggested branch),
reducing the upstream duplicate pressure on the existing FN-4918 / FN-4829
/ FN-5152 / FN-5220 dedup backstop. Emits a structured
room:coordination:branch run-audit event per decision. Single-agent rooms
and non-task-filing messages are unaffected.

View File

@@ -950,6 +950,32 @@ Heartbeat runs now surface both direct-message inbox traffic and recent room act
5. **Mark as Read**: After successful heartbeat completion, direct inbox messages are marked as read. 5. **Mark as Read**: After successful heartbeat completion, direct inbox messages are marked as read.
6. **Failed Runs**: If the heartbeat execution fails, inbox messages remain unread for retry on the next run. 6. **Failed Runs**: If the heartbeat execution fails, inbox messages remain unread for retry on the next run.
### Room Coordination Notices (FN-5425)
Heartbeat prompts may include a **Room Coordination Notices** section after **Room Ambiguity Notices** when both conditions are true:
1. A pending room message contains explicit task-filing intent (for example, "file a task" / "create a task").
2. The room currently has at least two active agent members.
This notice is advisory prompt-layer routing (not server-side serialization), with two branches:
- **claim**: the agent should post a one-line claim via `fn_post_room_message` first, then call `fn_task_create`, then post the resulting `FN-NNNN` task ID back to the room.
- **defer-suggested**: a peer claim or task announcement was already seen in recent room history; the agent should **not** call `fn_task_create`, and should instead acknowledge the prior claim/announcement via `fn_post_room_message`.
Deterministic duplicate prevention remains authoritative: FN-4918, FN-4829, FN-5152, and FN-5220 are still the hard intake backstop. The coordination notice reduces upstream duplicate pressure but does not replace those guards.
Each coordination decision emits a run-audit `mutationType` of `room:coordination:branch` with metadata:
- `roomId`
- `agentId`
- `branch` (`"claim" | "defer-suggested"`)
- `memberCount`
- `intentCue`
- `priorClaimMessageId`
- `priorTaskId`
Layering order is intentional: ambiguity guidance renders first, coordination guidance renders second.
### Message Response Modes ### Message Response Modes
The `messageResponseMode` runtime configuration controls when agents are triggered by incoming messages: The `messageResponseMode` runtime configuration controls when agents are triggered by incoming messages:

View File

@@ -98,6 +98,13 @@ function createMockStore(overrides: Partial<TaskStore> = {}): TaskStore {
isRunning: false, isRunning: false,
lastCheckedAt: null, lastCheckedAt: null,
}), }),
refreshDatabaseHealth: vi.fn().mockReturnValue({
healthy: true,
corruptionDetected: false,
corruptionErrors: [],
isRunning: false,
lastCheckedAt: null,
}),
getTaskIdIntegrityReport: vi.fn().mockReturnValue({ getTaskIdIntegrityReport: vi.fn().mockReturnValue({
status: "ok", status: "ok",
checkedAt: "2026-05-12T00:00:00.000Z", checkedAt: "2026-05-12T00:00:00.000Z",
@@ -506,7 +513,7 @@ describe("createServer health and headless mode", () => {
it("reports degraded status from /api/health/refresh when corruption is detected", async () => { it("reports degraded status from /api/health/refresh when corruption is detected", async () => {
const store = createMockStore({ const store = createMockStore({
getDatabaseHealth: vi.fn().mockReturnValue({ refreshDatabaseHealth: vi.fn().mockReturnValue({
healthy: false, healthy: false,
corruptionDetected: true, corruptionDetected: true,
corruptionErrors: ["bad row"], corruptionErrors: ["bad row"],

View File

@@ -109,7 +109,11 @@ exports[`agent-heartbeat procedure templates > keeps strict no-task procedure st
reply_to_message_id when answering. If Pending Room Messages are present, reply_to_message_id when answering. If Pending Room Messages are present,
review them in the prompt and use fn_post_room_message only when relevant. review them in the prompt and use fn_post_room_message only when relevant.
When Room Ambiguity Notices appear, follow the resolve/clarify branch and do When Room Ambiguity Notices appear, follow the resolve/clarify branch and do
not create tasks under clarification notices. not create tasks under clarification notices. If a Room Coordination Notices
section is present, follow its claim/defer branch exactly: under "claim" post
a one-line claim before calling fn_task_create; under "defer-suggested" do
NOT call fn_task_create and instead acknowledge the prior claim via
fn_post_room_message.
3. **Wake delta** — read the Wake Delta block above. The wake reason is the 3. **Wake delta** — read the Wake Delta block above. The wake reason is the
highest-priority change for this heartbeat. If you were woken by a comment highest-priority change for this heartbeat. If you were woken by a comment
or a message, acknowledge it before doing anything else. or a message, acknowledge it before doing anything else.
@@ -147,7 +151,11 @@ exports[`agent-heartbeat procedure templates > keeps strict task procedure stabl
reply_to_message_id when answering. If Pending Room Messages are present, reply_to_message_id when answering. If Pending Room Messages are present,
review them in the prompt and use fn_post_room_message only when relevant. review them in the prompt and use fn_post_room_message only when relevant.
When Room Ambiguity Notices appear, follow the resolve/clarify branch and do When Room Ambiguity Notices appear, follow the resolve/clarify branch and do
not create tasks under clarification notices. not create tasks under clarification notices. If a Room Coordination Notices
section is present, follow its claim/defer branch exactly: under "claim" post
a one-line claim before calling fn_task_create; under "defer-suggested" do
NOT call fn_task_create and instead acknowledge the prior claim via
fn_post_room_message.
3. **Wake delta** — read the Wake Delta block above. The wake reason is the 3. **Wake delta** — read the Wake Delta block above. The wake reason is the
highest-priority change for this heartbeat. If you were woken by a comment highest-priority change for this heartbeat. If you were woken by a comment
or a message, acknowledge it before doing anything else. or a message, acknowledge it before doing anything else.

View File

@@ -4,6 +4,7 @@ import { tmpdir } from "node:os";
import { join } from "node:path"; import { join } from "node:path";
import { AgentStore, ChatStore, TaskStore } from "@fusion/core"; import { AgentStore, ChatStore, TaskStore } from "@fusion/core";
import { HeartbeatMonitor } from "../agent-heartbeat.js"; import { HeartbeatMonitor } from "../agent-heartbeat.js";
import * as roomCoordination from "../room-coordination.js";
const sessionCapture = vi.hoisted(() => ({ const sessionCapture = vi.hoisted(() => ({
prompt: "", prompt: "",
@@ -283,6 +284,176 @@ describe("heartbeat room messages", () => {
expect(sessionCapture.prompt).not.toContain("Resolved Referent:"); expect(sessionCapture.prompt).not.toContain("Resolved Referent:");
}); });
describe("multi-agent room coordination (FN-5425)", () => {
async function seedMultiAgentRoom(
localHarness: Harness,
{ peerAgentId = "agent-peer", roomName }: { peerAgentId?: string; roomName: string },
): Promise<{ room: ReturnType<ChatStore["createRoom"]>; peerAgentId: string }> {
const peerAgent = await localHarness.agentStore.createAgent({
name: peerAgentId,
role: "executor",
soul: "Peer room member",
runtimeConfig: { enabled: true },
});
const room = localHarness.chatStore.createRoom({ name: roomName, memberAgentIds: [localHarness.agentId] });
localHarness.chatStore.addRoomMember(room.id, peerAgent.id);
return { room, peerAgentId: peerAgent.id };
}
it("renders claim branch and emits coordination audit in multi-agent room", async () => {
harness = await createHarness();
const { room } = await seedMultiAgentRoom(harness, { roomName: "coord-claim" });
const userMessage = harness.chatStore.addRoomMessage(room.id, {
role: "user",
content: "please file a task for the secrets-sync regression",
});
const monitor = new HeartbeatMonitor({
store: harness.agentStore,
taskStore: harness.taskStore,
rootDir: harness.rootDir,
chatStore: harness.chatStore,
});
const run = await monitor.executeHeartbeat({ agentId: harness.agentId, source: "timer" as any });
expect(sessionCapture.prompt).toContain("Room Coordination Notices:");
expect(sessionCapture.prompt).toContain("branch: claim");
expect(sessionCapture.prompt).toContain("Claiming:");
expect(sessionCapture.prompt).toContain("fn_task_create");
expect(sessionCapture.prompt).toContain("fn_post_room_message");
expect(sessionCapture.prompt).toContain("FN-4918");
const event = harness.taskStore
.getRunAuditEvents({ runId: run!.id })
.find((auditEvent) => auditEvent.mutationType === "room:coordination:branch" && auditEvent.target === userMessage.id);
expect(event?.metadata).toMatchObject({ branch: "claim" });
});
it("renders defer branch when peer already claimed", async () => {
harness = await createHarness();
const { room, peerAgentId } = await seedMultiAgentRoom(harness, { roomName: "coord-defer-claim" });
const priorClaim = harness.chatStore.addRoomMessage(room.id, {
role: "assistant",
senderAgentId: peerAgentId,
content: "Claiming: filing task for the secrets-sync regression",
});
harness.chatStore.addRoomMessage(room.id, {
role: "user",
content: "please file a task for the secrets-sync regression",
});
const monitor = new HeartbeatMonitor({
store: harness.agentStore,
taskStore: harness.taskStore,
rootDir: harness.rootDir,
chatStore: harness.chatStore,
});
const run = await monitor.executeHeartbeat({ agentId: harness.agentId, source: "timer" as any });
expect(sessionCapture.prompt).toContain("branch: defer-suggested");
expect(sessionCapture.prompt).toContain(peerAgentId);
expect(sessionCapture.prompt).toContain("Do NOT call fn_task_create");
const event = harness.taskStore
.getRunAuditEvents({ runId: run!.id })
.find((auditEvent) => auditEvent.mutationType === "room:coordination:branch");
expect(event?.metadata).toMatchObject({ branch: "defer-suggested", priorClaimMessageId: priorClaim.id });
});
it("captures prior task id from peer announcement", async () => {
harness = await createHarness();
const { room, peerAgentId } = await seedMultiAgentRoom(harness, { roomName: "coord-defer-task" });
harness.chatStore.addRoomMessage(room.id, {
role: "assistant",
senderAgentId: peerAgentId,
content: "Filed FN-9042 for the secrets-sync regression",
});
harness.chatStore.addRoomMessage(room.id, { role: "user", content: "please file a task for the secrets-sync regression" });
const monitor = new HeartbeatMonitor({
store: harness.agentStore,
taskStore: harness.taskStore,
rootDir: harness.rootDir,
chatStore: harness.chatStore,
});
const run = await monitor.executeHeartbeat({ agentId: harness.agentId, source: "timer" as any });
expect(sessionCapture.prompt).toContain("FN-9042");
const event = harness.taskStore.getRunAuditEvents({ runId: run!.id }).find((auditEvent) => auditEvent.mutationType === "room:coordination:branch");
expect(event?.metadata).toMatchObject({ priorTaskId: "FN-9042" });
});
it("does not render coordination notices for single-agent room", async () => {
harness = await createHarness();
const room = harness.chatStore.createRoom({ name: "single-agent", memberAgentIds: [harness.agentId] });
harness.chatStore.addRoomMessage(room.id, { role: "user", content: "file a task for X" });
const monitor = new HeartbeatMonitor({ store: harness.agentStore, taskStore: harness.taskStore, rootDir: harness.rootDir, chatStore: harness.chatStore });
const run = await monitor.executeHeartbeat({ agentId: harness.agentId, source: "timer" as any });
expect(sessionCapture.prompt).not.toContain("Room Coordination Notices:");
expect(harness.taskStore.getRunAuditEvents({ runId: run!.id }).some((event) => event.mutationType === "room:coordination:branch")).toBe(false);
});
it("does not render coordination notices for non-task-filing content", async () => {
harness = await createHarness();
const { room } = await seedMultiAgentRoom(harness, { roomName: "coord-non-intent" });
harness.chatStore.addRoomMessage(room.id, { role: "user", content: "what do you think about the secrets-sync regression?" });
const monitor = new HeartbeatMonitor({ store: harness.agentStore, taskStore: harness.taskStore, rootDir: harness.rootDir, chatStore: harness.chatStore });
const run = await monitor.executeHeartbeat({ agentId: harness.agentId, source: "timer" as any });
expect(sessionCapture.prompt).not.toContain("Room Coordination Notices:");
expect(harness.taskStore.getRunAuditEvents({ runId: run!.id }).some((event) => event.mutationType === "room:coordination:branch")).toBe(false);
});
it("keeps deictic-only messages in ambiguity layer only", async () => {
harness = await createHarness();
const { room } = await seedMultiAgentRoom(harness, { roomName: "coord-deictic-only" });
harness.chatStore.addRoomMessage(room.id, { role: "user", content: "we should investigate the secrets-sync regression" });
harness.chatStore.addRoomMessage(room.id, { role: "user", content: "yeah, create it" });
const monitor = new HeartbeatMonitor({ store: harness.agentStore, taskStore: harness.taskStore, rootDir: harness.rootDir, chatStore: harness.chatStore });
await monitor.executeHeartbeat({ agentId: harness.agentId, source: "timer" as any });
expect(sessionCapture.prompt).toContain("Room Ambiguity Notices:");
expect(sessionCapture.prompt).not.toContain("Room Coordination Notices:");
});
it("does not defer to a self-authored prior claim", async () => {
harness = await createHarness();
const { room } = await seedMultiAgentRoom(harness, { roomName: "coord-self-claim" });
harness.chatStore.addRoomMessage(room.id, {
role: "assistant",
senderAgentId: harness.agentId,
content: "Claiming: filing task for the secrets-sync regression",
});
harness.chatStore.addRoomMessage(room.id, { role: "user", content: "please file a task for the secrets-sync regression" });
const monitor = new HeartbeatMonitor({ store: harness.agentStore, taskStore: harness.taskStore, rootDir: harness.rootDir, chatStore: harness.chatStore });
await monitor.executeHeartbeat({ agentId: harness.agentId, source: "timer" as any });
expect(sessionCapture.prompt).toContain("branch: claim");
expect(sessionCapture.prompt).not.toContain("branch: defer-suggested");
});
it("fails open when coordination helper throws", async () => {
harness = await createHarness();
const { room } = await seedMultiAgentRoom(harness, { roomName: "coord-throw" });
harness.chatStore.addRoomMessage(room.id, { role: "user", content: "please file a task for the secrets-sync regression" });
const spy = vi.spyOn(roomCoordination, "decideRoomCoordination").mockImplementation(() => {
throw new Error("boom");
});
const monitor = new HeartbeatMonitor({ store: harness.agentStore, taskStore: harness.taskStore, rootDir: harness.rootDir, chatStore: harness.chatStore });
await expect(monitor.executeHeartbeat({ agentId: harness.agentId, source: "timer" as any })).resolves.toBeTruthy();
expect(sessionCapture.prompt).not.toContain("Room Coordination Notices:");
spy.mockRestore();
});
});
it("registers fn_post_room_message and posts through the real ChatStore under restrictive policy", async () => { it("registers fn_post_room_message and posts through the real ChatStore under restrictive policy", async () => {
harness = await createHarness({ harness = await createHarness({
presetId: "approval-required", presetId: "approval-required",

View File

@@ -0,0 +1,217 @@
import type { ChatRoomMember, ChatRoomMessage } from "@fusion/core";
import { describe, expect, it } from "vitest";
import {
countActiveAgentMembers,
decideRoomCoordination,
detectTaskFilingIntent,
renderRoomCoordinationPromptBlock,
} from "../room-coordination.js";
function roomMember(agentId: string | null): ChatRoomMember {
return {
roomId: "room-1",
// cast for test-only simulation of nullable rows from external callers
agentId: agentId as unknown as string,
role: "member",
addedAt: new Date().toISOString(),
};
}
function roomMessage(id: string, content: string, senderAgentId: string | null = "agent-peer"): ChatRoomMessage {
return {
id,
roomId: "room-1",
role: senderAgentId ? "assistant" : "user",
content,
thinkingOutput: null,
metadata: null,
senderAgentId,
mentions: [],
createdAt: new Date().toISOString(),
};
}
describe("room-coordination", () => {
describe("detectTaskFilingIntent", () => {
it("detects explicit task-filing intent with extracted subject", () => {
const result = detectTaskFilingIntent("please file a task for the secrets-sync regression");
expect(result.isTaskFilingIntent).toBe(true);
expect(result.subject).toBe("the secrets-sync regression");
expect(result.cues.length).toBeGreaterThan(0);
});
it.each([
["Can you create a task to fix the broken merge?", "fix the broken merge"],
["open a task about the new typecheck failure", "the new typecheck failure"],
["track this as a task", null],
["FILE A TASK: dashboard FAB regression.", "dashboard FAB regression"],
["create a task for it", "it"],
["create a task for FN-1234", "FN-1234"],
])("handles positive variants: %s", (content, expectedSubject) => {
const result = detectTaskFilingIntent(content);
expect(result.isTaskFilingIntent).toBe(true);
if (expectedSubject === null) {
expect([null, "this"]).toContain(result.subject);
} else {
expect(result.subject).toContain(expectedSubject);
}
});
it.each([
"file it",
"create it now",
"do that as a task",
"this is a task list",
"I filed the task report yesterday",
"",
" ",
])("rejects non-intent content: %s", (content) => {
const result = detectTaskFilingIntent(content);
expect(result).toEqual({ isTaskFilingIntent: false, cues: [], subject: null });
});
it("returns true for past-tense phrasing trade-off", () => {
const result = detectTaskFilingIntent("I filed a task earlier");
expect(result.isTaskFilingIntent).toBe(true);
});
it("rejects oversized content", () => {
const result = detectTaskFilingIntent(`create a task for ${"x".repeat(1000)}`);
expect(result).toEqual({ isTaskFilingIntent: false, cues: [], subject: null });
});
});
describe("countActiveAgentMembers", () => {
it("counts unique active agent members", () => {
expect(countActiveAgentMembers([roomMember("a1"), roomMember("a2"), roomMember("a3")])).toBe(3);
});
it("deduplicates duplicate agent ids", () => {
expect(countActiveAgentMembers([roomMember("a1"), roomMember("a1"), roomMember("a2")])).toBe(2);
});
it("ignores non-agent members", () => {
expect(countActiveAgentMembers([roomMember(null), roomMember("a1")])).toBe(1);
});
});
describe("decideRoomCoordination", () => {
const detection = detectTaskFilingIntent("please file a task for secrets sync");
it("returns null for non-intent detection", () => {
const result = decideRoomCoordination({
detection: { isTaskFilingIntent: false, cues: [], subject: null },
members: [roomMember("a1"), roomMember("a2")],
recentMessages: [],
pendingSenderAgentId: null,
});
expect(result).toBeNull();
});
it("returns null for single-agent rooms", () => {
const result = decideRoomCoordination({
detection,
members: [roomMember("a1")],
recentMessages: [],
pendingSenderAgentId: null,
});
expect(result).toBeNull();
});
it("returns claim when no prior peer claim exists", () => {
const result = decideRoomCoordination({
detection,
members: [roomMember("a1"), roomMember("a2")],
recentMessages: [roomMessage("m1", "hello", null)],
pendingSenderAgentId: null,
});
expect(result?.branch).toBe("claim");
expect(result?.priorClaimMessageId).toBeUndefined();
});
it("returns defer-suggested with prior peer claim", () => {
const result = decideRoomCoordination({
detection,
members: [roomMember("a1"), roomMember("a2")],
recentMessages: [roomMessage("m1", "Claiming: filing task for secrets sync", "agent-peer")],
pendingSenderAgentId: "agent-main",
});
expect(result?.branch).toBe("defer-suggested");
expect(result?.priorClaimMessageId).toBe("m1");
expect(result?.priorClaimSenderId).toBe("agent-peer");
});
it("returns defer-suggested with prior task announcement", () => {
const result = decideRoomCoordination({
detection,
members: [roomMember("a1"), roomMember("a2")],
recentMessages: [roomMessage("m2", "Filed FN-9001 for the secrets-sync regression", "agent-peer")],
pendingSenderAgentId: "agent-main",
});
expect(result?.branch).toBe("defer-suggested");
expect(result?.priorTaskId).toBe("FN-9001");
});
it("does not defer to self-authored claim", () => {
const result = decideRoomCoordination({
detection,
members: [roomMember("a1"), roomMember("a2")],
recentMessages: [roomMessage("m3", "Claiming: filing task for secrets sync", "agent-main")],
pendingSenderAgentId: "agent-main",
});
expect(result?.branch).toBe("claim");
});
it("skips pending message itself from prior lookup", () => {
const result = decideRoomCoordination({
detection,
members: [roomMember("a1"), roomMember("a2")],
recentMessages: [roomMessage("pending", "Claiming: filing task for secrets sync", "agent-main")],
pendingSenderAgentId: "agent-main",
});
expect(result?.branch).toBe("claim");
});
});
describe("renderRoomCoordinationPromptBlock", () => {
it("renders claim branch instructions with explicit guard references", () => {
const lines = renderRoomCoordinationPromptBlock(
{
branch: "claim",
memberCount: 3,
detection: { isTaskFilingIntent: true, cues: ["file a task"], subject: "secrets-sync regression" },
},
{ id: "m5" },
);
const block = lines.join("\n");
expect(block).toContain("fn_post_room_message");
expect(block).toContain("fn_task_create");
expect(block).toContain("FN-4918");
expect(block).toContain("FN-4829");
expect(block).toContain("FN-5152");
expect(block).toContain("FN-5220");
expect(block).toContain("reply_to_message_id = m5");
});
it("renders defer branch with peer, message, and prior task id", () => {
const lines = renderRoomCoordinationPromptBlock(
{
branch: "defer-suggested",
memberCount: 2,
detection: { isTaskFilingIntent: true, cues: ["create a task"], subject: "secrets-sync regression" },
priorClaimMessageId: "m4",
priorClaimSenderId: "agent-peer",
priorTaskId: "FN-9042",
},
{ id: "m6" },
);
const block = lines.join("\n");
expect(block).toContain("agent-peer");
expect(block).toContain("message m4");
expect(block).toContain("FN-9042");
expect(block).toContain("reply_to_message_id = m6");
});
});
});

View File

@@ -42,6 +42,7 @@ import { buildSessionSkillContextSync } from "./session-skill-context.js";
import type { AgentReflectionService } from "./agent-reflection.js"; import type { AgentReflectionService } from "./agent-reflection.js";
import { trimPromptMd, trimTaskDescription, trimTriggeringComments } from "./heartbeat-prompt-trim.js"; import { trimPromptMd, trimTaskDescription, trimTriggeringComments } from "./heartbeat-prompt-trim.js";
import { detectDeicticReference, extractAntecedentCandidates, renderAmbiguityPromptBlock, scoreReferentConfidence } from "./room-ambiguity.js"; import { detectDeicticReference, extractAntecedentCandidates, renderAmbiguityPromptBlock, scoreReferentConfidence } from "./room-ambiguity.js";
import { countActiveAgentMembers, decideRoomCoordination, detectTaskFilingIntent, renderRoomCoordinationPromptBlock } from "./room-coordination.js";
const promptSizeLog = createLogger("prompt-size"); const promptSizeLog = createLogger("prompt-size");
@@ -400,6 +401,7 @@ When you are woken by an incoming message (source includes "wake-on-message"), y
4. If a Pending Room Messages section is present, review it too: 4. If a Pending Room Messages section is present, review it too:
- Use fn_post_room_message only when the room content is relevant to your role, soul, or identity. - Use fn_post_room_message only when the room content is relevant to your role, soul, or identity.
- If a Room Ambiguity Notices section is present, follow it exactly: echo resolved referents before acting, and under clarification notices do not create tasks. - If a Room Ambiguity Notices section is present, follow it exactly: echo resolved referents before acting, and under clarification notices do not create tasks.
- If a Room Coordination Notices section is present, follow its claim/defer branch exactly: under "claim" post a one-line claim before calling fn_task_create; under "defer-suggested" do NOT call fn_task_create and instead acknowledge the prior claim via fn_post_room_message.
- Reference room message IDs when replying so humans can trace context. - Reference room message IDs when replying so humans can trace context.
5. After processing messages, continue with your normal heartbeat duties. 5. After processing messages, continue with your normal heartbeat duties.
@@ -486,7 +488,7 @@ When you are woken by an incoming message (source includes "wake-on-message"), y
- If the message is informational, acknowledge it and respond via fn_send_message when appropriate. - If the message is informational, acknowledge it and respond via fn_send_message when appropriate.
- If the message requests work, create a follow-up task with fn_task_create. - If the message requests work, create a follow-up task with fn_task_create.
- If the request has a clear owner and fn_delegate_task is available, delegate it directly. - If the request has a clear owner and fn_delegate_task is available, delegate it directly.
3. If a Pending Room Messages section is present, review it too and use fn_post_room_message only when the room content is relevant to your role or identity; if Room Ambiguity Notices are present, follow their resolve/clarify branch instructions exactly. 3. If a Pending Room Messages section is present, review it too and use fn_post_room_message only when the room content is relevant to your role or identity; if Room Ambiguity Notices are present, follow their resolve/clarify branch instructions exactly. If a Room Coordination Notices section is present, follow its claim/defer branch exactly: under "claim" post a one-line claim before calling fn_task_create; under "defer-suggested" do NOT call fn_task_create and instead acknowledge the prior claim via fn_post_room_message.
4. After processing messages, continue with your ambient work. 4. After processing messages, continue with your ambient work.
Example flow: Example flow:
@@ -518,7 +520,11 @@ export const HEARTBEAT_PROCEDURE_STRICT = `## Heartbeat Procedure (run every tic
reply_to_message_id when answering. If Pending Room Messages are present, reply_to_message_id when answering. If Pending Room Messages are present,
review them in the prompt and use fn_post_room_message only when relevant. review them in the prompt and use fn_post_room_message only when relevant.
When Room Ambiguity Notices appear, follow the resolve/clarify branch and do When Room Ambiguity Notices appear, follow the resolve/clarify branch and do
not create tasks under clarification notices. not create tasks under clarification notices. If a Room Coordination Notices
section is present, follow its claim/defer branch exactly: under "claim" post
a one-line claim before calling fn_task_create; under "defer-suggested" do
NOT call fn_task_create and instead acknowledge the prior claim via
fn_post_room_message.
3. **Wake delta** — read the Wake Delta block above. The wake reason is the 3. **Wake delta** — read the Wake Delta block above. The wake reason is the
highest-priority change for this heartbeat. If you were woken by a comment highest-priority change for this heartbeat. If you were woken by a comment
or a message, acknowledge it before doing anything else. or a message, acknowledge it before doing anything else.
@@ -618,7 +624,11 @@ export const HEARTBEAT_NO_TASK_PROCEDURE_STRICT = `## Heartbeat Procedure (run e
reply_to_message_id when answering. If Pending Room Messages are present, reply_to_message_id when answering. If Pending Room Messages are present,
review them in the prompt and use fn_post_room_message only when relevant. review them in the prompt and use fn_post_room_message only when relevant.
When Room Ambiguity Notices appear, follow the resolve/clarify branch and do When Room Ambiguity Notices appear, follow the resolve/clarify branch and do
not create tasks under clarification notices. not create tasks under clarification notices. If a Room Coordination Notices
section is present, follow its claim/defer branch exactly: under "claim" post
a one-line claim before calling fn_task_create; under "defer-suggested" do
NOT call fn_task_create and instead acknowledge the prior claim via
fn_post_room_message.
3. **Wake delta** — read the Wake Delta block above. The wake reason is the 3. **Wake delta** — read the Wake Delta block above. The wake reason is the
highest-priority change for this heartbeat. If you were woken by a comment highest-priority change for this heartbeat. If you were woken by a comment
or a message, acknowledge it before doing anything else. or a message, acknowledge it before doing anything else.
@@ -954,6 +964,86 @@ export class HeartbeatMonitor {
return lines; return lines;
} }
private async getRoomCoordinationNoticesSection(
agent: Agent,
runId: string,
entries: Array<{ room: ChatRoom; messages: ChatRoomMessage[] }>,
audit: ReturnType<typeof createRunAuditor>,
): Promise<string[]> {
if (!this.chatStore || entries.length === 0) {
return [];
}
const lines: string[] = [];
for (const entry of entries) {
for (const message of entry.messages) {
try {
const detection = detectTaskFilingIntent(message.content);
if (!detection.isTaskFilingIntent) {
continue;
}
const members = this.chatStore.listRoomMembers(entry.room.id);
if (countActiveAgentMembers(members) < 2) {
continue;
}
const roomTimeline = this.chatStore.getRoomMessages(entry.room.id, { limit: 100 });
const messageIndex = roomTimeline.findIndex((roomMessage) => roomMessage.id === message.id);
if (messageIndex < 0) {
continue;
}
// messageIndex === 0 means no prior messages; coordination correctly defaults to claim.
const recentMessages = messageIndex === 0
? []
: roomTimeline.slice(Math.max(0, messageIndex - 15), messageIndex);
const decision = decideRoomCoordination({
detection,
members,
recentMessages,
pendingSenderAgentId: message.senderAgentId ?? agent.id,
});
if (!decision) {
continue;
}
if (lines.length === 0) {
lines.push("", "Room Coordination Notices:");
}
lines.push(
`- [room: ${entry.room.name} (${entry.room.id})] [message: ${message.id}] [branch: ${decision.branch}]`,
...renderRoomCoordinationPromptBlock(decision, message).map((line) => ` - ${line}`),
);
await audit.database({
type: "room:coordination:branch",
target: message.id,
metadata: {
roomId: entry.room.id,
agentId: agent.id,
branch: decision.branch,
memberCount: decision.memberCount,
intentCue: detection.cues[0] ?? null,
priorClaimMessageId: decision.priorClaimMessageId ?? null,
priorTaskId: decision.priorTaskId ?? null,
},
});
heartbeatLog.log(
`[room-coordination] agent=${agent.id} run=${runId} room=${entry.room.id} messageId=${message.id} branch=${decision.branch} members=${decision.memberCount}`,
);
} catch (err) {
heartbeatLog.warn(`Room coordination notice failed for ${message.id}: ${err instanceof Error ? err.message : String(err)}`);
}
}
}
return lines;
}
private async getPendingRoomMessages(agent: Agent, sinceIso: string): Promise<{ private async getPendingRoomMessages(agent: Agent, sinceIso: string): Promise<{
entries: Array<{ room: ChatRoom; messages: ChatRoomMessage[] }>; entries: Array<{ room: ChatRoom; messages: ChatRoomMessage[] }>;
total: number; total: number;
@@ -2493,6 +2583,12 @@ export class HeartbeatMonitor {
pendingRoomMessages.entries, pendingRoomMessages.entries,
audit, audit,
); );
const roomCoordinationNoticesLines = await this.getRoomCoordinationNoticesSection(
agent,
run.id,
pendingRoomMessages.entries,
audit,
);
// Fetch unread messages when messageStore is available (for all trigger types) // Fetch unread messages when messageStore is available (for all trigger types)
if (this.messageStore) { if (this.messageStore) {
@@ -2679,6 +2775,7 @@ export class HeartbeatMonitor {
...pendingMessagesLines, ...pendingMessagesLines,
...pendingRoomMessagesLines, ...pendingRoomMessagesLines,
...roomAmbiguityNoticesLines, ...roomAmbiguityNoticesLines,
...roomCoordinationNoticesLines,
"", "",
"Your soul, instructions, and memory are already loaded in the system prompt.", "Your soul, instructions, and memory are already loaded in the system prompt.",
"Focus on work that benefits the project without requiring a specific task context.", "Focus on work that benefits the project without requiring a specific task context.",
@@ -2767,6 +2864,7 @@ export class HeartbeatMonitor {
...pendingMessagesLines, ...pendingMessagesLines,
...pendingRoomMessagesLines, ...pendingRoomMessagesLines,
...roomAmbiguityNoticesLines, ...roomAmbiguityNoticesLines,
...roomCoordinationNoticesLines,
...(reportsHealthSection ? ["", reportsHealthSection] : []), ...(reportsHealthSection ? ["", reportsHealthSection] : []),
"", "",
"Run the Heartbeat Procedure above. Call fn_heartbeat_done when finished.", "Run the Heartbeat Procedure above. Call fn_heartbeat_done when finished.",

View File

@@ -0,0 +1,205 @@
import type { ChatRoomMember, ChatRoomMessage } from "@fusion/core";
/**
* FN-5425: prompt-layer multi-agent coordination advisory for explicit task-filing room messages.
* This module stays deterministic and fail-open: it only detects intent + suggests claim/defer
* prompt branches, while authoritative duplicate prevention remains in intake safeguards.
*/
export interface TaskFilingIntentDetection {
isTaskFilingIntent: boolean;
cues: string[];
subject: string | null;
}
export interface RoomCoordinationDecision {
branch: "claim" | "defer-suggested";
memberCount: number;
detection: TaskFilingIntentDetection;
priorClaimMessageId?: string;
priorClaimSenderId?: string | null;
priorTaskId?: string;
}
export interface RenderCoordinationOptions {
deicticMessageId?: never;
pendingMessageId: string;
}
const MAX_INTENT_CONTENT_LENGTH = 800;
const MAX_SUBJECT_LENGTH = 140;
const DEICTIC_NOUN_FOLLOWUP_PATTERN = /\b(?:it|that|this)\s+(?:as|for|to)\b/i;
const TASK_ID_PATTERN = /\b(FN-\d{1,6})\b/i;
const PRIOR_CLAIM_RE = /^\s*claiming[:-]/i;
const TASK_ANNOUNCED_RE = /\b(?:filed|created|opened|tracked|added)\b[\s\S]{0,60}?\b(FN-\d{1,6})\b|\b(FN-\d{1,6})\b[\s\S]{0,60}?\b(?:filed|created|opened|tracked|added)\b/i;
const TASK_CONTEXT_SUFFIX = "(?=$|[.?!,:;-]|\\s+(?:for|about|to|on|regarding)\\b)";
const TASK_FILING_CUES: ReadonlyArray<[string, RegExp]> = [
["file a task", new RegExp(`\\bfile\\w*\\s+(?:a\\s+|the\\s+)?task${TASK_CONTEXT_SUFFIX}`, "i")],
["filed a task", /\bfiled\s+(?:a\s+|the\s+)?task\b/i],
["create a task", new RegExp(`\\bcreate\\s+(?:a\\s+|the\\s+)?task${TASK_CONTEXT_SUFFIX}`, "i")],
["open a task", new RegExp(`\\bopen\\s+(?:a\\s+|the\\s+)?task${TASK_CONTEXT_SUFFIX}`, "i")],
["add a task", new RegExp(`\\badd\\s+(?:a\\s+|the\\s+)?task${TASK_CONTEXT_SUFFIX}`, "i")],
["track this/that as a task", new RegExp(`\\btrack\\s+(?:this|that)\\s+(?:as\\s+)?(?:a\\s+)?task${TASK_CONTEXT_SUFFIX}`, "i")],
["start a task", new RegExp(`\\bstart\\s+(?:a\\s+|the\\s+)?task${TASK_CONTEXT_SUFFIX}`, "i")],
["make a task", new RegExp(`\\bmake\\s+(?:a\\s+|the\\s+)?task${TASK_CONTEXT_SUFFIX}`, "i")],
];
function normalizeMessageContent(content: string): string {
return content.replace(/\s+/g, " ").trim();
}
function truncateWithEllipsis(value: string, maxLength = MAX_SUBJECT_LENGTH): string {
if (value.length <= maxLength) {
return value;
}
return `${value.slice(0, maxLength - 1).trimEnd()}…`;
}
function cleanExtractedSubject(subject: string): string {
return truncateWithEllipsis(
subject
.replace(/^\s*["'`]+|["'`]+\s*$/g, "")
.replace(/\s+/g, " ")
.trim(),
);
}
function extractSubject(normalized: string): string | null {
const preferred = normalized.match(
/(?:file\w*|create|open|add|track|start|make)\s+(?:a\s+|the\s+)?task\s+(?:for|about|to|on|regarding)\s+(.+?)(?:[.?!]|$)/i,
);
if (preferred?.[1]) {
const cleaned = cleanExtractedSubject(preferred[1]);
return cleaned || null;
}
const fallback = normalized.match(
/(?:file\w*|create|open|add|track|start|make)\s+(?:a\s+|the\s+)?task(?:\s+|\s*[:-]\s*)(.+?)(?:[.?!]|$)/i,
);
if (fallback?.[1]) {
const cleaned = cleanExtractedSubject(fallback[1]);
return cleaned || null;
}
return null;
}
export function detectTaskFilingIntent(content: string): TaskFilingIntentDetection {
const normalized = normalizeMessageContent(content);
if (!normalized || normalized.length > MAX_INTENT_CONTENT_LENGTH) {
return { isTaskFilingIntent: false, cues: [], subject: null };
}
const cueMatches: Array<{ cue: string; index: number }> = [];
for (const [cue, pattern] of TASK_FILING_CUES) {
const match = pattern.exec(normalized);
if (match?.index != null) {
cueMatches.push({ cue, index: match.index });
}
}
if (cueMatches.length === 0) {
return { isTaskFilingIntent: false, cues: [], subject: null };
}
if (/\bfiled\s+(?:a\s+|the\s+)?task\s+report\b/i.test(normalized)) {
return { isTaskFilingIntent: false, cues: [], subject: null };
}
const earliestCueIndex = cueMatches.reduce((min, current) => Math.min(min, current.index), Number.POSITIVE_INFINITY);
const deicticMatch = DEICTIC_NOUN_FOLLOWUP_PATTERN.exec(normalized);
const deicticIndex = deicticMatch?.index ?? -1;
if (deicticIndex >= 0 && deicticIndex < earliestCueIndex) {
return { isTaskFilingIntent: false, cues: [], subject: null };
}
const subject = extractSubject(normalized);
// Intentional trade-off: the advisory is harmless under occasional tense false positives,
// and FN-5152/FN-5220 intake dedup remains the deterministic backstop.
return {
isTaskFilingIntent: true,
cues: cueMatches.sort((a, b) => a.index - b.index).map((entry) => entry.cue),
subject,
};
}
export function countActiveAgentMembers(members: ChatRoomMember[]): number {
// ChatStore.removeRoomMember hard-deletes rows, so listRoomMembers is already active-only.
const uniqueAgentIds = new Set(members.map((member) => member.agentId).filter((agentId): agentId is string => Boolean(agentId)));
return uniqueAgentIds.size;
}
export function decideRoomCoordination(args: {
detection: TaskFilingIntentDetection;
members: ChatRoomMember[];
recentMessages: ChatRoomMessage[];
pendingSenderAgentId: string | null;
}): RoomCoordinationDecision | null {
const { detection, members, recentMessages, pendingSenderAgentId } = args;
if (!detection.isTaskFilingIntent) {
return null;
}
const memberCount = countActiveAgentMembers(members);
if (memberCount < 2) {
return null;
}
const lookback = recentMessages.slice(-15);
let priorClaimMessageId: string | undefined;
let priorClaimSenderId: string | null | undefined;
let priorTaskId: string | undefined;
for (const message of lookback) {
if (!message.senderAgentId) {
continue;
}
if (pendingSenderAgentId && message.senderAgentId === pendingSenderAgentId) {
continue;
}
const claimMatch = PRIOR_CLAIM_RE.test(message.content);
const announcedMatch = TASK_ANNOUNCED_RE.exec(message.content);
if (!claimMatch && !announcedMatch) {
continue;
}
priorClaimMessageId = message.id;
priorClaimSenderId = message.senderAgentId;
const captured = announcedMatch?.[1] ?? announcedMatch?.[2] ?? message.content.match(TASK_ID_PATTERN)?.[1];
priorTaskId = captured?.toUpperCase();
break;
}
return {
branch: priorClaimMessageId ? "defer-suggested" : "claim",
memberCount,
detection,
priorClaimMessageId,
priorClaimSenderId,
priorTaskId,
};
}
export function renderRoomCoordinationPromptBlock(
decision: RoomCoordinationDecision,
pendingMessage: Pick<ChatRoomMessage, "id">,
): string[] {
if (decision.branch === "claim") {
return [
`Multi-agent room (${decision.memberCount} agents). Before calling fn_task_create for this request, coordinate:`,
`1. Post a ONE-LINE claim to the room first via fn_post_room_message: "Claiming: filing task for <subject>" (reply_to_message_id = ${pendingMessage.id}).`,
"2. Then call fn_task_create. The deterministic / near-duplicate / explicit-marker intake guards (FN-4918 / FN-4829 / FN-5152 / FN-5220) are your authoritative backstop — do NOT pass acknowledgedDuplicates or bypassDuplicateCheck to silence a duplicate match for a room request.",
`3. After fn_task_create returns, post the resulting FN-NNNN id back to the room via fn_post_room_message (reply_to_message_id = ${pendingMessage.id}).`,
`Detected subject: ${decision.detection.subject ?? "(none extracted — restate the user's request in your claim)"}.`,
];
}
return [
`Multi-agent room (${decision.memberCount} agents). A peer agent (${decision.priorClaimSenderId ?? "unknown"}) already posted a claim or task announcement in this room (message ${decision.priorClaimMessageId}${decision.priorTaskId ? `, task ${decision.priorTaskId}` : ""}).`,
`Do NOT call fn_task_create for this request. Reply once via fn_post_room_message (reply_to_message_id = ${pendingMessage.id}) acknowledging the existing claim${decision.priorTaskId ? ` and echoing ${decision.priorTaskId}` : ""}.`,
"If you believe the peer's task does NOT cover this request, say so explicitly in your reply and wait for human disambiguation rather than filing in parallel.",
];
}

View File

@@ -307,7 +307,8 @@ export type DatabaseMutationType =
* ``` * ```
*/ */
| "worktree:pool-double-lease-detected" | "worktree:pool-double-lease-detected"
| "room:ambiguity:branch"; | "room:ambiguity:branch"
| "room:coordination:branch";
// ── Filesystem mutation types ───────────────────────────────────────────────── // ── Filesystem mutation types ─────────────────────────────────────────────────