diff --git a/.changeset/calm-tasks-reuse-parent-intent.md b/.changeset/calm-tasks-reuse-parent-intent.md new file mode 100644 index 0000000000..45603b148f --- /dev/null +++ b/.changeset/calm-tasks-reuse-parent-intent.md @@ -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. diff --git a/packages/cli/src/__tests__/extension-experiment-finalize.test.ts b/packages/cli/src/__tests__/extension-experiment-finalize.test.ts index 515a9e797d..aa50319caa 100644 --- a/packages/cli/src/__tests__/extension-experiment-finalize.test.ts +++ b/packages/cli/src/__tests__/extension-experiment-finalize.test.ts @@ -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(() => ({})), diff --git a/packages/cli/src/__tests__/extension-gitlab-tracking.test.ts b/packages/cli/src/__tests__/extension-gitlab-tracking.test.ts index 597a6aa447..86f253c874 100644 --- a/packages/cli/src/__tests__/extension-gitlab-tracking.test.ts +++ b/packages/cli/src/__tests__/extension-gitlab-tracking.test.ts @@ -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(), diff --git a/packages/cli/src/__tests__/extension-web-fetch.test.ts b/packages/cli/src/__tests__/extension-web-fetch.test.ts index 6b0d782cd2..b6382dbc48 100644 --- a/packages/cli/src/__tests__/extension-web-fetch.test.ts +++ b/packages/cli/src/__tests__/extension-web-fetch.test.ts @@ -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(() => ({})), diff --git a/packages/cli/src/__tests__/extension.test.ts b/packages/cli/src/__tests__/extension.test.ts index 4309ff7c51..d3c4175ed2 100644 --- a/packages/cli/src/__tests__/extension.test.ts +++ b/packages/cli/src/__tests__/extension.test.ts @@ -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( diff --git a/packages/cli/src/extension.ts b/packages/cli/src/extension.ts index 9e2c4315b9..dfdc62d0b0 100644 --- a/packages/cli/src/extension.ts +++ b/packages/cli/src/extension.ts @@ -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, diff --git a/packages/core/src/__tests__/duplicate-guard.test.ts b/packages/core/src/__tests__/duplicate-guard.test.ts index e3307c38de..319319aca7 100644 --- a/packages/core/src/__tests__/duplicate-guard.test.ts +++ b/packages/core/src/__tests__/duplicate-guard.test.ts @@ -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([ diff --git a/packages/core/src/__tests__/duplicate-intake.test.ts b/packages/core/src/__tests__/duplicate-intake.test.ts index 51eb1b92f5..2985f69311 100644 --- a/packages/core/src/__tests__/duplicate-intake.test.ts +++ b/packages/core/src/__tests__/duplicate-intake.test.ts @@ -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( { diff --git a/packages/core/src/duplicate-guard.ts b/packages/core/src/duplicate-guard.ts index 078efb8a5d..d836440e83 100644 --- a/packages/core/src/duplicate-guard.ts +++ b/packages/core/src/duplicate-guard.ts @@ -15,6 +15,8 @@ export interface DeterministicGuardOptions { acknowledgedDuplicates?: readonly string[]; bypass?: boolean; logger?: { warn(msg: string, data?: Record): 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((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((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 { diff --git a/packages/core/src/duplicate-intake.ts b/packages/core/src/duplicate-intake.ts index 14c26147c4..0fc0a92325 100644 --- a/packages/core/src/duplicate-intake.ts +++ b/packages/core/src/duplicate-intake.ts @@ -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 { + 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( diff --git a/packages/core/src/index.gate.ts b/packages/core/src/index.gate.ts index 9963ee8268..abbb734a1d 100644 --- a/packages/core/src/index.gate.ts +++ b/packages/core/src/index.gate.ts @@ -716,6 +716,7 @@ export { export type { TaskDependencyMutation } from "./store.js"; export { findSameAgentDuplicates, + computeParentIntentClaimId, archiveAsSameAgentDuplicate, flagSameAgentDuplicate, flagTriageDuplicate, diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index 6f55d0ecc8..120098f1c2 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -739,6 +739,7 @@ export { export type { TaskDependencyMutation } from "./store.js"; export { findSameAgentDuplicates, + computeParentIntentClaimId, archiveAsSameAgentDuplicate, flagSameAgentDuplicate, flagTriageDuplicate, diff --git a/packages/core/src/store.ts b/packages/core/src/store.ts index a141320da0..21fe857674 100644 --- a/packages/core/src/store.ts +++ b/packages/core/src/store.ts @@ -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 { /** * FNXC:RuntimeTaskOrchestrationAsync 2026-06-24-13:15: */ - public async createTaskBackend( input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; }, ): Promise { + public async createTaskBackend( input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; }, ): Promise { 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 { + 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 { return _createTaskInternalBackendImpl(this, input, title, resolvedWorkflowSteps, id, options); } @@ -874,13 +875,13 @@ export class TaskStore extends EventEmitter { public async _maybeAutoArchiveSameAgentDuplicateBackend( task: Task, input: TaskCreateInput, ): Promise { return _maybeAutoArchiveSameAgentDuplicateBackendImpl(this, task, input); } - async createTask( input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; } ): Promise { + async createTask( input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; } ): Promise { return createTaskImpl(this, input, options); } async createTaskWithReservedId( input: TaskCreateInput, options: { taskId: string; createdAt?: string; updatedAt?: string; prompt?: string; applyDefaultWorkflowSteps?: boolean; invokeTaskCreatedHook?: boolean; }, ): Promise { 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 { + 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 { /* 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 { async findRecentTasksByContentFingerprint( fingerprint: string, options?: { windowMs?: number; includeArchived?: boolean }, ): Promise { return findRecentTasksByContentFingerprintImpl(this, fingerprint, options); } + async findRecentTasksBySourceParentTaskId(sourceParentTaskId: string, options?: { windowMs?: number }): Promise { + 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 { diff --git a/packages/core/src/task-store/remaining-ops-6.ts b/packages/core/src/task-store/remaining-ops-6.ts index c5aba1636d..9158f2e2a1 100644 --- a/packages/core/src/task-store/remaining-ops-6.ts +++ b/packages/core/src/task-store/remaining-ops-6.ts @@ -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 { + 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))); + } + 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 }, diff --git a/packages/core/src/task-store/task-creation.ts b/packages/core/src/task-store/task-creation.ts index a89e3892ad..08b4c69744 100644 --- a/packages/core/src/task-store/task-creation.ts +++ b/packages/core/src/task-store/task-creation.ts @@ -44,7 +44,7 @@ function ensureSqliteProposalClaimUniqueness(store: TaskStore): void { ); } -export async function createTaskBackendImpl(store: TaskStore, input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; },): Promise { +export async function createTaskBackendImpl(store: TaskStore, input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; },): Promise { 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 { +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 { 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; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; }): Promise { +export async function createTaskImpl(store: TaskStore, input: TaskCreateInput, options?: { onSummarize?: (description: string) => Promise; settings?: { autoSummarizeTitles?: boolean }; invokeTaskCreatedHook?: boolean; onProposalClaimConflict?: (task: Task) => void; }): Promise { // 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 { +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 { 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)}`); } } - diff --git a/packages/engine/src/__tests__/agent-tools-delegation.test.ts b/packages/engine/src/__tests__/agent-tools-delegation.test.ts index 46af001526..02a31cd0d7 100644 --- a/packages/engine/src/__tests__/agent-tools-delegation.test.ts +++ b/packages/engine/src/__tests__/agent-tools-delegation.test.ts @@ -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 { return { @@ -14,6 +14,7 @@ function createMockAgentStore(overrides: Partial = {}): AgentStore { function createMockTaskStore(overrides: Partial = {}): 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", diff --git a/packages/engine/src/__tests__/heartbeat-scheduler.test.ts b/packages/engine/src/__tests__/heartbeat-scheduler.test.ts index bbd84bc62c..5cb7a8d2ba 100644 --- a/packages/engine/src/__tests__/heartbeat-scheduler.test.ts +++ b/packages/engine/src/__tests__/heartbeat-scheduler.test.ts @@ -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", () => { }); }); }); - diff --git a/packages/engine/src/agent-heartbeat.ts b/packages/engine/src/agent-heartbeat.ts index d6011ad1b3..05b23279f2 100644 --- a/packages/engine/src/agent-heartbeat.ts +++ b/packages/engine/src/agent-heartbeat.ts @@ -3718,15 +3718,25 @@ export class HeartbeatMonitor { execute: async (id: string, params: Static, 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 }); diff --git a/packages/engine/src/agent-tools.ts b/packages/engine/src/agent-tools.ts index 782c813162..3913bb2858 100644 --- a/packages/engine/src/agent-tools.ts +++ b/packages/engine/src/agent-tools.ts @@ -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:")) { diff --git a/packages/engine/src/index.ts b/packages/engine/src/index.ts index 7e355cbb7b..39dfce2d62 100644 --- a/packages/engine/src/index.ts +++ b/packages/engine/src/index.ts @@ -60,6 +60,7 @@ export { createWorkflowSettingsTool, createTraitListTool, createWorkflowAuthoringTools, + createAgentTask, taskCreateParams, taskListParams, taskShowParams,