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:
gsxdsm
2026-07-18 14:18:25 -07:00
committed by GitHub
parent 89b42ed1d6
commit 9bc0eb6943
20 changed files with 458 additions and 43 deletions

View 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.

View File

@@ -78,6 +78,7 @@ vi.mock("@fusion/engine", () => ({
installBaselineArchiveWorktreeDisposer: vi.fn(), installBaselineArchiveWorktreeDisposer: vi.fn(),
...workflowAuthoringEngineMock, ...workflowAuthoringEngineMock,
createFnAgent: vi.fn(), createFnAgent: vi.fn(),
createAgentTask: vi.fn(),
fetchWebContent: vi.fn(), fetchWebContent: vi.fn(),
emitGoalRetrievalAudit: vi.fn(), emitGoalRetrievalAudit: vi.fn(),
createWorkflowAuthoringTools: vi.fn(() => ({})), createWorkflowAuthoringTools: vi.fn(() => ({})),

View File

@@ -38,6 +38,7 @@ vi.mock("@fusion/dashboard", () => {
vi.mock("@fusion/engine", () => ({ vi.mock("@fusion/engine", () => ({
installBaselineArchiveWorktreeDisposer: vi.fn(), installBaselineArchiveWorktreeDisposer: vi.fn(),
createFnAgent: vi.fn(), createFnAgent: vi.fn(),
createAgentTask: vi.fn(),
fetchWebContent: vi.fn(), fetchWebContent: vi.fn(),
assertNoSecretPlaintext: vi.fn(), assertNoSecretPlaintext: vi.fn(),
emitGoalRetrievalAudit: vi.fn(), emitGoalRetrievalAudit: vi.fn(),

View File

@@ -17,6 +17,7 @@ vi.mock("@fusion/engine", () => ({
installBaselineArchiveWorktreeDisposer: vi.fn(), installBaselineArchiveWorktreeDisposer: vi.fn(),
...workflowAuthoringEngineMock, ...workflowAuthoringEngineMock,
createFnAgent: vi.fn(), createFnAgent: vi.fn(),
createAgentTask: vi.fn(),
fetchWebContent: fetchWebContentMock, fetchWebContent: fetchWebContentMock,
emitGoalRetrievalAudit: vi.fn(), emitGoalRetrievalAudit: vi.fn(),
createWorkflowAuthoringTools: vi.fn(() => ({})), createWorkflowAuthoringTools: vi.fn(() => ({})),

View File

@@ -362,6 +362,21 @@ legacyDescribe("fn pi extension (legacy exhaustive suite)", () => {
expect(result.details.priority).toBe("normal"); 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 () => { it("creates a task with explicit priority", async () => {
const tool = api.tools.get("fn_task_create")!; const tool = api.tools.get("fn_task_create")!;
const result = await tool.execute( const result = await tool.execute(

View File

@@ -69,6 +69,7 @@ import {
isInReviewMissingWorktreeSessionStartFailure, isInReviewMissingWorktreeSessionStartFailure,
normalizeAgentLogPaging, normalizeAgentLogPaging,
renderAgentLogEntries, renderAgentLogEntries,
createAgentTask,
} from "@fusion/engine"; } from "@fusion/engine";
import * as dashboard from "@fusion/dashboard"; import * as dashboard from "@fusion/dashboard";
import { resolve, relative, isAbsolute, sep, basename, extname, join } from "node:path"; 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: 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. 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 projectSettingsForGate = await store.getSettings();
const callerIsEphemeral = await isEphemeralCallerAgent(ctx.cwd ?? process.cwd(), fnCtx.agentId); const callerIsEphemeral = await isEphemeralCallerAgent(ctx.cwd ?? process.cwd(), fnCtx.agentId);
if (callerIsEphemeral) { if (callerIsEphemeral) {
@@ -1213,13 +1214,13 @@ export default function kbExtension(pi: ExtensionAPI) {
); );
const workflowId = params.workflow_id?.trim() || undefined; const workflowId = params.workflow_id?.trim() || undefined;
const task = await store.createTask({ const { task, wasDuplicate } = await createAgentTask(store, {
description: params.description.trim(), description: params.description.trim(),
dependencies: params.depends, dependencies: params.depends,
assignedAgentId: normalizedAgentId === null ? undefined : normalizedAgentId, assignedAgentId: normalizedAgentId === null ? undefined : normalizedAgentId,
priority: params.priority as TaskPriority | undefined, priority: params.priority as TaskPriority | undefined,
...(workflowId ? { workflowId } : {}), ...(workflowId ? { workflowId } : {}),
source: { sourceType: "api" }, source: { sourceType: "api", sourceAgentId: fnCtx.agentId, sourceParentTaskId: fnCtx.taskId },
githubTracking: resolvedTracking.enabled githubTracking: resolvedTracking.enabled
? { ? {
enabled: true, enabled: true,
@@ -1228,7 +1229,7 @@ export default function kbExtension(pi: ExtensionAPI) {
: {}), : {}),
} }
: undefined, : undefined,
}); }, { rootDir: ctx.cwd, sourceAgentId: fnCtx.agentId, sourceTaskId: fnCtx.taskId });
const label = const label =
task.description.length > 80 task.description.length > 80
@@ -1248,7 +1249,7 @@ export default function kbExtension(pi: ExtensionAPI) {
{ {
type: "text", type: "text",
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` + `Column: ${task.column}\n` +
(task.dependencies.length (task.dependencies.length
? `Dependencies: ${task.dependencies.join(", ")}\n` ? `Dependencies: ${task.dependencies.join(", ")}\n`
@@ -1262,6 +1263,7 @@ export default function kbExtension(pi: ExtensionAPI) {
], ],
details: { details: {
taskId: task.id, taskId: task.id,
wasDuplicate,
column: task.column, column: task.column,
dependencies: task.dependencies, dependencies: task.dependencies,
assignedAgentId: task.assignedAgentId, assignedAgentId: task.assignedAgentId,

View File

@@ -122,6 +122,25 @@ describe("runDeterministicDuplicateGuard", () => {
secondResult.releaseLock(); 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 () => { it("allows concurrent proceed without shared lock scope", async () => {
const { store } = makeStore(); const { store } = makeStore();
const [a, b] = await Promise.all([ const [a, b] = await Promise.all([

View File

@@ -1,6 +1,6 @@
import { describe, expect, it, vi } from "vitest"; 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"; import type { TaskStore } from "../store.js";
describe("findSameAgentDuplicates", () => { describe("findSameAgentDuplicates", () => {
@@ -79,6 +79,49 @@ describe("findSameAgentDuplicates", () => {
expect(matches[0]?.id).toBe("FN-5544"); 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", () => { it("does not match sibling with different parent task", () => {
const matches = findSameAgentDuplicates( const matches = findSameAgentDuplicates(
{ {

View File

@@ -15,6 +15,8 @@ export interface DeterministicGuardOptions {
acknowledgedDuplicates?: readonly string[]; acknowledgedDuplicates?: readonly string[];
bypass?: boolean; bypass?: boolean;
logger?: { warn(msg: string, data?: Record<string, unknown>): void }; logger?: { warn(msg: string, data?: Record<string, unknown>): void };
/** Serialize related creates even when their exact-content fingerprints differ. */
serializationKey?: string;
} }
export interface DeterministicGuardOutcome { export interface DeterministicGuardOutcome {
@@ -67,8 +69,16 @@ export async function runDeterministicDuplicateGuard(
return { action: "proceed", fingerprint, releaseLock: noop }; return { action: "proceed", fingerprint, releaseLock: noop };
} }
const lockKey = `${opts.lockScope}:${fingerprint}`; const lockKey = `${opts.lockScope}:${opts.serializationKey ?? fingerprint}`;
const existingLock = deterministicGuardLocks.get(lockKey); 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) { if (existingLock) {
try { try {
await existingLock; await existingLock;
@@ -78,24 +88,17 @@ export async function runDeterministicDuplicateGuard(
contentFingerprint: fingerprint, contentFingerprint: fingerprint,
error: error instanceof Error ? error.message : String(error), 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 = () => { const releaseLock = () => {
if (releaseCalled) { if (releaseCalled) {
return; return;
} }
releaseCalled = true; releaseCalled = true;
resolveGate?.(); resolveGate?.();
deterministicGuardLocks.delete(lockKey); if (deterministicGuardLocks.get(lockKey) === gate) deterministicGuardLocks.delete(lockKey);
}; };
try { try {

View File

@@ -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 { ColumnId } from "./types.js";
import type { TaskStore } from "./store.js"; import type { TaskStore } from "./store.js";
@@ -34,6 +34,50 @@ export interface SameAgentDuplicateMatch {
allowResurrection?: boolean; 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. * 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])); 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); const candidate = metadataById.get(match.id);
return { return {
id: match.id, id: match.id,
@@ -86,6 +130,36 @@ export function findSameAgentDuplicates(
allowResurrection: candidate?.allowResurrection, 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( export async function archiveAsSameAgentDuplicate(

View File

@@ -716,6 +716,7 @@ export {
export type { TaskDependencyMutation } from "./store.js"; export type { TaskDependencyMutation } from "./store.js";
export { export {
findSameAgentDuplicates, findSameAgentDuplicates,
computeParentIntentClaimId,
archiveAsSameAgentDuplicate, archiveAsSameAgentDuplicate,
flagSameAgentDuplicate, flagSameAgentDuplicate,
flagTriageDuplicate, flagTriageDuplicate,

View File

@@ -739,6 +739,7 @@ export {
export type { TaskDependencyMutation } from "./store.js"; export type { TaskDependencyMutation } from "./store.js";
export { export {
findSameAgentDuplicates, findSameAgentDuplicates,
computeParentIntentClaimId,
archiveAsSameAgentDuplicate, archiveAsSameAgentDuplicate,
flagSameAgentDuplicate, flagSameAgentDuplicate,
flagTriageDuplicate, flagTriageDuplicate,

View File

@@ -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 { 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 { 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 { 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 { 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 { 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"; 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: * 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); return createTaskBackendImpl(this, input, options);
} }
/** /**
* FNXC:RuntimeTaskOrchestrationAsync 2026-06-24-13:25: * 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); 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> { public async _maybeAutoArchiveSameAgentDuplicateBackend( task: Task, input: TaskCreateInput, ): Promise<void> {
return _maybeAutoArchiveSameAgentDuplicateBackendImpl(this, task, input); 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); return createTaskImpl(this, input, options);
} }
async createTaskWithReservedId( input: TaskCreateInput, options: { taskId: string; createdAt?: string; updatedAt?: string; prompt?: string; applyDefaultWorkflowSteps?: boolean; invokeTaskCreatedHook?: boolean; }, ): Promise<Task> { async createTaskWithReservedId( input: TaskCreateInput, options: { taskId: string; createdAt?: string; updatedAt?: string; prompt?: string; applyDefaultWorkflowSteps?: boolean; invokeTaskCreatedHook?: boolean; }, ): Promise<Task> {
return createTaskWithReservedIdImpl(this, input, options); 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: FNXC:SqliteFinalRemoval 2026-06-25-10:35:
Route to the async backend variant when the store is in backend mode so 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[]> { async findRecentTasksByContentFingerprint( fingerprint: string, options?: { windowMs?: number; includeArchived?: boolean }, ): Promise<Task[]> {
return findRecentTasksByContentFingerprintImpl(this, fingerprint, options); 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. */ /** 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[]> { public async clearNearDuplicateReferencesTo( canonicalId: string, inactiveState: { column?: ColumnId | null; deletedAt?: string | null; reason: string }, ): Promise<Task[]> {

View File

@@ -727,6 +727,35 @@ export async function findRecentTasksByContentFingerprintImpl(store: TaskStore,
return rows.map((row) => store.rowToTask(row)); 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, export async function clearNearDuplicateReferencesToFailSoftImpl(store: TaskStore,
canonicalId: string, canonicalId: string,
inactiveState: { column?: ColumnId | null; deletedAt?: string | null; reason: string }, inactiveState: { column?: ColumnId | null; deletedAt?: string | null; reason: string },

View File

@@ -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()) { if (!input.description?.trim()) {
throw new Error("Description is required and cannot be empty"); throw new Error("Description is required and cannot be empty");
} }
@@ -205,7 +205,7 @@ export async function createTaskBackendImpl(store: TaskStore, input: TaskCreateI
title, title,
resolvedWorkflowSteps, resolvedWorkflowSteps,
reservation.taskId, reservation.taskId,
{ invokeTaskCreatedHook: shouldInvokeTaskCreatedHook && !hasPendingSummarization, resolvedEntryColumn }, { invokeTaskCreatedHook: shouldInvokeTaskCreatedHook && !hasPendingSummarization, resolvedEntryColumn, onProposalClaimConflict: options?.onProposalClaimConflict },
); );
await allocator.commitDistributedTaskIdReservation({ await allocator.commitDistributedTaskIdReservation({
reservationId: reservation.reservationId, reservationId: reservation.reservationId,
@@ -289,7 +289,7 @@ export async function createTaskBackendImpl(store: TaskStore, input: TaskCreateI
return task; 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 layer = store.asyncLayer!;
const now = options?.createdAt ?? new Date().toISOString(); const now = options?.createdAt ?? new Date().toISOString();
const normalizedTitle = normalizeTitleForTaskId(title, id); const normalizedTitle = normalizeTitleForTaskId(title, id);
@@ -388,7 +388,10 @@ export async function _createTaskInternalBackendImpl(store: TaskStore, input: Ta
*/ */
if (input.proposalClaimId && isTaskIdConflictError(error)) { if (input.proposalClaimId && isTaskIdConflictError(error)) {
const existing = (await store.listTasks()).find((candidate) => candidate.proposalClaimId === input.proposalClaimId); 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)) { if (isTaskIdConflictError(error)) {
throw new Error(`Task ID already exists: ${task.id}`); throw new Error(`Task ID already exists: ${task.id}`);
@@ -445,7 +448,7 @@ export async function _createTaskInternalBackendImpl(store: TaskStore, input: Ta
return task; 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: // FNXC:RuntimeTaskOrchestrationAsync 2026-06-24-13:10:
// Backend-mode createTask: delegates to createTaskBackend which uses the // Backend-mode createTask: delegates to createTaskBackend which uses the
// async DistributedTaskIdAllocator (now wired for backend mode) and 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) { if (input.proposalClaimId) {
ensureSqliteProposalClaimUniqueness(store); ensureSqliteProposalClaimUniqueness(store);
const existing = (await store.listTasks()).find((task) => task.proposalClaimId === input.proposalClaimId); 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 ?? []); const selfDefeatingDep = detectSelfDefeatingDependency(input.title, input.dependencies ?? []);
@@ -626,7 +632,7 @@ export async function createTaskImpl(store: TaskStore, input: TaskCreateInput, o
title, title,
resolvedWorkflowSteps, resolvedWorkflowSteps,
taskId, 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); await store.cleanupOrphanedMaterializedSteps(pendingWorkflowSelection?.stepIds);
if (input.proposalClaimId && isTaskIdConflictError(err)) { if (input.proposalClaimId && isTaskIdConflictError(err)) {
const existing = (await store.listTasks()).find((candidate) => candidate.proposalClaimId === input.proposalClaimId); const existing = (await store.listTasks()).find((candidate) => candidate.proposalClaimId === input.proposalClaimId);
if (existing) return existing; if (existing) {
options?.onProposalClaimConflict?.(existing);
return existing;
}
} }
throw err; throw err;
} }
@@ -869,7 +878,7 @@ export async function createTaskWithReservedIdImpl(store: TaskStore, input: Task
return createdTask; 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(); const now = options?.createdAt ?? new Date().toISOString();
// FN-5077: null normalized titles are treated as "no title" and allow standard fallback/summarization behavior. // FN-5077: null normalized titles are treated as "no title" and allow standard fallback/summarization behavior.
const normalizedTitle = normalizeTitleForTaskId(title, id); 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)}`); storeLog.warn(`FN-4892 same-agent duplicate intake failed open for ${task.id}: ${getErrorMessage(error)}`);
} }
} }

View File

@@ -1,6 +1,6 @@
import { describe, it, expect, vi, beforeEach } from "vitest"; import { describe, it, expect, vi, beforeEach } from "vitest";
import type { Agent, AgentStore, TaskStore, Task } from "@fusion/core"; 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 { function createMockAgentStore(overrides: Partial<AgentStore> = {}): AgentStore {
return { return {
@@ -14,6 +14,7 @@ function createMockAgentStore(overrides: Partial<AgentStore> = {}): AgentStore {
function createMockTaskStore(overrides: Partial<TaskStore> = {}): TaskStore { function createMockTaskStore(overrides: Partial<TaskStore> = {}): TaskStore {
return { return {
getSettings: vi.fn().mockResolvedValue({ autoSummarizeTitles: false }), getSettings: vi.fn().mockResolvedValue({ autoSummarizeTitles: false }),
findRecentTasksBySourceParentTaskId: vi.fn().mockResolvedValue([]),
findRecentTasksByContentFingerprint: vi.fn().mockResolvedValue([]), findRecentTasksByContentFingerprint: vi.fn().mockResolvedValue([]),
updateTask: vi.fn(), updateTask: vi.fn(),
moveTask: vi.fn(), moveTask: vi.fn(),
@@ -301,6 +302,93 @@ describe("createDelegateTaskTool", () => {
expect(taskStore.moveTask).not.toHaveBeenCalled(); 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 () => { it("carries delegation routing onto the reconcile canonical task", async () => {
const created = { const created = {
id: "FN-new", id: "FN-new",

View File

@@ -2988,6 +2988,7 @@ describe("HeartbeatTriggerScheduler", () => {
// Create a minimal mock TaskStore // Create a minimal mock TaskStore
const mockTaskStore = { const mockTaskStore = {
createTask: vi.fn().mockResolvedValue({ id: "FN-200", description: "New task created", dependencies: [] }), createTask: vi.fn().mockResolvedValue({ id: "FN-200", description: "New task created", dependencies: [] }),
findRecentTasksBySourceParentTaskId: vi.fn().mockResolvedValue([]),
logEntry: vi.fn().mockResolvedValue({}), logEntry: vi.fn().mockResolvedValue({}),
getTask: vi.fn().mockResolvedValue({ getTask: vi.fn().mockResolvedValue({
id: "FN-001", 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 () => { it("createHeartbeatTools works without runContext (backward compat)", async () => {
// Create a minimal mock TaskStore // Create a minimal mock TaskStore
const mockTaskStore = { const mockTaskStore = {
createTask: vi.fn().mockResolvedValue({ id: "FN-NEW", description: "New task" }), createTask: vi.fn().mockResolvedValue({ id: "FN-NEW", description: "New task" }),
findRecentTasksBySourceParentTaskId: vi.fn().mockResolvedValue([]),
logEntry: vi.fn().mockResolvedValue({}), logEntry: vi.fn().mockResolvedValue({}),
getTask: vi.fn().mockResolvedValue({ getTask: vi.fn().mockResolvedValue({
id: "FN-001", id: "FN-001",
@@ -3268,4 +3325,3 @@ describe("HeartbeatTriggerScheduler", () => {
}); });
}); });
}); });

View File

@@ -3718,15 +3718,25 @@ export class HeartbeatMonitor {
execute: async (id: string, params: Static<typeof taskCreateParams>, signal, onUpdate, ctx) => { execute: async (id: string, params: Static<typeof taskCreateParams>, signal, onUpdate, ctx) => {
const result = await baseCreateTool.execute(id, params, 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 // Log agent link on the created task with run context for correlation
try { 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) { } catch (taskCreateLogErr) {
heartbeatLog.warn(`Task ${createdTaskId} agent-link log failed: ${taskCreateLogErr instanceof Error ? taskCreateLogErr.message : String(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) // Audit trail: record task creation (FN-1404)
await audit?.database({ type: "task:create", target: createdTaskId }); await audit?.database({ type: "task:create", target: createdTaskId });

View File

@@ -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 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 { 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 { 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 { ResearchOrchestrator } from "./research-orchestrator.js";
import { ResearchProviderRegistry } from "./research/provider-registry.js"; import { ResearchProviderRegistry } from "./research/provider-registry.js";
import { ResearchStepRunner } from "./research-step-runner.js"; import { ResearchStepRunner } from "./research-step-runner.js";
@@ -973,6 +973,16 @@ export async function createAgentTask(
? await store.getSettings() ? await store.getSettings()
: {} as Settings; : {} as Settings;
const rootDir = options?.rootDir; 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, { const guard = await runDeterministicDuplicateGuard(store, {
title: input.title, title: input.title,
description: input.description, description: input.description,
@@ -980,6 +990,7 @@ export async function createAgentTask(
lockScope: rootDir ?? store.getRootDir?.() ?? "agent-tools", lockScope: rootDir ?? store.getRootDir?.() ?? "agent-tools",
bypass: options?.bypassDuplicateCheck === true, bypass: options?.bypassDuplicateCheck === true,
acknowledgedDuplicates: options?.acknowledgedDuplicates, acknowledgedDuplicates: options?.acknowledgedDuplicates,
serializationKey: sourceParentTaskId ? `parent:${sourceParentTaskId}` : undefined,
logger: log, 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 = { const sourceMetadata = {
...(input.source?.sourceMetadata ?? {}), ...(effectiveSource?.sourceMetadata ?? {}),
...(guard.fingerprint ? { contentFingerprint: guard.fingerprint } : {}), ...(guard.fingerprint ? { contentFingerprint: guard.fingerprint } : {}),
}; };
const nextSource = input.source const nextSource = effectiveSource
? { ? {
...input.source, ...effectiveSource,
sourceMetadata: Object.keys(sourceMetadata).length > 0 ? sourceMetadata : undefined, sourceMetadata: Object.keys(sourceMetadata).length > 0 ? sourceMetadata : undefined,
} }
: undefined; : undefined;
@@ -1014,6 +1056,11 @@ export async function createAgentTask(
input.githubTracking?.enabled !== false && resolvedTracking.enabled; input.githubTracking?.enabled !== false && resolvedTracking.enabled;
const createInput: TaskCreateInput = { const createInput: TaskCreateInput = {
...input, ...input,
proposalClaimId: input.proposalClaimId ?? (
options?.bypassDuplicateCheck === true
? undefined
: computeParentIntentClaimId({ title: input.title, description: input.description, sourceParentTaskId }) ?? undefined
),
summarize: !input.title?.trim() ? true : undefined, summarize: !input.title?.trim() ? true : undefined,
source: nextSource, source: nextSource,
githubTracking: shouldPrefillGithubTrackingEnabled githubTracking: shouldPrefillGithubTrackingEnabled
@@ -1027,8 +1074,10 @@ export async function createAgentTask(
: input.githubTracking, : input.githubTracking,
}; };
let proposalClaimConflict = false;
const createdTask = await store.createTask(createInput, { const createdTask = await store.createTask(createInput, {
settings, settings,
onProposalClaimConflict: () => { proposalClaimConflict = true; },
}); });
const reconcile = await reconcileDeterministicDuplicate(store, { const reconcile = await reconcileDeterministicDuplicate(store, {
@@ -1038,10 +1087,12 @@ export async function createAgentTask(
}); });
return { return {
task: reconcile.outcome === "archived" task: proposalClaimConflict
? await carryCanonicalTaskRouting(store, createdTask, input)
: reconcile.outcome === "archived"
? await carryCanonicalTaskRouting(store, reconcile.canonical, input) ? await carryCanonicalTaskRouting(store, reconcile.canonical, input)
: reconcile.canonical, : reconcile.canonical,
wasDuplicate: reconcile.outcome === "archived", wasDuplicate: proposalClaimConflict || reconcile.outcome === "archived",
}; };
} finally { } finally {
guard.releaseLock(); guard.releaseLock();
@@ -1137,7 +1188,7 @@ export function createTaskCreateTool(
type: "text" as const, type: "text" as const,
text: `${wasDuplicate ? "Linked existing" : "Created"} ${task.id}: ${params.description}${deps}${workflow}`, text: `${wasDuplicate ? "Linked existing" : "Created"} ${task.id}: ${params.description}${deps}${workflow}`,
}], }],
details: { taskId: task.id }, details: { taskId: task.id, wasDuplicate },
}; };
} catch (err) { } catch (err) {
if (err instanceof Error && err.message.startsWith("Task ID already exists:")) { if (err instanceof Error && err.message.startsWith("Task ID already exists:")) {

View File

@@ -60,6 +60,7 @@ export {
createWorkflowSettingsTool, createWorkflowSettingsTool,
createTraitListTool, createTraitListTool,
createWorkflowAuthoringTools, createWorkflowAuthoringTools,
createAgentTask,
taskCreateParams, taskCreateParams,
taskListParams, taskListParams,
taskShowParams, taskShowParams,