fix(FN-8277): prevent replayed duplicate follow-up tasks (#2319)
## Summary Repeated review and executor sessions could replay a follow-up creation step and produce another live task whenever the wording changed. In the incident behind this fix, 21 creation calls for three intended follow-ups left 18 duplicate tasks. Agent-created tasks now retain their parent and agent provenance across step sessions, heartbeats, and the published CLI surface. Same-parent paraphrases converge on the existing task through a serialized pre-check and a database-backed intent claim, while distinct sibling actions remain separate. Candidate lookup is parent-indexed, uniqueness failures abort creation, and canonical reuse no longer emits misleading creation audit events or workflow claims. Related: FN-8277 ## Validation - Core duplicate guard and intake: 24 tests passed - Engine task creation and heartbeat: 154 tests passed - CLI extension: 66 tests passed, 95 skipped - Core, engine, and CLI typechecks passed - Core, engine, and CLI builds passed before rebase; the rebase was conflict-free - Scoped lint, strict changeset validation, and diff checks passed <!-- This is an auto-generated comment: release notes by coderabbit.ai --> ## Summary by CodeRabbit * **Bug Fixes** * Prevented retried agent steps from creating duplicate follow-up tasks. * Improved parent-scoped deduplication for paraphrased follow-ups while preserving distinct sibling actions. * Preserved parent-task and agent context for created follow-ups. * Concurrent follow-up requests are now serialized/deduplicated so duplicates link to the existing task instead of showing as newly created. * Updated agent follow-up/heartbeat activity so reused follow-ups no longer appear in run results as fresh creations. <!-- end of auto-generated comment: release notes by coderabbit.ai -->
This commit is contained in:
7
.changeset/calm-tasks-reuse-parent-intent.md
Normal file
7
.changeset/calm-tasks-reuse-parent-intent.md
Normal file
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@runfusion/fusion": patch
|
||||
---
|
||||
|
||||
summary: Prevent retried agent steps from creating duplicate follow-up tasks.
|
||||
category: fix
|
||||
dev: Persists parent provenance and serializes parent-scoped intent deduplication across task-create surfaces.
|
||||
@@ -78,6 +78,7 @@ vi.mock("@fusion/engine", () => ({
|
||||
installBaselineArchiveWorktreeDisposer: vi.fn(),
|
||||
...workflowAuthoringEngineMock,
|
||||
createFnAgent: vi.fn(),
|
||||
createAgentTask: vi.fn(),
|
||||
fetchWebContent: vi.fn(),
|
||||
emitGoalRetrievalAudit: vi.fn(),
|
||||
createWorkflowAuthoringTools: vi.fn(() => ({})),
|
||||
|
||||
@@ -38,6 +38,7 @@ vi.mock("@fusion/dashboard", () => {
|
||||
vi.mock("@fusion/engine", () => ({
|
||||
installBaselineArchiveWorktreeDisposer: vi.fn(),
|
||||
createFnAgent: vi.fn(),
|
||||
createAgentTask: vi.fn(),
|
||||
fetchWebContent: vi.fn(),
|
||||
assertNoSecretPlaintext: vi.fn(),
|
||||
emitGoalRetrievalAudit: vi.fn(),
|
||||
|
||||
@@ -17,6 +17,7 @@ vi.mock("@fusion/engine", () => ({
|
||||
installBaselineArchiveWorktreeDisposer: vi.fn(),
|
||||
...workflowAuthoringEngineMock,
|
||||
createFnAgent: vi.fn(),
|
||||
createAgentTask: vi.fn(),
|
||||
fetchWebContent: fetchWebContentMock,
|
||||
emitGoalRetrievalAudit: vi.fn(),
|
||||
createWorkflowAuthoringTools: vi.fn(() => ({})),
|
||||
|
||||
@@ -362,6 +362,21 @@ legacyDescribe("fn pi extension (legacy exhaustive suite)", () => {
|
||||
expect(result.details.priority).toBe("normal");
|
||||
});
|
||||
|
||||
it("persists ambient task provenance for agent-created follow-ups", async () => {
|
||||
const tool = api.tools.get("fn_task_create")!;
|
||||
const result = await tool.execute("call-parented", { description: "Capture optional report screenshots" }, undefined, undefined,
|
||||
{ cwd: tmpDir, taskId: "FN-PARENT", agentId: "agent-worker" } as ToolExecuteContext);
|
||||
const task = await h.store().getTask(result.details.taskId);
|
||||
expect(task.sourceParentTaskId).toBe("FN-PARENT");
|
||||
expect(task.sourceAgentId).toBe("agent-worker");
|
||||
const replay = await tool.execute("call-parented-replay", {
|
||||
description: "Capture optional report screenshots with privacy context",
|
||||
}, undefined, undefined, { cwd: tmpDir, taskId: "FN-PARENT", agentId: "agent-worker" } as ToolExecuteContext);
|
||||
expect(replay.details.taskId).toBe(result.details.taskId);
|
||||
expect(replay.details.wasDuplicate).toBe(true);
|
||||
expect(replay.content[0].text).toContain("Linked existing");
|
||||
});
|
||||
|
||||
it("creates a task with explicit priority", async () => {
|
||||
const tool = api.tools.get("fn_task_create")!;
|
||||
const result = await tool.execute(
|
||||
|
||||
@@ -69,6 +69,7 @@ import {
|
||||
isInReviewMissingWorktreeSessionStartFailure,
|
||||
normalizeAgentLogPaging,
|
||||
renderAgentLogEntries,
|
||||
createAgentTask,
|
||||
} from "@fusion/engine";
|
||||
import * as dashboard from "@fusion/dashboard";
|
||||
import { resolve, relative, isAbsolute, sep, basename, extname, join } from "node:path";
|
||||
@@ -1164,7 +1165,7 @@ export default function kbExtension(pi: ExtensionAPI) {
|
||||
FNXC:EphemeralAgentTaskCreation 2026-07-01-00:00:
|
||||
Gate ephemeral task-worker callers behind the project `ephemeralAgentsCanCreateTasks` toggle (default true). Only runtime-managed agents are affected; humans/CLI/dashboard calls carry no ctx.agentId and pass through.
|
||||
*/
|
||||
const fnCtx = ctx as typeof ctx & { agentId?: string };
|
||||
const fnCtx = ctx as typeof ctx & { agentId?: string; taskId?: string };
|
||||
const projectSettingsForGate = await store.getSettings();
|
||||
const callerIsEphemeral = await isEphemeralCallerAgent(ctx.cwd ?? process.cwd(), fnCtx.agentId);
|
||||
if (callerIsEphemeral) {
|
||||
@@ -1213,13 +1214,13 @@ export default function kbExtension(pi: ExtensionAPI) {
|
||||
);
|
||||
const workflowId = params.workflow_id?.trim() || undefined;
|
||||
|
||||
const task = await store.createTask({
|
||||
const { task, wasDuplicate } = await createAgentTask(store, {
|
||||
description: params.description.trim(),
|
||||
dependencies: params.depends,
|
||||
assignedAgentId: normalizedAgentId === null ? undefined : normalizedAgentId,
|
||||
priority: params.priority as TaskPriority | undefined,
|
||||
...(workflowId ? { workflowId } : {}),
|
||||
source: { sourceType: "api" },
|
||||
source: { sourceType: "api", sourceAgentId: fnCtx.agentId, sourceParentTaskId: fnCtx.taskId },
|
||||
githubTracking: resolvedTracking.enabled
|
||||
? {
|
||||
enabled: true,
|
||||
@@ -1228,7 +1229,7 @@ export default function kbExtension(pi: ExtensionAPI) {
|
||||
: {}),
|
||||
}
|
||||
: undefined,
|
||||
});
|
||||
}, { rootDir: ctx.cwd, sourceAgentId: fnCtx.agentId, sourceTaskId: fnCtx.taskId });
|
||||
|
||||
const label =
|
||||
task.description.length > 80
|
||||
@@ -1248,7 +1249,7 @@ export default function kbExtension(pi: ExtensionAPI) {
|
||||
{
|
||||
type: "text",
|
||||
text:
|
||||
`Created ${task.id}: ${label}${workflowId ? ` (workflow: ${workflowId})` : ""}\n` +
|
||||
`${wasDuplicate ? "Linked existing" : "Created"} ${task.id}: ${label}${workflowId && !wasDuplicate ? ` (workflow: ${workflowId})` : ""}\n` +
|
||||
`Column: ${task.column}\n` +
|
||||
(task.dependencies.length
|
||||
? `Dependencies: ${task.dependencies.join(", ")}\n`
|
||||
@@ -1262,6 +1263,7 @@ export default function kbExtension(pi: ExtensionAPI) {
|
||||
],
|
||||
details: {
|
||||
taskId: task.id,
|
||||
wasDuplicate,
|
||||
column: task.column,
|
||||
dependencies: task.dependencies,
|
||||
assignedAgentId: task.assignedAgentId,
|
||||
|
||||
@@ -122,6 +122,25 @@ describe("runDeterministicDuplicateGuard", () => {
|
||||
secondResult.releaseLock();
|
||||
});
|
||||
|
||||
it("queues three different fingerprints behind one serialization key", async () => {
|
||||
const { store } = makeStore();
|
||||
const calls = [INPUT, { title: "Second", description: "second description" }, { title: "Third", description: "third description" }];
|
||||
const first = await runDeterministicDuplicateGuard(store, calls[0]!, { lockScope: "p-1", serializationKey: "parent:FN-8277" });
|
||||
const order: number[] = [];
|
||||
const waiting = calls.slice(1).map((input, index) => runDeterministicDuplicateGuard(store, input, {
|
||||
lockScope: "p-1", serializationKey: "parent:FN-8277",
|
||||
}).then((result) => { order.push(index + 2); return result; }));
|
||||
await Promise.resolve();
|
||||
expect(order).toEqual([]);
|
||||
first.releaseLock();
|
||||
const second = await waiting[0]!;
|
||||
expect(order).toEqual([2]);
|
||||
second.releaseLock();
|
||||
const third = await waiting[1]!;
|
||||
expect(order).toEqual([2, 3]);
|
||||
third.releaseLock();
|
||||
});
|
||||
|
||||
it("allows concurrent proceed without shared lock scope", async () => {
|
||||
const { store } = makeStore();
|
||||
const [a, b] = await Promise.all([
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import { describe, expect, it, vi } from "vitest";
|
||||
|
||||
import { findSameAgentDuplicates, flagSameAgentDuplicate } from "../duplicate-intake.js";
|
||||
import { computeParentIntentClaimId, findSameAgentDuplicates, flagSameAgentDuplicate } from "../duplicate-intake.js";
|
||||
import type { TaskStore } from "../store.js";
|
||||
|
||||
describe("findSameAgentDuplicates", () => {
|
||||
@@ -79,6 +79,49 @@ describe("findSameAgentDuplicates", () => {
|
||||
expect(matches[0]?.id).toBe("FN-5544");
|
||||
});
|
||||
|
||||
it("recognizes parent-scoped paraphrases through stable intent phrases", () => {
|
||||
const input = {
|
||||
description: "Add screenshot and activity-trace context capture with privacy scrub coverage before GitHub egress.",
|
||||
sourceParentTaskId: "FN-8277",
|
||||
};
|
||||
const candidate = {
|
||||
id: "FN-8309", title: "", description: "Add optional screenshot and short activity-trace context capture, preserving scrub-before-egress.",
|
||||
column: "triage", createdAt: nowMs - 60_000, sourceAgentId: null, sourceParentTaskId: "FN-8277",
|
||||
} as const;
|
||||
const matches = findSameAgentDuplicates(input, [candidate], { nowMs });
|
||||
expect(matches.map((match) => match.id)).toEqual(["FN-8309"]);
|
||||
expect(findSameAgentDuplicates(input, [{ ...candidate, tombstoned: true }], { nowMs })).toEqual([]);
|
||||
});
|
||||
|
||||
it("keeps distinct actions and sibling integrations separate", () => {
|
||||
const candidates = [
|
||||
{ id: "FN-UPLOAD", title: "", description: "Add screenshot upload support", column: "triage" as const, createdAt: nowMs - 60_000, sourceAgentId: null, sourceParentTaskId: "FN-PARENT" },
|
||||
{ id: "FN-DISCUSS", title: "", description: "Add GitHub Discussions as a filing target", column: "triage" as const, createdAt: nowMs - 60_000, sourceAgentId: null, sourceParentTaskId: "FN-PARENT" },
|
||||
];
|
||||
expect(findSameAgentDuplicates({ description: "Add screenshot deletion support", sourceParentTaskId: "FN-PARENT" }, candidates, { nowMs })).toEqual([]);
|
||||
expect(findSameAgentDuplicates({ description: "Add GitHub Issues as a filing target", sourceParentTaskId: "FN-PARENT" }, candidates, { nowMs })).toEqual([]);
|
||||
expect(findSameAgentDuplicates({
|
||||
description: "Add OAuth token revocation for the GitHub API",
|
||||
sourceParentTaskId: "FN-PARENT",
|
||||
}, [{
|
||||
id: "FN-ROTATE", title: "", description: "Add OAuth token rotation for the GitHub API",
|
||||
column: "triage", createdAt: nowMs - 60_000, sourceAgentId: null, sourceParentTaskId: "FN-PARENT",
|
||||
}], { nowMs })).toEqual([]);
|
||||
});
|
||||
|
||||
it("derives stable database claims for paraphrases and distinct claims for sibling actions", () => {
|
||||
const claim = (description: string) => computeParentIntentClaimId({ description, sourceParentTaskId: "fn-8277" });
|
||||
expect(claim("Add screenshot and short activity-trace context capture")).toBe(
|
||||
claim("Capture screenshots with activity-trace context"),
|
||||
);
|
||||
expect(claim("Add GitHub Discussions as a filing target")).toBe(
|
||||
claim("Use GitHub Discussions for optional filing"),
|
||||
);
|
||||
expect(claim("Add screenshot upload support")).not.toBe(claim("Add screenshot deletion support"));
|
||||
expect(claim("Add GitHub Discussions as a target")).not.toBe(claim("Add GitHub Issues as a target"));
|
||||
expect(claim("Add new support")).toMatch(/^agent-parent-intent:FN-8277:[a-f0-9]{64}$/);
|
||||
});
|
||||
|
||||
it("does not match sibling with different parent task", () => {
|
||||
const matches = findSameAgentDuplicates(
|
||||
{
|
||||
|
||||
@@ -15,6 +15,8 @@ export interface DeterministicGuardOptions {
|
||||
acknowledgedDuplicates?: readonly string[];
|
||||
bypass?: boolean;
|
||||
logger?: { warn(msg: string, data?: Record<string, unknown>): void };
|
||||
/** Serialize related creates even when their exact-content fingerprints differ. */
|
||||
serializationKey?: string;
|
||||
}
|
||||
|
||||
export interface DeterministicGuardOutcome {
|
||||
@@ -67,8 +69,16 @@ export async function runDeterministicDuplicateGuard(
|
||||
return { action: "proceed", fingerprint, releaseLock: noop };
|
||||
}
|
||||
|
||||
const lockKey = `${opts.lockScope}:${fingerprint}`;
|
||||
const lockKey = `${opts.lockScope}:${opts.serializationKey ?? fingerprint}`;
|
||||
const existingLock = deterministicGuardLocks.get(lockKey);
|
||||
let releaseCalled = false;
|
||||
let resolveGate: (() => void) | undefined;
|
||||
const gate = new Promise<void>((resolve) => {
|
||||
resolveGate = resolve;
|
||||
});
|
||||
// Install our tail before waiting so three or more callers form a queue
|
||||
// instead of all waking and proceeding when the first holder releases.
|
||||
deterministicGuardLocks.set(lockKey, gate);
|
||||
if (existingLock) {
|
||||
try {
|
||||
await existingLock;
|
||||
@@ -78,24 +88,17 @@ export async function runDeterministicDuplicateGuard(
|
||||
contentFingerprint: fingerprint,
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
});
|
||||
deterministicGuardLocks.delete(lockKey);
|
||||
if (deterministicGuardLocks.get(lockKey) === gate) deterministicGuardLocks.delete(lockKey);
|
||||
}
|
||||
}
|
||||
|
||||
let releaseCalled = false;
|
||||
let resolveGate: (() => void) | undefined;
|
||||
const gate = new Promise<void>((resolve) => {
|
||||
resolveGate = resolve;
|
||||
});
|
||||
deterministicGuardLocks.set(lockKey, gate);
|
||||
|
||||
const releaseLock = () => {
|
||||
if (releaseCalled) {
|
||||
return;
|
||||
}
|
||||
releaseCalled = true;
|
||||
resolveGate?.();
|
||||
deterministicGuardLocks.delete(lockKey);
|
||||
if (deterministicGuardLocks.get(lockKey) === gate) deterministicGuardLocks.delete(lockKey);
|
||||
};
|
||||
|
||||
try {
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { findDuplicateMatches } from "./duplicate-detection.js";
|
||||
import { computeContentFingerprint, findDuplicateMatches, tokenize } from "./duplicate-detection.js";
|
||||
import type { ColumnId } from "./types.js";
|
||||
import type { TaskStore } from "./store.js";
|
||||
|
||||
@@ -34,6 +34,50 @@ export interface SameAgentDuplicateMatch {
|
||||
allowResurrection?: boolean;
|
||||
}
|
||||
|
||||
const INTENT_BOILERPLATE_TOKENS = new Set([
|
||||
"add", "app", "agentic", "as", "bug", "feedback", "filing", "help", "idea", "ideas",
|
||||
"optional", "privacy", "report", "reports", "reporting", "pipeline", "existing", "new",
|
||||
"selectable", "support", "target", "use",
|
||||
]);
|
||||
|
||||
function intentTokens(title: string | null | undefined, description: string): string[] {
|
||||
return tokenize(`${title ?? ""} ${description}`)
|
||||
.filter((token) => token.length >= 2 && !INTENT_BOILERPLATE_TOKENS.has(token))
|
||||
.map((token) => token.length > 4 && token.endsWith("s") && !token.endsWith("ss") ? token.slice(0, -1) : token);
|
||||
}
|
||||
|
||||
function intentBigrams(title: string | null | undefined, description: string): Set<string> {
|
||||
const tokens = intentTokens(title, description);
|
||||
return new Set(tokens.slice(0, -1).map((token, index) => `${token}:${tokens[index + 1]}`));
|
||||
}
|
||||
|
||||
function hasSingleTokenReplacement(left: readonly string[], right: readonly string[]): boolean {
|
||||
return left.length === right.length
|
||||
&& left.filter((token, index) => token !== right[index]).length === 1;
|
||||
}
|
||||
|
||||
function stableIntentAnchor(title: string | null | undefined, description: string): string | null {
|
||||
const text = `${title ?? ""} ${description}`;
|
||||
const namedAnchors = [
|
||||
...(text.match(/[A-Za-z0-9]+(?:-[A-Za-z0-9]+)+/g) ?? []),
|
||||
...(text.match(/\b[A-Z][A-Za-z0-9]+(?:\s+[A-Z][A-Za-z0-9]+)+\b/g) ?? []),
|
||||
].map((value) => value.toLowerCase().split(/[\s-]+/)
|
||||
.filter((token) => token && !INTENT_BOILERPLATE_TOKENS.has(token))
|
||||
.map((token) => token.length > 4 && token.endsWith("s") && !token.endsWith("ss") ? token.slice(0, -1) : token)
|
||||
.join(":"))
|
||||
.filter(Boolean)
|
||||
.sort();
|
||||
return namedAnchors[0] ?? [...intentBigrams(title, description)].sort()[0] ?? null;
|
||||
}
|
||||
|
||||
/** Stable database idempotency claim for one parent-scoped follow-up intent. */
|
||||
export function computeParentIntentClaimId(input: SameAgentDuplicateInput): string | null {
|
||||
const parentId = input.sourceParentTaskId?.trim().toUpperCase();
|
||||
const anchor = stableIntentAnchor(input.title, input.description)
|
||||
?? computeContentFingerprint({ title: input.title, description: input.description });
|
||||
return parentId && anchor ? `agent-parent-intent:${parentId}:${anchor}` : null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Find candidate tasks that look like duplicates spawned by the same caller.
|
||||
*
|
||||
@@ -76,7 +120,7 @@ export function findSameAgentDuplicates(
|
||||
);
|
||||
|
||||
const metadataById = new Map(recent.map((candidate) => [candidate.id, candidate]));
|
||||
return matches.map((match) => {
|
||||
const mappedMatches = matches.map((match) => {
|
||||
const candidate = metadataById.get(match.id);
|
||||
return {
|
||||
id: match.id,
|
||||
@@ -86,6 +130,36 @@ export function findSameAgentDuplicates(
|
||||
allowResurrection: candidate?.allowResurrection,
|
||||
};
|
||||
});
|
||||
if (mappedMatches.length > 0 || !inputParentId) return mappedMatches;
|
||||
|
||||
/*
|
||||
FNXC:TaskCreationDeduplication 2026-07-18-12:36:
|
||||
Exact content fingerprints cannot contain a retried agent step that paraphrases
|
||||
its follow-up. Within one parent and the existing 24-hour window, reuse a live
|
||||
sibling only when a meaningful adjacent intent phrase survives the rewrite.
|
||||
Bigrams keep distinct actions such as "screenshot upload" and "screenshot delete"
|
||||
separate while recognizing stable concepts such as "GitHub Discussions".
|
||||
*/
|
||||
const sourceTokens = intentTokens(input.title, input.description);
|
||||
const sourceBigrams = intentBigrams(input.title, input.description);
|
||||
const sourceAnchor = stableIntentAnchor(input.title, input.description);
|
||||
if (sourceBigrams.size === 0) return [];
|
||||
return recent.flatMap((candidate) => {
|
||||
if (candidate.tombstoned || candidate.column === "done" || candidate.column === "archived") return [];
|
||||
const candidateTokens = intentTokens(candidate.title, candidate.description);
|
||||
if (sourceAnchor !== stableIntentAnchor(candidate.title, candidate.description)
|
||||
|| hasSingleTokenReplacement(sourceTokens, candidateTokens)) return [];
|
||||
const candidateBigrams = intentBigrams(candidate.title, candidate.description);
|
||||
const sharedCount = [...sourceBigrams].filter((bigram) => candidateBigrams.has(bigram)).length;
|
||||
const identicalIntent = sourceTokens.join(":") === candidateTokens.join(":");
|
||||
const diceScore = (2 * sharedCount) / (sourceBigrams.size + candidateBigrams.size);
|
||||
return identicalIntent || (sharedCount >= 2 && diceScore >= 0.3)
|
||||
? [{ id: candidate.id, score: diceScore }]
|
||||
: [];
|
||||
}).sort((left, right) =>
|
||||
(metadataById.get(left.id)?.createdAt ?? Number.POSITIVE_INFINITY)
|
||||
- (metadataById.get(right.id)?.createdAt ?? Number.POSITIVE_INFINITY),
|
||||
);
|
||||
}
|
||||
|
||||
export async function archiveAsSameAgentDuplicate(
|
||||
|
||||
@@ -716,6 +716,7 @@ export {
|
||||
export type { TaskDependencyMutation } from "./store.js";
|
||||
export {
|
||||
findSameAgentDuplicates,
|
||||
computeParentIntentClaimId,
|
||||
archiveAsSameAgentDuplicate,
|
||||
flagSameAgentDuplicate,
|
||||
flagTriageDuplicate,
|
||||
|
||||
@@ -739,6 +739,7 @@ export {
|
||||
export type { TaskDependencyMutation } from "./store.js";
|
||||
export {
|
||||
findSameAgentDuplicates,
|
||||
computeParentIntentClaimId,
|
||||
archiveAsSameAgentDuplicate,
|
||||
flagSameAgentDuplicate,
|
||||
flagTriageDuplicate,
|
||||
|
||||
@@ -103,6 +103,7 @@ import { claimNextToolFailureRetryImpl, clearNearDuplicateReferencesToFailSoftIm
|
||||
import { addPrInfoImpl, addSteeringCommentImpl, archiveAllDoneImpl, cleanupStaleMergeQueueRowsImpl, clearCompletionHandoffAcceptedMarkerImpl, clearDoneTransientFieldsImpl, clearStaleExecutionStartBranchReferencesImpl, computeWorkflowColumnsGraduationReportImpl, deleteTaskCommentImpl, deleteTaskDocumentImpl, emitUsageEventImpl, enqueueMergeQueueImpl, getAgentLogCountImpl, getAgentLogsImpl, getArtifactImpl, getArtifactsImpl, getAttachmentImpl, getCompletionHandoffAcceptedMarkerImpl, getTaskDocumentImpl, getTaskDocumentRevisionsImpl, getTaskDocumentsImpl, insertArtifactRowImpl, linkGithubIssueImpl, listWorkflowWorkItemsForTaskSyncImpl, moveToDoneImpl, parseDependenciesFromPromptImpl, parseFileScopeFromPromptImpl, parseStepsFromPromptImpl, peekMergeQueueHeadImpl, peekMergeQueueImpl, readPreArchiveColumnFromTaskFileImpl, recordPluginActivationImpl, recordRunAuditEventBackendImpl, removePrInfoByNumberImpl, resolvePrimaryPrInfoImpl, resolveUnarchiveTargetColumnImpl, rewriteLineageChildrenForRemovalImpl, runGitCommandImpl, stopWatchingImpl, syncAgentTaskLinkOnReassignmentImpl, updateArtifactImpl, updateGithubTrackingImpl, updatePrInfoByNumberImpl, updateTaskCommentImpl, upsertPrInfoByNumberImpl, writeArtifactDataImpl } from "./task-store/remaining-ops-7.js";
|
||||
import { approveCliAutonomyImpl, approveWorkflowCliCommandImpl, cleanupOrphanedMaterializedStepsImpl, consumePluginGateVerdictsImpl, getAgentLogsByTimeRangeImpl, getDatabaseHealthImpl, getDistributedTaskIdAllocatorImpl, getExperimentSessionStoreImpl, getInReviewDurationEventsImpl, getMissionStoreImpl, getIdeationStoreImpl, getPluginStoreImpl, getSecretsStoreImpl, getSettingsSyncImpl, getTaskMergedTaskIdsImpl, getTaskWorkflowSelectionImpl, getImportTranslationImpl, recordImportTranslationImpl, pruneImportTranslationsImpl, type ImportTranslationCacheKey, type ImportTranslationCacheEntry, getVerificationCacheHitImpl, getWorkflowDefinitionImpl, healthCheckImpl, importLegacyAgentLogsOnceImpl, insertWorkflowDefinitionSyncImpl, isCliAutonomyApprovedImpl, isPluginInstalledImpl, isWorkflowCliCommandApprovedImpl, listWorkflowDefinitionsImpl, materializeExplicitWorkflowStepsImpl, materializeWorkflowStepsImpl, migrateActiveArchivedTasksToArchiveDbImpl, migrateLegacyArchiveEntriesToArchiveDbImpl, nextWorkflowDefinitionIdImpl, occupantsByColumnForWorkflowImpl, parseWorkflowLayoutImpl, pruneAgentLogFilesImpl, purgeTaskWorkflowSelectionRowsImpl, readAllWorkflowDefinitionsImpl, readRawProjectSettingsImpl, recordPluginGateVerdictImpl, recordVerificationCachePassImpl, removeMaterializedSelectionImpl, resolvePluginWorkflowStepImpl, resolveTaskWorkflowIrSyncImpl, revokeCliAutonomyImpl, selectTaskWorkflowAndReconcileImpl, writeTaskWorkflowSelectionImpl, getTaskWorkflowSelectionAsyncImpl, } from "./task-store/remaining-ops-8.js";
|
||||
import { getTaskCommitAssociationsByLineageIdImpl, replaceLegacyTaskCommitAssociationsImpl } from "./task-store/task-commit-associations.js";
|
||||
import { findRecentTasksBySourceParentTaskIdImpl } from "./task-store/remaining-ops-6.js";
|
||||
import { addTaskCommentImpl, applyBuiltInPromptOverridesSyncImpl, areAllDependenciesDoneImpl, artifactStoredNameImpl, assertWorkflowIrTraitsValidImpl, clearActivityLogImpl, clearTaskWorkflowSelectionImpl, deleteTaskByIdImpl, getDefaultWorkflowIdImpl, getInsightStoreImpl, getMergeQueuedTaskIdsImpl, getMergeRequestRecordImpl, getMergeRequestRecordAsyncImpl, getResearchStoreImpl, getTaskIdFromDirImpl, getTodoStoreImpl, getWorkflowWorkItemByIdentityImpl, hasActiveTaskImpl, invalidateConfigCacheAfterMigrationImpl, isTaskIdConflictErrorImpl, listLegacyAutoMergeStampCandidatesImpl, readTaskRowFromDbImpl, recordBranchGroupMemberLandedImpl, refreshDatabaseHealthImpl, resolveEffectiveWorkflowIdSyncImpl, resolveTaskCustomFieldDefsSyncImpl, resolveWorkflowBypassGuardsImpl, serializeConfigForDiskImpl, setPluginWorkflowStepTemplatesImpl, shouldSkipWorkflowMovePoliciesImpl, suppressWatcherImpl, upsertTaskWithFtsRecoveryImpl } from "./task-store/task-store-helpers.js";
|
||||
import { getTaskSelectClauseImpl2, createTaskPersistSerializationContextImpl, getTaskPersistValuesImpl, getTaskPatchDescriptorsImpl, normalizeTaskFromDiskImpl, writeTaskJsonFileImpl, rowToPrEntityImpl, generatePrEntityIdImpl, readTaskForMoveImpl, rowToMergeQueueEntryImpl, rowToMergeRequestRecordImpl, rowToCompletionHandoffMarkerImpl, rowToWorkflowWorkItemImpl, rowToRunAuditEventImpl } from "./task-store/task-row-mappers.js";
|
||||
import { getTaskSelectClauseWithActivityLogLimitImpl, getChangedTaskColumnsImpl, getSoftDeletedWriteConflictImpl, readTaskJsonImpl, writeConfigImpl, _maybeAutoArchiveSameAgentDuplicateBackendImpl, updateBranchGroupImpl, updatePrEntityImpl, listTasksForGithubTrackingReconcileImpl, listTasksForGitlabTrackingReconcileImpl, renewCheckoutLeaseImpl, updateTaskAtomicImpl, getWorkflowPromptOverridesImpl, updateWorkflowSettingValuesImpl, rollbackConfigurationImpl, cancelActiveWorkflowWorkItemsForTaskImpl, setCompletionHandoffAcceptedMarkerImpl, reconcileLegacyAutoMergeStampsImpl, recoverExpiredMergeQueueLeasesImpl, rewriteDependentsForRemovalImpl, cleanupBranchForTaskImpl, addAttachmentImpl, deleteAttachmentImpl, registerArtifactImpl, updatePrInfoImpl, unlinkGithubIssueImpl, cleanupArchivedTasksImpl, generatePromptFromArchiveEntryImpl, listWorkflowOccupantTaskIdsImpl, evacuateCustomColumnsToLegacyImpl, listApprovedCliAutonomyAdaptersImpl, closeImpl, getActivityLogImpl } from "./task-store/remaining-ops-2.js";
|
||||
@@ -857,14 +858,14 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
|
||||
/**
|
||||
* FNXC:RuntimeTaskOrchestrationAsync 2026-06-24-13:15:
|
||||
*/
|
||||
public async createTaskBackend( input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise<string | null>; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; }, ): Promise<Task> {
|
||||
public async createTaskBackend( input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise<string | null>; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; }, ): Promise<Task> {
|
||||
return createTaskBackendImpl(this, input, options);
|
||||
}
|
||||
|
||||
/**
|
||||
* FNXC:RuntimeTaskOrchestrationAsync 2026-06-24-13:25:
|
||||
*/
|
||||
public async _createTaskInternalBackend( input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; }, ): Promise<Task> {
|
||||
public async _createTaskInternalBackend( input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; onProposalClaimConflict?: (task: Task) => void; }, ): Promise<Task> {
|
||||
return _createTaskInternalBackendImpl(this, input, title, resolvedWorkflowSteps, id, options);
|
||||
}
|
||||
|
||||
@@ -874,13 +875,13 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
|
||||
public async _maybeAutoArchiveSameAgentDuplicateBackend( task: Task, input: TaskCreateInput, ): Promise<void> {
|
||||
return _maybeAutoArchiveSameAgentDuplicateBackendImpl(this, task, input);
|
||||
}
|
||||
async createTask( input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise<string | null>; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; } ): Promise<Task> {
|
||||
async createTask( input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise<string | null>; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; } ): Promise<Task> {
|
||||
return createTaskImpl(this, input, options);
|
||||
}
|
||||
async createTaskWithReservedId( input: TaskCreateInput, options: { taskId: string; createdAt?: string; updatedAt?: string; prompt?: string; applyDefaultWorkflowSteps?: boolean; invokeTaskCreatedHook?: boolean; }, ): Promise<Task> {
|
||||
return createTaskWithReservedIdImpl(this, input, options);
|
||||
}
|
||||
public async _createTaskInternal( input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; }, ): Promise<Task> {
|
||||
public async _createTaskInternal( input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; onProposalClaimConflict?: (task: Task) => void; }, ): Promise<Task> {
|
||||
/*
|
||||
FNXC:SqliteFinalRemoval 2026-06-25-10:35:
|
||||
Route to the async backend variant when the store is in backend mode so
|
||||
@@ -1108,6 +1109,9 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
|
||||
async findRecentTasksByContentFingerprint( fingerprint: string, options?: { windowMs?: number; includeArchived?: boolean }, ): Promise<Task[]> {
|
||||
return findRecentTasksByContentFingerprintImpl(this, fingerprint, options);
|
||||
}
|
||||
async findRecentTasksBySourceParentTaskId(sourceParentTaskId: string, options?: { windowMs?: number }): Promise<Task[]> {
|
||||
return findRecentTasksBySourceParentTaskIdImpl(this, sourceParentTaskId, options);
|
||||
}
|
||||
|
||||
/** FNXC:NearDuplicateDetection 2026-06-14-12:00: FN-6439 requires the store to reconcile persisted duplicate flags after a canonical becomes inactive. */
|
||||
public async clearNearDuplicateReferencesTo( canonicalId: string, inactiveState: { column?: ColumnId | null; deletedAt?: string | null; reason: string }, ): Promise<Task[]> {
|
||||
|
||||
@@ -727,6 +727,35 @@ export async function findRecentTasksByContentFingerprintImpl(store: TaskStore,
|
||||
return rows.map((row) => store.rowToTask(row));
|
||||
}
|
||||
|
||||
export async function findRecentTasksBySourceParentTaskIdImpl(
|
||||
store: TaskStore,
|
||||
sourceParentTaskId: string,
|
||||
options?: { windowMs?: number },
|
||||
): Promise<Task[]> {
|
||||
const parentId = sourceParentTaskId.trim();
|
||||
if (!parentId) return [];
|
||||
const dayMs = 24 * 60 * 60 * 1000;
|
||||
const windowMs = Math.max(1, Math.min(dayMs, Math.trunc(options?.windowMs ?? dayMs)));
|
||||
const cutoffIso = new Date(Date.now() - windowMs).toISOString();
|
||||
if (store.backendMode) {
|
||||
const layer = store.asyncLayer!;
|
||||
const rows = await layer.db.select().from(schema.project.tasks).where(and(
|
||||
isNull(schema.project.tasks.deletedAt), taskProjectScope(layer),
|
||||
eq(schema.project.tasks.sourceParentTaskId, parentId),
|
||||
sql`${schema.project.tasks.createdAt} >= ${cutoffIso}`,
|
||||
ne(schema.project.tasks.column, "archived"), ne(schema.project.tasks.column, "done"),
|
||||
)).orderBy(asc(schema.project.tasks.createdAt));
|
||||
return rows.map((row) => store.rowToTask(store.pgRowToTaskRow(row as unknown as Record<string, unknown>)));
|
||||
}
|
||||
const selectClause = store.getTaskSelectClause(false, "t");
|
||||
const rows = store.db.prepare(`
|
||||
SELECT ${selectClause} FROM tasks t
|
||||
WHERE t."deletedAt" IS NULL AND t.sourceParentTaskId = ? AND t.createdAt >= ?
|
||||
AND t."column" NOT IN ('archived', 'done') ORDER BY t.createdAt ASC
|
||||
`).all(parentId, cutoffIso) as TaskRow[];
|
||||
return rows.map((row) => store.rowToTask(row));
|
||||
}
|
||||
|
||||
export async function clearNearDuplicateReferencesToFailSoftImpl(store: TaskStore,
|
||||
canonicalId: string,
|
||||
inactiveState: { column?: ColumnId | null; deletedAt?: string | null; reason: string },
|
||||
|
||||
@@ -44,7 +44,7 @@ function ensureSqliteProposalClaimUniqueness(store: TaskStore): void {
|
||||
);
|
||||
}
|
||||
|
||||
export async function createTaskBackendImpl(store: TaskStore, input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise<string | null>; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; },): Promise<Task> {
|
||||
export async function createTaskBackendImpl(store: TaskStore, input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise<string | null>; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; },): Promise<Task> {
|
||||
if (!input.description?.trim()) {
|
||||
throw new Error("Description is required and cannot be empty");
|
||||
}
|
||||
@@ -205,7 +205,7 @@ export async function createTaskBackendImpl(store: TaskStore, input: TaskCreateI
|
||||
title,
|
||||
resolvedWorkflowSteps,
|
||||
reservation.taskId,
|
||||
{ invokeTaskCreatedHook: shouldInvokeTaskCreatedHook && !hasPendingSummarization, resolvedEntryColumn },
|
||||
{ invokeTaskCreatedHook: shouldInvokeTaskCreatedHook && !hasPendingSummarization, resolvedEntryColumn, onProposalClaimConflict: options?.onProposalClaimConflict },
|
||||
);
|
||||
await allocator.commitDistributedTaskIdReservation({
|
||||
reservationId: reservation.reservationId,
|
||||
@@ -289,7 +289,7 @@ export async function createTaskBackendImpl(store: TaskStore, input: TaskCreateI
|
||||
return task;
|
||||
}
|
||||
|
||||
export async function _createTaskInternalBackendImpl(store: TaskStore, input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; },): Promise<Task> {
|
||||
export async function _createTaskInternalBackendImpl(store: TaskStore, input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; onProposalClaimConflict?: (task: Task) => void; },): Promise<Task> {
|
||||
const layer = store.asyncLayer!;
|
||||
const now = options?.createdAt ?? new Date().toISOString();
|
||||
const normalizedTitle = normalizeTitleForTaskId(title, id);
|
||||
@@ -388,7 +388,10 @@ export async function _createTaskInternalBackendImpl(store: TaskStore, input: Ta
|
||||
*/
|
||||
if (input.proposalClaimId && isTaskIdConflictError(error)) {
|
||||
const existing = (await store.listTasks()).find((candidate) => candidate.proposalClaimId === input.proposalClaimId);
|
||||
if (existing) return existing;
|
||||
if (existing) {
|
||||
options?.onProposalClaimConflict?.(existing);
|
||||
return existing;
|
||||
}
|
||||
}
|
||||
if (isTaskIdConflictError(error)) {
|
||||
throw new Error(`Task ID already exists: ${task.id}`);
|
||||
@@ -445,7 +448,7 @@ export async function _createTaskInternalBackendImpl(store: TaskStore, input: Ta
|
||||
return task;
|
||||
}
|
||||
|
||||
export async function createTaskImpl(store: TaskStore, input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise<string | null>; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; }): Promise<Task> {
|
||||
export async function createTaskImpl(store: TaskStore, input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise<string | null>; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; }): Promise<Task> {
|
||||
// FNXC:RuntimeTaskOrchestrationAsync 2026-06-24-13:10:
|
||||
// Backend-mode createTask: delegates to createTaskBackend which uses the
|
||||
// async DistributedTaskIdAllocator (now wired for backend mode) and the
|
||||
@@ -462,7 +465,10 @@ export async function createTaskImpl(store: TaskStore, input: TaskCreateInput, o
|
||||
if (input.proposalClaimId) {
|
||||
ensureSqliteProposalClaimUniqueness(store);
|
||||
const existing = (await store.listTasks()).find((task) => task.proposalClaimId === input.proposalClaimId);
|
||||
if (existing) return existing;
|
||||
if (existing) {
|
||||
options?.onProposalClaimConflict?.(existing);
|
||||
return existing;
|
||||
}
|
||||
}
|
||||
|
||||
const selfDefeatingDep = detectSelfDefeatingDependency(input.title, input.dependencies ?? []);
|
||||
@@ -626,7 +632,7 @@ export async function createTaskImpl(store: TaskStore, input: TaskCreateInput, o
|
||||
title,
|
||||
resolvedWorkflowSteps,
|
||||
taskId,
|
||||
{ invokeTaskCreatedHook: shouldInvokeTaskCreatedHook && !hasPendingSummarization, resolvedEntryColumn },
|
||||
{ invokeTaskCreatedHook: shouldInvokeTaskCreatedHook && !hasPendingSummarization, resolvedEntryColumn, onProposalClaimConflict: options?.onProposalClaimConflict },
|
||||
);
|
||||
},
|
||||
});
|
||||
@@ -636,7 +642,10 @@ export async function createTaskImpl(store: TaskStore, input: TaskCreateInput, o
|
||||
await store.cleanupOrphanedMaterializedSteps(pendingWorkflowSelection?.stepIds);
|
||||
if (input.proposalClaimId && isTaskIdConflictError(err)) {
|
||||
const existing = (await store.listTasks()).find((candidate) => candidate.proposalClaimId === input.proposalClaimId);
|
||||
if (existing) return existing;
|
||||
if (existing) {
|
||||
options?.onProposalClaimConflict?.(existing);
|
||||
return existing;
|
||||
}
|
||||
}
|
||||
throw err;
|
||||
}
|
||||
@@ -869,7 +878,7 @@ export async function createTaskWithReservedIdImpl(store: TaskStore, input: Task
|
||||
return createdTask;
|
||||
}
|
||||
|
||||
export async function _createTaskInternalImpl(store: TaskStore, input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; },): Promise<Task> {
|
||||
export async function _createTaskInternalImpl(store: TaskStore, input: TaskCreateInput, title: string | undefined, resolvedWorkflowSteps: string[] | undefined, id: string, options?: { createdAt?: string; updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; resolvedEntryColumn?: string; onProposalClaimConflict?: (task: Task) => void; },): Promise<Task> {
|
||||
const now = options?.createdAt ?? new Date().toISOString();
|
||||
// FN-5077: null normalized titles are treated as "no title" and allow standard fallback/summarization behavior.
|
||||
const normalizedTitle = normalizeTitleForTaskId(title, id);
|
||||
@@ -1128,4 +1137,3 @@ export async function _maybeAutoArchiveSameAgentDuplicateImpl(store: TaskStore,
|
||||
storeLog.warn(`FN-4892 same-agent duplicate intake failed open for ${task.id}: ${getErrorMessage(error)}`);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import { describe, it, expect, vi, beforeEach } from "vitest";
|
||||
import type { Agent, AgentStore, TaskStore, Task } from "@fusion/core";
|
||||
import { createAgentTask, createListAgentsTool, createDelegateTaskTool } from "../agent-tools.js";
|
||||
import { createAgentTask, createListAgentsTool, createDelegateTaskTool, createTaskCreateTool } from "../agent-tools.js";
|
||||
|
||||
function createMockAgentStore(overrides: Partial<AgentStore> = {}): AgentStore {
|
||||
return {
|
||||
@@ -14,6 +14,7 @@ function createMockAgentStore(overrides: Partial<AgentStore> = {}): AgentStore {
|
||||
function createMockTaskStore(overrides: Partial<TaskStore> = {}): TaskStore {
|
||||
return {
|
||||
getSettings: vi.fn().mockResolvedValue({ autoSummarizeTitles: false }),
|
||||
findRecentTasksBySourceParentTaskId: vi.fn().mockResolvedValue([]),
|
||||
findRecentTasksByContentFingerprint: vi.fn().mockResolvedValue([]),
|
||||
updateTask: vi.fn(),
|
||||
moveTask: vi.fn(),
|
||||
@@ -301,6 +302,93 @@ describe("createDelegateTaskTool", () => {
|
||||
expect(taskStore.moveTask).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("reuses paraphrased follow-ups from the same parent without collapsing distinct sibling intents", async () => {
|
||||
const tasks: Task[] = [];
|
||||
vi.mocked(taskStore.findRecentTasksBySourceParentTaskId).mockImplementation(async () => tasks);
|
||||
vi.mocked(taskStore.createTask).mockImplementation(async (input) => {
|
||||
const created = {
|
||||
id: `FN-${tasks.length + 1}`, title: input.title, description: input.description,
|
||||
dependencies: input.dependencies ?? [], column: "triage" as const,
|
||||
sourceType: input.source?.sourceType, sourceAgentId: input.source?.sourceAgentId,
|
||||
sourceParentTaskId: input.source?.sourceParentTaskId,
|
||||
steps: [], currentStep: 0, log: [], createdAt: new Date().toISOString(), updatedAt: new Date().toISOString(),
|
||||
} as Task;
|
||||
tasks.push(created);
|
||||
return created;
|
||||
});
|
||||
const descriptions = [
|
||||
"Add optional screenshot and short activity-trace context capture to in-app agentic reports, preserving the scrub-before-egress boundary.",
|
||||
"Add GitHub Discussions as a selectable filing target for in-app agentic reports.",
|
||||
"Add public-roadmap (FR-30) as a deduplication source for in-app agentic reports.",
|
||||
];
|
||||
for (const description of descriptions) {
|
||||
expect((await createAgentTask(taskStore, { description }, { sourceTaskId: "FN-PARENT" })).wasDuplicate).toBe(false);
|
||||
}
|
||||
const replay = await createAgentTask(taskStore, {
|
||||
description: "Add screenshot and activity-trace context capture to in-app Bug/Feedback/Idea/Help reports, with privacy scrub coverage before GitHub egress.",
|
||||
}, { sourceTaskId: "FN-PARENT" });
|
||||
expect(replay.wasDuplicate).toBe(true);
|
||||
expect(replay.task.id).toBe("FN-1");
|
||||
expect(tasks).toHaveLength(3);
|
||||
});
|
||||
|
||||
it("persists option-based parent provenance on the step-session fn_task_create surface", async () => {
|
||||
const tool = createTaskCreateTool(taskStore, undefined, { sourceTaskId: "FN-PARENT", sourceAgentId: "agent-worker" });
|
||||
await tool.execute("call-1", { description: "Capture optional report screenshots" }, undefined as any, undefined as any, undefined as any);
|
||||
expect(taskStore.createTask).toHaveBeenCalledWith(expect.objectContaining({
|
||||
source: expect.objectContaining({ sourceType: "api", sourceAgentId: "agent-worker", sourceParentTaskId: "FN-PARENT" }),
|
||||
}), expect.anything());
|
||||
});
|
||||
|
||||
it("serializes three concurrent paraphrased creates from one parent", async () => {
|
||||
const tasks: Task[] = [];
|
||||
vi.mocked(taskStore.findRecentTasksBySourceParentTaskId).mockImplementation(async () => tasks);
|
||||
vi.mocked(taskStore.createTask).mockImplementation(async (input) => {
|
||||
await Promise.resolve();
|
||||
const created = { id: `FN-${tasks.length + 1}`, description: input.description, dependencies: [], column: "triage" as const,
|
||||
sourceParentTaskId: input.source?.sourceParentTaskId, steps: [], currentStep: 0, log: [],
|
||||
createdAt: new Date().toISOString(), updatedAt: new Date().toISOString() } as Task;
|
||||
tasks.push(created); return created;
|
||||
});
|
||||
const results = await Promise.all([
|
||||
"Add screenshot and activity-trace capture with privacy scrub coverage.",
|
||||
"Add screenshot and activity-trace capture while preserving privacy scrubbing.",
|
||||
"Capture screenshots and activity traces with mandatory privacy scrubbing.",
|
||||
].map((description) => createAgentTask(taskStore, { description }, { sourceTaskId: "FN-PARENT" })));
|
||||
expect(results.filter((result) => result.wasDuplicate)).toHaveLength(2);
|
||||
expect(tasks).toHaveLength(1);
|
||||
});
|
||||
|
||||
it("fails closed when parent-scoped duplicate lookup is unavailable", async () => {
|
||||
vi.mocked(taskStore.findRecentTasksBySourceParentTaskId).mockRejectedValue(new Error("database unavailable"));
|
||||
await expect(createAgentTask(taskStore, {
|
||||
description: "Add screenshot and activity-trace capture",
|
||||
}, { sourceTaskId: "FN-PARENT" })).rejects.toThrow("Unable to verify parent-scoped task uniqueness");
|
||||
expect(taskStore.createTask).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("normalizes parent provenance and reports database claim reuse as a duplicate", async () => {
|
||||
const canonical = {
|
||||
id: "FN-1", title: "", description: "Add new support", dependencies: [], column: "triage" as const,
|
||||
sourceParentTaskId: "FN-PARENT", proposalClaimId: "claim", steps: [], currentStep: 0, log: [],
|
||||
createdAt: new Date().toISOString(), updatedAt: new Date().toISOString(),
|
||||
} as Task;
|
||||
vi.mocked(taskStore.findRecentTasksBySourceParentTaskId).mockResolvedValue([]);
|
||||
vi.mocked(taskStore.createTask).mockImplementation(async (_input, options) => {
|
||||
options?.onProposalClaimConflict?.(canonical);
|
||||
return canonical;
|
||||
});
|
||||
|
||||
const result = await createAgentTask(taskStore, { description: "Add new support" }, { sourceTaskId: "fn-parent" });
|
||||
|
||||
expect(taskStore.findRecentTasksBySourceParentTaskId).toHaveBeenCalledWith("FN-PARENT");
|
||||
expect(taskStore.createTask).toHaveBeenCalledWith(expect.objectContaining({
|
||||
source: expect.objectContaining({ sourceParentTaskId: "FN-PARENT" }),
|
||||
proposalClaimId: expect.stringMatching(/^agent-parent-intent:FN-PARENT:/),
|
||||
}), expect.anything());
|
||||
expect(result).toMatchObject({ task: canonical, wasDuplicate: true });
|
||||
});
|
||||
|
||||
it("carries delegation routing onto the reconcile canonical task", async () => {
|
||||
const created = {
|
||||
id: "FN-new",
|
||||
|
||||
@@ -2988,6 +2988,7 @@ describe("HeartbeatTriggerScheduler", () => {
|
||||
// Create a minimal mock TaskStore
|
||||
const mockTaskStore = {
|
||||
createTask: vi.fn().mockResolvedValue({ id: "FN-200", description: "New task created", dependencies: [] }),
|
||||
findRecentTasksBySourceParentTaskId: vi.fn().mockResolvedValue([]),
|
||||
logEntry: vi.fn().mockResolvedValue({}),
|
||||
getTask: vi.fn().mockResolvedValue({
|
||||
id: "FN-001",
|
||||
@@ -3023,10 +3024,66 @@ describe("HeartbeatTriggerScheduler", () => {
|
||||
);
|
||||
});
|
||||
|
||||
it("does not audit a reused heartbeat follow-up as a new task", async () => {
|
||||
const canonical = {
|
||||
id: "FN-200", title: "", description: "Add GitHub Discussions as a filing target",
|
||||
dependencies: [], column: "triage", sourceParentTaskId: "FN-001",
|
||||
steps: [], currentStep: 0, log: [], createdAt: new Date().toISOString(), updatedAt: new Date().toISOString(),
|
||||
};
|
||||
const mockTaskStore = {
|
||||
getSettings: vi.fn().mockResolvedValue({}),
|
||||
findRecentTasksByContentFingerprint: vi.fn().mockResolvedValue([]),
|
||||
findRecentTasksBySourceParentTaskId: vi.fn().mockResolvedValue([canonical]),
|
||||
createTask: vi.fn(),
|
||||
logEntry: vi.fn().mockResolvedValue({}),
|
||||
} as unknown as import("@fusion/core").TaskStore;
|
||||
const audit = { database: vi.fn() } as any;
|
||||
const monitor = new HeartbeatMonitor({ store, taskStore: mockTaskStore, rootDir: "/tmp" });
|
||||
const tool = monitor.createHeartbeatTools("agent-abc", mockTaskStore, "FN-001", undefined, audit)
|
||||
.find((candidate) => candidate.name === "fn_task_create")!;
|
||||
|
||||
const result = await tool.execute("call-replay", {
|
||||
description: "Add GitHub Discussions as an optional filing target",
|
||||
}, undefined as any, undefined as any, undefined as any);
|
||||
|
||||
expect((result.details as { wasDuplicate?: boolean }).wasDuplicate).toBe(true);
|
||||
expect(mockTaskStore.createTask).not.toHaveBeenCalled();
|
||||
expect(audit.database).not.toHaveBeenCalled();
|
||||
expect(mockTaskStore.logEntry).toHaveBeenCalledWith(
|
||||
"FN-200", "Linked existing task by agent agent-abc during heartbeat run", undefined, undefined,
|
||||
);
|
||||
});
|
||||
|
||||
it("does not log or audit a failed heartbeat task creation without a task id", async () => {
|
||||
const createError = Object.assign(new Error("self-defeating dependency"), {
|
||||
code: "SELF_DEFEATING_DEPENDENCY",
|
||||
});
|
||||
const mockTaskStore = {
|
||||
getSettings: vi.fn().mockResolvedValue({}),
|
||||
findRecentTasksByContentFingerprint: vi.fn().mockResolvedValue([]),
|
||||
findRecentTasksBySourceParentTaskId: vi.fn().mockResolvedValue([]),
|
||||
createTask: vi.fn().mockRejectedValue(createError),
|
||||
logEntry: vi.fn(),
|
||||
} as unknown as import("@fusion/core").TaskStore;
|
||||
const audit = { database: vi.fn() } as any;
|
||||
const monitor = new HeartbeatMonitor({ store, taskStore: mockTaskStore, rootDir: "/tmp" });
|
||||
const tool = monitor.createHeartbeatTools("agent-abc", mockTaskStore, "FN-001", undefined, audit)
|
||||
.find((candidate) => candidate.name === "fn_task_create")!;
|
||||
|
||||
const result = await tool.execute("call-error", {
|
||||
description: "Create an invalid follow-up",
|
||||
}, undefined as any, undefined as any, undefined as any);
|
||||
|
||||
expect(result.isError).toBe(true);
|
||||
expect(mockTaskStore.logEntry).not.toHaveBeenCalled();
|
||||
expect(audit.database).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("createHeartbeatTools works without runContext (backward compat)", async () => {
|
||||
// Create a minimal mock TaskStore
|
||||
const mockTaskStore = {
|
||||
createTask: vi.fn().mockResolvedValue({ id: "FN-NEW", description: "New task" }),
|
||||
findRecentTasksBySourceParentTaskId: vi.fn().mockResolvedValue([]),
|
||||
logEntry: vi.fn().mockResolvedValue({}),
|
||||
getTask: vi.fn().mockResolvedValue({
|
||||
id: "FN-001",
|
||||
@@ -3268,4 +3325,3 @@ describe("HeartbeatTriggerScheduler", () => {
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
|
||||
@@ -3718,15 +3718,25 @@ export class HeartbeatMonitor {
|
||||
execute: async (id: string, params: Static<typeof taskCreateParams>, signal, onUpdate, ctx) => {
|
||||
const result = await baseCreateTool.execute(id, params, signal, onUpdate, ctx);
|
||||
|
||||
const createdTaskId = (result.details as { taskId?: string })?.taskId ?? "unknown";
|
||||
const resultDetails = result.details as { taskId?: string; wasDuplicate?: boolean };
|
||||
if (!resultDetails.taskId) return result;
|
||||
const createdTaskId = resultDetails.taskId;
|
||||
const wasDuplicate = resultDetails.wasDuplicate === true;
|
||||
|
||||
// Log agent link on the created task with run context for correlation
|
||||
try {
|
||||
await taskStore.logEntry(createdTaskId, `Created by agent ${agentId} during heartbeat run`, undefined, runContext);
|
||||
await taskStore.logEntry(
|
||||
createdTaskId,
|
||||
`${wasDuplicate ? "Linked existing task" : "Created"} by agent ${agentId} during heartbeat run`,
|
||||
undefined,
|
||||
runContext,
|
||||
);
|
||||
} catch (taskCreateLogErr) {
|
||||
heartbeatLog.warn(`Task ${createdTaskId} agent-link log failed: ${taskCreateLogErr instanceof Error ? taskCreateLogErr.message : String(taskCreateLogErr)}`);
|
||||
}
|
||||
|
||||
if (wasDuplicate) return result;
|
||||
|
||||
// Audit trail: record task creation (FN-1404)
|
||||
await audit?.database({ type: "task:create", target: createdTaskId });
|
||||
|
||||
|
||||
@@ -16,7 +16,7 @@ import * as fusionCore from "@fusion/core";
|
||||
import type { AgentState, AgentCapability, AgentUpdateInput, AgentLogEntry, Artifact, ArtifactCreateInput, ArtifactWithTask, Task, TaskDocument, TaskDocumentCreateInput, TaskStore, RunMutationContext, MessageStore, Message, SourceType, Settings, ResearchRun, ResearchRunStatus, TaskCreateInput, ReflectionStore, ApprovalRequestStore, ProjectSettings, ChatStore, WorkflowSettingDefinition, GoalStatus, WorkflowIrNode } from "@fusion/core";
|
||||
import { listTraits, isBuiltinWorkflowId, AgentStore, validateColumnAgentBindings, ColumnAgentBindingError, stripApprovalBypassFlags, WorkflowSettingRejectionError, resolveEffectiveSettingsById, resolveWorkflowIrById, findOrphanedSettingValues, BUILTIN_WORKFLOW_SETTINGS, MAX_TASK_LIST_TEXT_CHARS, formatCurrentTaskLine, normalizeWorkflowIcon, parseWorkflowIr, WorkflowIrError, assertColumnTraitsValid, ColumnTraitValidationError } from "@fusion/core";
|
||||
import { promoteHeldTask } from "./hold-release.js";
|
||||
import { DASHBOARD_USER_ID, dailyMemoryPath, ensureOpenClawMemoryFiles, evaluateImplementationTaskBind, extractAgentProvisioningRequest, getMemoryBackendCapabilities, getProjectMemory, isEphemeralAgent, memoryLongTermPath, normalizeMessageParticipant, reconcileDeterministicDuplicate, resolveAgentProvisioningPolicy, resolveMemoryBackend, resolveResearchSettings, resolveTaskGithubTracking, runDeterministicDuplicateGuard, scheduleQmdProjectMemoryRefresh, searchProjectMemory, shouldSkipBackgroundQmdRefresh } from "@fusion/core";
|
||||
import { computeParentIntentClaimId, DASHBOARD_USER_ID, dailyMemoryPath, ensureOpenClawMemoryFiles, evaluateImplementationTaskBind, extractAgentProvisioningRequest, findSameAgentDuplicates, getMemoryBackendCapabilities, getProjectMemory, isEphemeralAgent, memoryLongTermPath, normalizeMessageParticipant, reconcileDeterministicDuplicate, resolveAgentProvisioningPolicy, resolveMemoryBackend, resolveResearchSettings, resolveTaskGithubTracking, runDeterministicDuplicateGuard, scheduleQmdProjectMemoryRefresh, searchProjectMemory, shouldSkipBackgroundQmdRefresh } from "@fusion/core";
|
||||
import { ResearchOrchestrator } from "./research-orchestrator.js";
|
||||
import { ResearchProviderRegistry } from "./research/provider-registry.js";
|
||||
import { ResearchStepRunner } from "./research-step-runner.js";
|
||||
@@ -973,6 +973,16 @@ export async function createAgentTask(
|
||||
? await store.getSettings()
|
||||
: {} as Settings;
|
||||
const rootDir = options?.rootDir;
|
||||
const sourceParentTaskId = (input.source?.sourceParentTaskId ?? options?.sourceTaskId)?.trim().toUpperCase();
|
||||
const sourceAgentId = input.source?.sourceAgentId ?? options?.sourceAgentId;
|
||||
const effectiveSource = input.source || sourceParentTaskId || sourceAgentId
|
||||
? {
|
||||
sourceType: input.source?.sourceType ?? "api" as const,
|
||||
...input.source,
|
||||
sourceAgentId,
|
||||
sourceParentTaskId,
|
||||
}
|
||||
: undefined;
|
||||
const guard = await runDeterministicDuplicateGuard(store, {
|
||||
title: input.title,
|
||||
description: input.description,
|
||||
@@ -980,6 +990,7 @@ export async function createAgentTask(
|
||||
lockScope: rootDir ?? store.getRootDir?.() ?? "agent-tools",
|
||||
bypass: options?.bypassDuplicateCheck === true,
|
||||
acknowledgedDuplicates: options?.acknowledgedDuplicates,
|
||||
serializationKey: sourceParentTaskId ? `parent:${sourceParentTaskId}` : undefined,
|
||||
logger: log,
|
||||
});
|
||||
|
||||
@@ -991,13 +1002,44 @@ export async function createAgentTask(
|
||||
};
|
||||
}
|
||||
|
||||
if (sourceParentTaskId && options?.bypassDuplicateCheck !== true) {
|
||||
try {
|
||||
const acknowledged = new Set(options?.acknowledgedDuplicates ?? []);
|
||||
const candidates = await store.findRecentTasksBySourceParentTaskId(sourceParentTaskId);
|
||||
const matches = findSameAgentDuplicates({
|
||||
title: input.title,
|
||||
description: input.description,
|
||||
sourceParentTaskId,
|
||||
}, candidates.map((candidate) => ({
|
||||
id: candidate.id,
|
||||
title: candidate.title ?? "",
|
||||
description: candidate.description,
|
||||
column: candidate.column,
|
||||
createdAt: Date.parse(candidate.createdAt),
|
||||
sourceAgentId: candidate.sourceAgentId ?? null,
|
||||
sourceParentTaskId: candidate.sourceParentTaskId ?? null,
|
||||
})), { sourceAgentId: sourceAgentId ?? null });
|
||||
const match = matches.find((candidate) => !acknowledged.has(candidate.id));
|
||||
const canonical = match ? candidates.find((candidate) => candidate.id === match.id) : undefined;
|
||||
if (canonical) {
|
||||
return { task: await carryCanonicalTaskRouting(store, canonical, input), wasDuplicate: true };
|
||||
}
|
||||
} catch (error) {
|
||||
log.warn("Parent-scoped task duplicate pre-check failed; aborting creation", {
|
||||
sourceParentTaskId,
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
});
|
||||
throw new Error(`Unable to verify parent-scoped task uniqueness for ${sourceParentTaskId}`, { cause: error });
|
||||
}
|
||||
}
|
||||
|
||||
const sourceMetadata = {
|
||||
...(input.source?.sourceMetadata ?? {}),
|
||||
...(effectiveSource?.sourceMetadata ?? {}),
|
||||
...(guard.fingerprint ? { contentFingerprint: guard.fingerprint } : {}),
|
||||
};
|
||||
const nextSource = input.source
|
||||
const nextSource = effectiveSource
|
||||
? {
|
||||
...input.source,
|
||||
...effectiveSource,
|
||||
sourceMetadata: Object.keys(sourceMetadata).length > 0 ? sourceMetadata : undefined,
|
||||
}
|
||||
: undefined;
|
||||
@@ -1014,6 +1056,11 @@ export async function createAgentTask(
|
||||
input.githubTracking?.enabled !== false && resolvedTracking.enabled;
|
||||
const createInput: TaskCreateInput = {
|
||||
...input,
|
||||
proposalClaimId: input.proposalClaimId ?? (
|
||||
options?.bypassDuplicateCheck === true
|
||||
? undefined
|
||||
: computeParentIntentClaimId({ title: input.title, description: input.description, sourceParentTaskId }) ?? undefined
|
||||
),
|
||||
summarize: !input.title?.trim() ? true : undefined,
|
||||
source: nextSource,
|
||||
githubTracking: shouldPrefillGithubTrackingEnabled
|
||||
@@ -1027,8 +1074,10 @@ export async function createAgentTask(
|
||||
: input.githubTracking,
|
||||
};
|
||||
|
||||
let proposalClaimConflict = false;
|
||||
const createdTask = await store.createTask(createInput, {
|
||||
settings,
|
||||
onProposalClaimConflict: () => { proposalClaimConflict = true; },
|
||||
});
|
||||
|
||||
const reconcile = await reconcileDeterministicDuplicate(store, {
|
||||
@@ -1038,10 +1087,12 @@ export async function createAgentTask(
|
||||
});
|
||||
|
||||
return {
|
||||
task: reconcile.outcome === "archived"
|
||||
task: proposalClaimConflict
|
||||
? await carryCanonicalTaskRouting(store, createdTask, input)
|
||||
: reconcile.outcome === "archived"
|
||||
? await carryCanonicalTaskRouting(store, reconcile.canonical, input)
|
||||
: reconcile.canonical,
|
||||
wasDuplicate: reconcile.outcome === "archived",
|
||||
wasDuplicate: proposalClaimConflict || reconcile.outcome === "archived",
|
||||
};
|
||||
} finally {
|
||||
guard.releaseLock();
|
||||
@@ -1137,7 +1188,7 @@ export function createTaskCreateTool(
|
||||
type: "text" as const,
|
||||
text: `${wasDuplicate ? "Linked existing" : "Created"} ${task.id}: ${params.description}${deps}${workflow}`,
|
||||
}],
|
||||
details: { taskId: task.id },
|
||||
details: { taskId: task.id, wasDuplicate },
|
||||
};
|
||||
} catch (err) {
|
||||
if (err instanceof Error && err.message.startsWith("Task ID already exists:")) {
|
||||
|
||||
@@ -60,6 +60,7 @@ export {
|
||||
createWorkflowSettingsTool,
|
||||
createTraitListTool,
|
||||
createWorkflowAuthoringTools,
|
||||
createAgentTask,
|
||||
taskCreateParams,
|
||||
taskListParams,
|
||||
taskShowParams,
|
||||
|
||||
Reference in New Issue
Block a user