From 632ae67de3fa33487e1ad513043cf740cb32bb1a Mon Sep 17 00:00:00 2001 From: Albert Alises Date: Mon, 27 Apr 2026 17:17:54 +0200 Subject: [PATCH] fix(ai-builder): Hide and reap intermediate AI-created workflows (#29066) --- .../src/schemas/instance-ai.schema.ts | 7 + .../entities/ai-builder-temporary-workflow.ts | 28 ++ packages/@n8n/db/src/entities/index.ts | 3 + ...3-CreateAiBuilderTemporaryWorkflowTable.ts | 26 ++ .../db/src/migrations/postgresdb/index.ts | 2 + .../@n8n/db/src/migrations/sqlite/index.ts | 2 + .../__tests__/workflow.repository.test.ts | 23 +- ...i-builder-temporary-workflow.repository.ts | 37 +++ packages/@n8n/db/src/repositories/index.ts | 1 + .../src/repositories/workflow.repository.ts | 18 ++ .../__tests__/planned-task-storage.test.ts | 70 +++++ .../__tests__/parse-file.tool.test.ts | 2 + .../build-workflow-agent.tool.test.ts | 36 ++- .../build-workflow-agent.tool.ts | 70 ++++- .../instance-ai/src/tools/workflows.tool.ts | 5 +- .../__tests__/setup-workflow.service.test.ts | 2 + .../tools/workflows/build-workflow.tool.ts | 9 +- .../tools/workflows/submit-workflow.tool.ts | 9 +- packages/@n8n/instance-ai/src/types.ts | 25 +- ...stance-ai.adapter.service.security.test.ts | 3 + .../instance-ai.adapter.service.test.ts | 134 ++++++++- .../instance-ai.adapter.service.ts | 68 ++++- .../instance-ai/instance-ai.service.ts | 258 +++++++++++++++--- .../frontend/@n8n/i18n/src/locales/en.json | 1 + .../components/AgentActivityTree.vue | 1 + .../instanceAi/components/AgentTimeline.vue | 1 + .../ai/instanceAi/components/ArtifactCard.vue | 22 +- .../components/InstanceAiArtifactsPanel.vue | 24 +- .../ai/instanceAi/instanceAi.store.ts | 12 + .../ai/instanceAi/useResourceRegistry.ts | 17 ++ 30 files changed, 852 insertions(+), 64 deletions(-) create mode 100644 packages/@n8n/db/src/entities/ai-builder-temporary-workflow.ts create mode 100644 packages/@n8n/db/src/migrations/common/1777281990043-CreateAiBuilderTemporaryWorkflowTable.ts create mode 100644 packages/@n8n/db/src/repositories/ai-builder-temporary-workflow.repository.ts create mode 100644 packages/@n8n/instance-ai/src/storage/__tests__/planned-task-storage.test.ts diff --git a/packages/@n8n/api-types/src/schemas/instance-ai.schema.ts b/packages/@n8n/api-types/src/schemas/instance-ai.schema.ts index c0161d32439..cdac3ec0b90 100644 --- a/packages/@n8n/api-types/src/schemas/instance-ai.schema.ts +++ b/packages/@n8n/api-types/src/schemas/instance-ai.schema.ts @@ -112,6 +112,13 @@ export const runStartPayloadSchema = z.object({ export const runFinishPayloadSchema = z.object({ status: instanceAiRunStatusSchema, reason: z.string().optional(), + /** + * Workflow IDs the run-finish reap soft-deleted — intermediate + * stepping-stones the agent created but never promoted to the main + * deliverable. Surfaced to the UI so the artifacts panel can dim these + * entries and label them as archived. + */ + archivedWorkflowIds: z.array(z.string()).optional(), }); export const agentSpawnedTargetResourceSchema = z.object({ diff --git a/packages/@n8n/db/src/entities/ai-builder-temporary-workflow.ts b/packages/@n8n/db/src/entities/ai-builder-temporary-workflow.ts new file mode 100644 index 00000000000..5a54443e2f0 --- /dev/null +++ b/packages/@n8n/db/src/entities/ai-builder-temporary-workflow.ts @@ -0,0 +1,28 @@ +import { + Column, + Entity, + Index, + JoinColumn, + ManyToOne, + PrimaryColumn, + Relation, +} from '@n8n/typeorm'; + +import { WithTimestamps } from './abstract-entity'; +import type { WorkflowEntity } from './workflow-entity'; + +@Entity({ name: 'ai_builder_temporary_workflow' }) +export class AiBuilderTemporaryWorkflow extends WithTimestamps { + @PrimaryColumn({ type: 'varchar', length: 36 }) + workflowId: string; + + @Index() + @Column({ type: 'uuid' }) + threadId: string; + + @ManyToOne('WorkflowEntity', { + onDelete: 'CASCADE', + }) + @JoinColumn({ name: 'workflowId' }) + workflow: Relation; +} diff --git a/packages/@n8n/db/src/entities/index.ts b/packages/@n8n/db/src/entities/index.ts index f05294750fc..f7753a800e4 100644 --- a/packages/@n8n/db/src/entities/index.ts +++ b/packages/@n8n/db/src/entities/index.ts @@ -1,3 +1,4 @@ +import { AiBuilderTemporaryWorkflow } from './ai-builder-temporary-workflow'; import { AnnotationTagEntity } from './annotation-tag-entity.ee'; import { AnnotationTagMapping } from './annotation-tag-mapping.ee'; import { ApiKey } from './api-key'; @@ -46,6 +47,7 @@ import { WorkflowTagMapping } from './workflow-tag-mapping'; export { InvalidAuthToken, + AiBuilderTemporaryWorkflow, ProcessedData, Settings, Variables, @@ -94,6 +96,7 @@ export { export const entities = { InvalidAuthToken, + AiBuilderTemporaryWorkflow, ProcessedData, Settings, Variables, diff --git a/packages/@n8n/db/src/migrations/common/1777281990043-CreateAiBuilderTemporaryWorkflowTable.ts b/packages/@n8n/db/src/migrations/common/1777281990043-CreateAiBuilderTemporaryWorkflowTable.ts new file mode 100644 index 00000000000..0773a38a5c1 --- /dev/null +++ b/packages/@n8n/db/src/migrations/common/1777281990043-CreateAiBuilderTemporaryWorkflowTable.ts @@ -0,0 +1,26 @@ +import type { MigrationContext, ReversibleMigration } from '../migration-types'; + +export class CreateAiBuilderTemporaryWorkflowTable1777281990043 implements ReversibleMigration { + async up({ schemaBuilder: { createTable, column } }: MigrationContext) { + await createTable('ai_builder_temporary_workflow') + .withColumns( + column('workflowId').varchar(36).primary.notNull, + column('threadId').uuid.notNull, + ) + .withForeignKey('workflowId', { + tableName: 'workflow_entity', + columnName: 'id', + onDelete: 'CASCADE', + }) + .withForeignKey('threadId', { + tableName: 'instance_ai_threads', + columnName: 'id', + onDelete: 'CASCADE', + }) + .withIndexOn('threadId').withTimestamps; + } + + async down({ schemaBuilder: { dropTable } }: MigrationContext) { + await dropTable('ai_builder_temporary_workflow'); + } +} diff --git a/packages/@n8n/db/src/migrations/postgresdb/index.ts b/packages/@n8n/db/src/migrations/postgresdb/index.ts index 8f124926a91..628c7c33328 100644 --- a/packages/@n8n/db/src/migrations/postgresdb/index.ts +++ b/packages/@n8n/db/src/migrations/postgresdb/index.ts @@ -164,6 +164,7 @@ import { CreateFavoritesTable1776150756000 } from '../common/1776150756000-Creat import { CreateDeploymentKeyTable1777000000000 } from '../common/1777000000000-CreateDeploymentKeyTable'; import { AddTracingContextToExecution1777045000000 } from '../common/1777045000000-AddTracingContextToExecution'; import { AddLangsmithIdsToInstanceAiRunSnapshots1777100000000 } from '../common/1777100000000-AddLangsmithIdsToInstanceAiRunSnapshots'; +import { CreateAiBuilderTemporaryWorkflowTable1777281990043 } from '../common/1777281990043-CreateAiBuilderTemporaryWorkflowTable'; import { AddExecutionDeduplicationKey1778000000000 } from '../common/1778000000000-AddExecutionDeduplicationKey'; import type { Migration } from '../migration-types'; @@ -335,4 +336,5 @@ export const postgresMigrations: Migration[] = [ AddLangsmithIdsToInstanceAiRunSnapshots1777100000000, AddExecutionDeduplicationKey1778000000000, AddTracingContextToExecution1777045000000, + CreateAiBuilderTemporaryWorkflowTable1777281990043, ]; diff --git a/packages/@n8n/db/src/migrations/sqlite/index.ts b/packages/@n8n/db/src/migrations/sqlite/index.ts index 36113389d36..fe04f1965bf 100644 --- a/packages/@n8n/db/src/migrations/sqlite/index.ts +++ b/packages/@n8n/db/src/migrations/sqlite/index.ts @@ -158,6 +158,7 @@ import { CreateFavoritesTable1776150756000 } from '../common/1776150756000-Creat import { CreateDeploymentKeyTable1777000000000 } from '../common/1777000000000-CreateDeploymentKeyTable'; import { AddTracingContextToExecution1777045000000 } from '../common/1777045000000-AddTracingContextToExecution'; import { AddLangsmithIdsToInstanceAiRunSnapshots1777100000000 } from '../common/1777100000000-AddLangsmithIdsToInstanceAiRunSnapshots'; +import { CreateAiBuilderTemporaryWorkflowTable1777281990043 } from '../common/1777281990043-CreateAiBuilderTemporaryWorkflowTable'; import { AddExecutionDeduplicationKey1778000000000 } from '../common/1778000000000-AddExecutionDeduplicationKey'; import type { Migration } from '../migration-types'; @@ -323,6 +324,7 @@ const sqliteMigrations: Migration[] = [ AddLangsmithIdsToInstanceAiRunSnapshots1777100000000, AddExecutionDeduplicationKey1778000000000, AddTracingContextToExecution1777045000000, + CreateAiBuilderTemporaryWorkflowTable1777281990043, ]; export { sqliteMigrations }; diff --git a/packages/@n8n/db/src/repositories/__tests__/workflow.repository.test.ts b/packages/@n8n/db/src/repositories/__tests__/workflow.repository.test.ts index 88c490b56f1..467ddee0857 100644 --- a/packages/@n8n/db/src/repositories/__tests__/workflow.repository.test.ts +++ b/packages/@n8n/db/src/repositories/__tests__/workflow.repository.test.ts @@ -2,7 +2,7 @@ import { GlobalConfig } from '@n8n/config'; import { In, type SelectQueryBuilder } from '@n8n/typeorm'; import { mock } from 'jest-mock-extended'; -import { WorkflowEntity } from '../../entities'; +import { AiBuilderTemporaryWorkflow, WorkflowEntity } from '../../entities'; import { mockEntityManager } from '../../utils/test-utils/mock-entity-manager'; import { mockInstance } from '../../utils/test-utils/mock-instance'; import { FolderRepository } from '../folder.repository'; @@ -32,8 +32,17 @@ describe('WorkflowRepository', () => { jest.resetAllMocks(); queryBuilder = mock>(); + const subQueryBuilder = mock>(); + subQueryBuilder.select.mockReturnThis(); + subQueryBuilder.from.mockReturnThis(); + subQueryBuilder.where.mockReturnThis(); + subQueryBuilder.getQuery.mockReturnValue( + '(SELECT 1 FROM "ai_builder_temporary_workflow" "aitw" WHERE aitw."workflowId" = workflow.id)', + ); + queryBuilder.where.mockReturnThis(); queryBuilder.andWhere.mockReturnThis(); + queryBuilder.subQuery.mockReturnValue(subQueryBuilder); queryBuilder.orWhere.mockReturnThis(); queryBuilder.select.mockReturnThis(); queryBuilder.addSelect.mockReturnThis(); @@ -56,6 +65,18 @@ describe('WorkflowRepository', () => { jest.spyOn(workflowRepository, 'createQueryBuilder').mockReturnValue(queryBuilder); }); + describe('applyAiBuilderTemporaryFilter', () => { + it('hides marker-table rows through a prefix-safe entity subquery', async () => { + await workflowRepository.getMany(['workflow1']); + + expect(queryBuilder.subQuery).toHaveBeenCalled(); + expect(queryBuilder.subQuery().from).toHaveBeenCalledWith(AiBuilderTemporaryWorkflow, 'aitw'); + expect(queryBuilder.andWhere).toHaveBeenCalledWith( + expect.stringContaining('NOT EXISTS (SELECT 1 FROM "ai_builder_temporary_workflow"'), + ); + }); + }); + describe('applyNameFilter', () => { it('should search for workflows containing all words from the query', async () => { const workflowIds = ['workflow1']; diff --git a/packages/@n8n/db/src/repositories/ai-builder-temporary-workflow.repository.ts b/packages/@n8n/db/src/repositories/ai-builder-temporary-workflow.repository.ts new file mode 100644 index 00000000000..ae96d5e7b8d --- /dev/null +++ b/packages/@n8n/db/src/repositories/ai-builder-temporary-workflow.repository.ts @@ -0,0 +1,37 @@ +import { Service } from '@n8n/di'; +import { DataSource, Repository } from '@n8n/typeorm'; +import type { EntityManager } from '@n8n/typeorm'; + +import { AiBuilderTemporaryWorkflow } from '../entities'; + +@Service() +export class AiBuilderTemporaryWorkflowRepository extends Repository { + constructor(dataSource: DataSource) { + super(AiBuilderTemporaryWorkflow, dataSource.manager); + } + + async mark( + workflowId: string, + threadId: string, + entityManager: EntityManager = this.manager, + ): Promise { + await entityManager.upsert(AiBuilderTemporaryWorkflow, { workflowId, threadId }, [ + 'workflowId', + ]); + } + + async unmark(workflowId: string, entityManager: EntityManager = this.manager): Promise { + await entityManager.delete(AiBuilderTemporaryWorkflow, { workflowId }); + } + + async findByThread(threadId: string): Promise { + return await this.find({ + where: { threadId }, + select: ['workflowId', 'threadId'], + }); + } + + async existsForWorkflow(workflowId: string): Promise { + return await this.existsBy({ workflowId }); + } +} diff --git a/packages/@n8n/db/src/repositories/index.ts b/packages/@n8n/db/src/repositories/index.ts index 09fa5442cd9..872d0276fb6 100644 --- a/packages/@n8n/db/src/repositories/index.ts +++ b/packages/@n8n/db/src/repositories/index.ts @@ -1,5 +1,6 @@ export { AnnotationTagMappingRepository } from './annotation-tag-mapping.repository.ee'; export { AnnotationTagRepository } from './annotation-tag.repository.ee'; +export { AiBuilderTemporaryWorkflowRepository } from './ai-builder-temporary-workflow.repository'; export { ApiKeyRepository } from './api-key.repository'; export { AuthIdentityRepository } from './auth-identity.repository'; export { AuthProviderSyncHistoryRepository } from './auth-provider-sync-history.repository'; diff --git a/packages/@n8n/db/src/repositories/workflow.repository.ts b/packages/@n8n/db/src/repositories/workflow.repository.ts index 8791570ab77..481fba2d1df 100644 --- a/packages/@n8n/db/src/repositories/workflow.repository.ts +++ b/packages/@n8n/db/src/repositories/workflow.repository.ts @@ -18,6 +18,7 @@ import { SharedWorkflowRepository } from './shared-workflow.repository'; import { WorkflowHistoryRepository } from './workflow-history.repository'; import { WebhookEntity, + AiBuilderTemporaryWorkflow, TagEntity, WorkflowEntity, WorkflowTagMapping, @@ -882,6 +883,23 @@ export class WorkflowRepository extends Repository { this.applyParentFolderFilter(qb, filter); this.applyNodeTypesFilter(qb, filter); this.applyAvailableInMCPFilter(qb, filter); + this.applyAiBuilderTemporaryFilter(qb); + } + + /** + * Hide workflows the AI builder created and has not yet promoted to the + * main deliverable. The orchestrator clears the marker on the main at + * build-time and reaps the rest at run-finish, but in the window between + * create and reap, marked rows must not surface in the workflows list. + */ + private applyAiBuilderTemporaryFilter(qb: SelectQueryBuilder): void { + const markerSubquery = qb + .subQuery() + .select('1') + .from(AiBuilderTemporaryWorkflow, 'aitw') + .where('aitw."workflowId" = workflow.id') + .getQuery(); + qb.andWhere(`NOT EXISTS ${markerSubquery}`); } private applyAvailableInMCPFilter( diff --git a/packages/@n8n/instance-ai/src/storage/__tests__/planned-task-storage.test.ts b/packages/@n8n/instance-ai/src/storage/__tests__/planned-task-storage.test.ts new file mode 100644 index 00000000000..126ce609451 --- /dev/null +++ b/packages/@n8n/instance-ai/src/storage/__tests__/planned-task-storage.test.ts @@ -0,0 +1,70 @@ +import type { Memory } from '@mastra/memory'; + +import type { PlannedTaskGraph } from '../../types'; +import { PlannedTaskStorage } from '../planned-task-storage'; + +jest.mock('../thread-patch', () => ({ + patchThread: jest.fn( + ( + _memory: Memory, + opts: { + threadId: string; + update: (thread: { metadata?: Record }) => { + metadata: Record; + }; + }, + ) => { + const currentMetadata = metadataByThread.get(opts.threadId) ?? {}; + const next = opts.update({ metadata: currentMetadata }); + metadataByThread.set(opts.threadId, next.metadata); + }, + ), +})); + +const metadataByThread = new Map>(); + +function makeMemory(): Memory { + return { + getThreadById: jest.fn(({ threadId }: { threadId: string }) => ({ + id: threadId, + title: 'Test', + metadata: metadataByThread.get(threadId), + resourceId: 'res-1', + createdAt: new Date(), + updatedAt: new Date(), + })), + } as unknown as Memory; +} + +function baseGraph(): PlannedTaskGraph { + return { + planRunId: 'run-1', + status: 'active', + tasks: [ + { + id: 'task-1', + title: 'Build', + kind: 'build-workflow', + spec: 'spec', + deps: [], + status: 'running', + }, + ], + }; +} + +describe('PlannedTaskStorage', () => { + beforeEach(() => { + metadataByThread.clear(); + }); + + it('round-trips a graph through save -> get', async () => { + const storage = new PlannedTaskStorage(makeMemory()); + const graph = baseGraph(); + + await storage.save('thread-1', graph); + const loaded = await storage.get('thread-1'); + + expect(loaded?.tasks[0].id).toBe('task-1'); + }); +}); diff --git a/packages/@n8n/instance-ai/src/tools/attachments/__tests__/parse-file.tool.test.ts b/packages/@n8n/instance-ai/src/tools/attachments/__tests__/parse-file.tool.test.ts index 081b50de19e..d1d84aebf06 100644 --- a/packages/@n8n/instance-ai/src/tools/attachments/__tests__/parse-file.tool.test.ts +++ b/packages/@n8n/instance-ai/src/tools/attachments/__tests__/parse-file.tool.test.ts @@ -22,6 +22,8 @@ function createMockContext(overrides?: Partial): InstanceAiCo delete: jest.fn(), publish: jest.fn(), unpublish: jest.fn(), + clearAiTemporary: jest.fn(), + archiveIfAiTemporary: jest.fn(), }, executionService: { list: jest.fn(), diff --git a/packages/@n8n/instance-ai/src/tools/orchestration/__tests__/build-workflow-agent.tool.test.ts b/packages/@n8n/instance-ai/src/tools/orchestration/__tests__/build-workflow-agent.tool.test.ts index 238f19c5976..46b2ae4c75b 100644 --- a/packages/@n8n/instance-ai/src/tools/orchestration/__tests__/build-workflow-agent.tool.test.ts +++ b/packages/@n8n/instance-ai/src/tools/orchestration/__tests__/build-workflow-agent.tool.test.ts @@ -8,7 +8,7 @@ jest.mock('@mastra/core/mastra', () => ({ import type { SubmitWorkflowAttempt } from '../../workflows/submit-workflow.tool'; -const { resultFromPostStreamError } = +const { recordSuccessfulWorkflowBuilds, resultFromPostStreamError } = // eslint-disable-next-line @typescript-eslint/no-require-imports, @typescript-eslint/consistent-type-imports require('../build-workflow-agent.tool') as typeof import('../build-workflow-agent.tool'); @@ -128,3 +128,37 @@ describe('resultFromPostStreamError', () => { }); }); }); + +describe('recordSuccessfulWorkflowBuilds', () => { + it('records workflow IDs returned from successful build-workflow executions', async () => { + const onWorkflowId = jest.fn(); + const input = { prompt: 'build it' }; + const context = { runId: 'run-1' }; + const result = { success: true, workflowId: 'wf-main', displayName: 'Main' }; + const execute = jest.fn( + async (_input: unknown, _context?: unknown) => await Promise.resolve(result), + ); + const tool = { execute }; + + recordSuccessfulWorkflowBuilds(tool, onWorkflowId); + + await expect(tool.execute(input, context)).resolves.toBe(result); + expect(execute).toHaveBeenCalledWith(input, context); + expect(onWorkflowId).toHaveBeenCalledWith('wf-main'); + }); + + it('does not record failed or incomplete build-workflow results', async () => { + const onWorkflowId = jest.fn(); + const execute = jest + .fn() + .mockResolvedValueOnce({ success: false, workflowId: 'wf-failed' }) + .mockResolvedValueOnce({ success: true }); + const tool = { execute }; + + recordSuccessfulWorkflowBuilds(tool, onWorkflowId); + + await tool.execute({}); + await tool.execute({}); + expect(onWorkflowId).not.toHaveBeenCalled(); + }); +}); diff --git a/packages/@n8n/instance-ai/src/tools/orchestration/build-workflow-agent.tool.ts b/packages/@n8n/instance-ai/src/tools/orchestration/build-workflow-agent.tool.ts index a0a5966c173..a2979ba03f0 100644 --- a/packages/@n8n/instance-ai/src/tools/orchestration/build-workflow-agent.tool.ts +++ b/packages/@n8n/instance-ai/src/tools/orchestration/build-workflow-agent.tool.ts @@ -29,6 +29,7 @@ import { createVerifyBuiltWorkflowTool } from './verify-built-workflow.tool'; import { registerWithMastra } from '../../agent/register-with-mastra'; import { buildSubAgentBriefing } from '../../agent/sub-agent-briefing'; import { MAX_STEPS } from '../../constants/max-steps'; +import type { Logger } from '../../logger'; import { createLlmStepTraceHooks } from '../../runtime/resumable-stream-executor'; import { consumeStreamWithHitl } from '../../stream/consume-with-hitl'; import { @@ -37,7 +38,7 @@ import { mergeTraceRunInputs, withTraceParentContext, } from '../../tracing/langsmith-tracing'; -import type { BackgroundTaskResult, OrchestrationContext } from '../../types'; +import type { BackgroundTaskResult, InstanceAiContext, OrchestrationContext } from '../../types'; import { SDK_IMPORT_STATEMENT } from '../../workflow-builder/extract-code'; import type { TriggerType, WorkflowBuildOutcome } from '../../workflow-loop'; import type { BuilderWorkspace } from '../../workspace/builder-sandbox-factory'; @@ -61,6 +62,57 @@ function triggerLabel(nodeType: string): string { return short.replace(/Trigger$/i, '').toLowerCase() || short.toLowerCase(); } +/** + * Clear the AI-builder temporary marker from the build's main workflow so the + * run-finish reap leaves it alone. Best-effort: a failure here means the + * main workflow gets archived at run-finish, which the user can recover + * from the archive view. + */ +async function promoteMainWorkflow( + context: InstanceAiContext | undefined, + logger: Logger, + workflowId: string | undefined, +): Promise { + if (!workflowId || !context) return; + try { + await context.workflowService.clearAiTemporary(workflowId); + } catch (error) { + logger.warn( + `Failed to clear AI-builder temporary marker on main workflow ${workflowId}: ${ + error instanceof Error ? error.message : String(error) + }`, + ); + } +} + +function isRecord(value: unknown): value is Record { + return typeof value === 'object' && value !== null; +} + +type ExecutableTool = Record & { + execute: (...args: unknown[]) => unknown; +}; + +function isExecutableTool(tool: unknown): tool is ExecutableTool { + return isRecord(tool) && typeof tool.execute === 'function'; +} + +export function recordSuccessfulWorkflowBuilds( + tool: unknown, + onWorkflowId: (workflowId: string) => void, +): void { + if (!isExecutableTool(tool)) return; + + const execute = tool.execute.bind(tool); + tool.execute = async (...args: unknown[]) => { + const result = await execute(...args); + if (isRecord(result) && result.success === true && typeof result.workflowId === 'string') { + onWorkflowId(result.workflowId); + } + return result; + }; +} + const UNTESTABLE_TRIGGER_LABELS = [...UNTESTABLE_TRIGGERS].map(triggerLabel).join(', '); function detectTriggerType(attempt: SubmitWorkflowAttempt | undefined): TriggerType { @@ -508,6 +560,11 @@ export async function startBuildWorkflowAgentTask( const refreshedAttempt = submitAttempts.get(mainWorkflowPath); if (refreshedAttempt?.success) { + await promoteMainWorkflow( + domainContext, + context.logger, + refreshedAttempt.workflowId, + ); return { text: finalText, outcome: buildOutcome(workItemId, taskId, refreshedAttempt, finalText), @@ -527,12 +584,22 @@ export async function startBuildWorkflowAgentTask( } } + await promoteMainWorkflow( + domainContext, + context.logger, + mainWorkflowAttempt.workflowId, + ); return { text: finalText, outcome: buildOutcome(workItemId, taskId, mainWorkflowAttempt, finalText), }; } + let fallbackMainWorkflowId: string | undefined; + recordSuccessfulWorkflowBuilds(builderTools['build-workflow'], (workflowId) => { + fallbackMainWorkflowId = workflowId; + }); + const tracedBuilderTools = traceSubAgentTools(context, builderTools, 'workflow-builder'); const subAgent = new Agent({ @@ -592,6 +659,7 @@ export async function startBuildWorkflowAgentTask( }); const toolFinalText = await hitlResult.text; + await promoteMainWorkflow(domainContext, context.logger, fallbackMainWorkflowId); return { text: toolFinalText }; } finally { await builderWs?.cleanup(); diff --git a/packages/@n8n/instance-ai/src/tools/workflows.tool.ts b/packages/@n8n/instance-ai/src/tools/workflows.tool.ts index b8b52ca1769..b8ed3a03cc1 100644 --- a/packages/@n8n/instance-ai/src/tools/workflows.tool.ts +++ b/packages/@n8n/instance-ai/src/tools/workflows.tool.ts @@ -36,7 +36,9 @@ const getAsCodeAction = z.object({ }); const deleteAction = z.object({ - action: z.literal('delete').describe('Archive a workflow by ID (soft delete)'), + action: z + .literal('delete') + .describe('Archive a workflow by ID (soft delete — recoverable by the user)'), workflowId: z.string().describe('ID of the workflow'), }); @@ -230,7 +232,6 @@ async function handleDelete( return { success: false, denied: true, reason: 'User denied the action' }; } - // Approved or always_allow — execute await context.workflowService.archive(input.workflowId); return { success: true }; } diff --git a/packages/@n8n/instance-ai/src/tools/workflows/__tests__/setup-workflow.service.test.ts b/packages/@n8n/instance-ai/src/tools/workflows/__tests__/setup-workflow.service.test.ts index 6d45072c127..954bcb0acdc 100644 --- a/packages/@n8n/instance-ai/src/tools/workflows/__tests__/setup-workflow.service.test.ts +++ b/packages/@n8n/instance-ai/src/tools/workflows/__tests__/setup-workflow.service.test.ts @@ -27,6 +27,8 @@ function createMockContext(overrides?: Partial): InstanceAiCo delete: jest.fn(), publish: jest.fn(), unpublish: jest.fn(), + clearAiTemporary: jest.fn(), + archiveIfAiTemporary: jest.fn(), }, executionService: { list: jest.fn(), diff --git a/packages/@n8n/instance-ai/src/tools/workflows/build-workflow.tool.ts b/packages/@n8n/instance-ai/src/tools/workflows/build-workflow.tool.ts index 79fc5d042d0..0850178799b 100644 --- a/packages/@n8n/instance-ai/src/tools/workflows/build-workflow.tool.ts +++ b/packages/@n8n/instance-ai/src/tools/workflows/build-workflow.tool.ts @@ -176,12 +176,11 @@ export function createBuildWorkflowTool(context: InstanceAiContext) { await ensureWebhookIds(json, workflowId, context); try { - const opts = projectId ? { projectId } : undefined; if (workflowId) { const updated = await context.workflowService.updateFromWorkflowJSON( workflowId, json, - opts, + projectId ? { projectId } : undefined, ); return { success: true, @@ -192,7 +191,11 @@ export function createBuildWorkflowTool(context: InstanceAiContext) { : undefined, }; } else { - const created = await context.workflowService.createFromWorkflowJSON(json, opts); + const created = await context.workflowService.createFromWorkflowJSON(json, { + ...(projectId ? { projectId } : {}), + markAsAiTemporary: true, + }); + (context.aiCreatedWorkflowIds ??= new Set()).add(created.id); return { success: true, workflowId: created.id, diff --git a/packages/@n8n/instance-ai/src/tools/workflows/submit-workflow.tool.ts b/packages/@n8n/instance-ai/src/tools/workflows/submit-workflow.tool.ts index a9ed7ad86ae..898fa5d8132 100644 --- a/packages/@n8n/instance-ai/src/tools/workflows/submit-workflow.tool.ts +++ b/packages/@n8n/instance-ai/src/tools/workflows/submit-workflow.tool.ts @@ -341,18 +341,21 @@ export function createSubmitWorkflowTool( // Save let savedId: string; - const opts = projectId ? { projectId } : undefined; try { if (workflowId) { const updated = await context.workflowService.updateFromWorkflowJSON( workflowId, json, - opts, + projectId ? { projectId } : undefined, ); savedId = updated.id; } else { - const created = await context.workflowService.createFromWorkflowJSON(json, opts); + const created = await context.workflowService.createFromWorkflowJSON(json, { + ...(projectId ? { projectId } : {}), + markAsAiTemporary: true, + }); savedId = created.id; + (context.aiCreatedWorkflowIds ??= new Set()).add(created.id); } } catch (error) { const errors = [ diff --git a/packages/@n8n/instance-ai/src/types.ts b/packages/@n8n/instance-ai/src/types.ts index 6e95c46682c..89809cb988a 100644 --- a/packages/@n8n/instance-ai/src/types.ts +++ b/packages/@n8n/instance-ai/src/types.ts @@ -153,7 +153,7 @@ export interface InstanceAiWorkflowService { /** Create a workflow from SDK-produced WorkflowJSON (full NodeJSON with typeVersion, credentials, etc.). */ createFromWorkflowJSON( json: WorkflowJSON, - options?: { projectId?: string }, + options?: { projectId?: string; markAsAiTemporary?: boolean }, ): Promise; /** Update a workflow from SDK-produced WorkflowJSON. */ updateFromWorkflowJSON( @@ -163,6 +163,18 @@ export interface InstanceAiWorkflowService { ): Promise; archive(workflowId: string): Promise; delete(workflowId: string): Promise; + /** + * Clear the AI-builder temporary marker on a workflow — used to promote the + * main deliverable so the run-finish reap leaves it alone. + */ + clearAiTemporary(workflowId: string): Promise; + /** + * Archive the workflow only if it still carries the AI-builder temporary + * marker. Check-and-archive used by the run-finish reap so a main + * workflow whose marker was cleared survives even if its id is still in + * the orchestrator's in-memory created-set. + */ + archiveIfAiTemporary(workflowId: string): Promise; publish( workflowId: string, options?: { versionId?: string; name?: string; description?: string }, @@ -535,6 +547,17 @@ export interface InstanceAiContext { domainAccessTracker?: DomainAccessTracker; /** Current run ID — used for transient (allow_once) domain approvals. */ runId?: string; + /** + * IDs of workflows the agent created during the **currently active plan + * cycle**. Populated by build-workflow and submit-workflow on every + * successful create, and hydrated at run start from the persisted plan + * graph when — and only when — the plan is still `active` or + * `awaiting_replan`, so replan follow-up runs keep the bypass active but + * the window closes as soon as the plan settles. Consumed by the delete + * handler to skip the confirmation gate when the agent cleans up its own + * in-flight artifacts. Lazily initialized on first create. + */ + aiCreatedWorkflowIds?: Set; /** * Attachments from the current user message. Runtime-only — not persisted. * Used to register `parse-file` and supply data to the parser. diff --git a/packages/cli/src/modules/instance-ai/__tests__/instance-ai.adapter.service.security.test.ts b/packages/cli/src/modules/instance-ai/__tests__/instance-ai.adapter.service.security.test.ts index 3cde3792681..d4869282a65 100644 --- a/packages/cli/src/modules/instance-ai/__tests__/instance-ai.adapter.service.security.test.ts +++ b/packages/cli/src/modules/instance-ai/__tests__/instance-ai.adapter.service.security.test.ts @@ -11,6 +11,7 @@ jest.mock('@n8n/instance-ai', () => ({ import { mock } from 'jest-mock-extended'; import type { + AiBuilderTemporaryWorkflowRepository, User, ExecutionRepository, ProjectRepository, @@ -90,6 +91,7 @@ const executionPersistence = mock(); const eventService = mock(); const roleService = mock(); const telemetry = mock(); +const aiBuilderTemporaryWorkflowRepository = mock(); const service = new InstanceAiAdapterService( logger, @@ -121,6 +123,7 @@ const service = new InstanceAiAdapterService( eventService, roleService, telemetry, + aiBuilderTemporaryWorkflowRepository, ); const user = mock({ diff --git a/packages/cli/src/modules/instance-ai/__tests__/instance-ai.adapter.service.test.ts b/packages/cli/src/modules/instance-ai/__tests__/instance-ai.adapter.service.test.ts index 7fecc664381..5dfdee59713 100644 --- a/packages/cli/src/modules/instance-ai/__tests__/instance-ai.adapter.service.test.ts +++ b/packages/cli/src/modules/instance-ai/__tests__/instance-ai.adapter.service.test.ts @@ -684,6 +684,7 @@ jest.mock('@/permissions.ee/check-access', () => ({ })); import type { + AiBuilderTemporaryWorkflowRepository, User, ExecutionRepository, ProjectRepository, @@ -746,6 +747,7 @@ function createNodeAdapterForTests(nodes: Array>) { {} as unknown as ConstructorParameters[26], {} as unknown as ConstructorParameters[27], {} as unknown as ConstructorParameters[28], + {} as unknown as ConstructorParameters[29], ); ( @@ -874,6 +876,7 @@ function createDataTableAdapterForTests(overrides?: { {} as unknown as ConstructorParameters[26], {} as unknown as ConstructorParameters[27], {} as unknown as ConstructorParameters[28], + {} as unknown as ConstructorParameters[29], ); const adapter = service.createContext(mockUser).dataTableService; @@ -1063,14 +1066,38 @@ function createWorkflowAdapterForTests(overrides?: { const mockWorkflowRepository = { create: jest.fn().mockImplementation((data: Record) => data), save: jest.fn().mockResolvedValue(savedWorkflow), + update: jest.fn().mockResolvedValue(undefined), + manager: { + transaction: jest.fn( + async ( + fn: (transactionManager: { save: jest.Mock }) => Promise, + ): Promise => { + return await fn({ + save: jest.fn().mockResolvedValue(savedWorkflow), + }); + }, + ), + }, + }; + + const mockWorkflowFinderService = { + findWorkflowForUser: jest.fn().mockResolvedValue(savedWorkflow), }; const mockSharedWorkflowRepository = { create: jest.fn().mockImplementation((data: Record) => data), save: jest.fn().mockResolvedValue(undefined), + makeOwner: jest.fn().mockResolvedValue(undefined), + }; + + const mockAiBuilderTemporaryWorkflowRepository = { + mark: jest.fn().mockResolvedValue(undefined), + unmark: jest.fn().mockResolvedValue(undefined), + existsForWorkflow: jest.fn().mockResolvedValue(false), }; const mockWorkflowService = { + archive: jest.fn().mockResolvedValue(undefined), update: jest.fn().mockResolvedValue(savedWorkflow), }; @@ -1084,7 +1111,9 @@ function createWorkflowAdapterForTests(overrides?: { typeof InstanceAiAdapterService >[1], mockWorkflowService as unknown as WorkflowService, - {} as unknown as ConstructorParameters[3], + mockWorkflowFinderService as unknown as ConstructorParameters< + typeof InstanceAiAdapterService + >[3], mockWorkflowRepository as unknown as WorkflowRepository, mockSharedWorkflowRepository as unknown as SharedWorkflowRepository, mockProjectRepository as unknown as ProjectRepository, @@ -1122,10 +1151,11 @@ function createWorkflowAdapterForTests(overrides?: { {} as unknown as ConstructorParameters[25], {} as unknown as ConstructorParameters[26], {} as unknown as ConstructorParameters[27], - {} as unknown as ConstructorParameters[28], + { track: jest.fn() } as unknown as ConstructorParameters[28], + mockAiBuilderTemporaryWorkflowRepository as unknown as AiBuilderTemporaryWorkflowRepository, ); - const context = service.createContext(mockUser); + const context = service.createContext(mockUser, { threadId: 'thread-1' }); const adapter = context.workflowService; return { @@ -1133,7 +1163,9 @@ function createWorkflowAdapterForTests(overrides?: { context, mockProjectRepository, mockWorkflowRepository, + mockWorkflowFinderService, mockSharedWorkflowRepository, + mockAiBuilderTemporaryWorkflowRepository, mockWorkflowService, mockUser, }; @@ -1158,8 +1190,10 @@ describe('createWorkflowAdapter', () => { await adapter.createFromWorkflowJSON(minimalWorkflowJSON); expect(mockProjectRepository.getPersonalProjectForUserOrFail).toHaveBeenCalledWith('user-1'); - expect(mockSharedWorkflowRepository.create).toHaveBeenCalledWith( - expect.objectContaining({ projectId: 'personal-project-id' }), + expect(mockSharedWorkflowRepository.makeOwner).toHaveBeenCalledWith( + ['wf-new'], + 'personal-project-id', + expect.any(Object), ); }); @@ -1172,8 +1206,10 @@ describe('createWorkflowAdapter', () => { }); expect(mockProjectRepository.getPersonalProjectForUserOrFail).not.toHaveBeenCalled(); - expect(mockSharedWorkflowRepository.create).toHaveBeenCalledWith( - expect.objectContaining({ projectId: 'team-project-id' }), + expect(mockSharedWorkflowRepository.makeOwner).toHaveBeenCalledWith( + ['wf-new'], + 'team-project-id', + expect.any(Object), ); }); @@ -1188,6 +1224,89 @@ describe('createWorkflowAdapter', () => { ).rejects.toThrow('User does not have the required permissions in this project'); }); + it('marks the workflow as AI-builder temporary when markAsAiTemporary is true', async () => { + const { + adapter, + mockWorkflowRepository, + mockSharedWorkflowRepository, + mockAiBuilderTemporaryWorkflowRepository, + } = createWorkflowAdapterForTests(); + + await adapter.createFromWorkflowJSON(minimalWorkflowJSON, { + markAsAiTemporary: true, + }); + + expect(mockWorkflowRepository.create).toHaveBeenCalledWith( + expect.not.objectContaining({ meta: expect.anything() }), + ); + expect(mockWorkflowRepository.manager.transaction).toHaveBeenCalled(); + expect(mockSharedWorkflowRepository.makeOwner).toHaveBeenCalledWith( + ['wf-new'], + 'personal-project-id', + expect.any(Object), + ); + expect(mockAiBuilderTemporaryWorkflowRepository.mark).toHaveBeenCalledWith( + 'wf-new', + 'thread-1', + expect.any(Object), + ); + }); + + it('does not mark the workflow as AI-builder temporary when markAsAiTemporary is omitted', async () => { + const { adapter, mockWorkflowRepository } = createWorkflowAdapterForTests(); + + await adapter.createFromWorkflowJSON(minimalWorkflowJSON); + + expect(mockWorkflowRepository.create).toHaveBeenCalledWith( + expect.not.objectContaining({ meta: expect.anything() }), + ); + }); + + it('clears the AI-builder temporary marker when promoting the main workflow', async () => { + const { adapter, mockAiBuilderTemporaryWorkflowRepository, mockWorkflowRepository } = + createWorkflowAdapterForTests(); + mockAiBuilderTemporaryWorkflowRepository.existsForWorkflow.mockResolvedValue(true); + + await adapter.clearAiTemporary('wf-new'); + + expect(mockAiBuilderTemporaryWorkflowRepository.unmark).toHaveBeenCalledWith('wf-new'); + expect(mockWorkflowRepository.update).not.toHaveBeenCalled(); + }); + + it('archives and unmarks an unpromoted AI-builder temporary workflow', async () => { + const { adapter, mockAiBuilderTemporaryWorkflowRepository, mockWorkflowService } = + createWorkflowAdapterForTests(); + mockAiBuilderTemporaryWorkflowRepository.existsForWorkflow.mockResolvedValue(true); + + await expect(adapter.archiveIfAiTemporary('wf-new')).resolves.toBe(true); + + expect(mockWorkflowService.archive).toHaveBeenCalledWith( + expect.objectContaining({ id: 'user-1' }), + 'wf-new', + { skipArchived: true }, + ); + expect(mockAiBuilderTemporaryWorkflowRepository.unmark).toHaveBeenCalledWith('wf-new'); + }); + + it('unmarks already-archived temporary workflows without archiving again', async () => { + const { + adapter, + mockAiBuilderTemporaryWorkflowRepository, + mockWorkflowFinderService, + mockWorkflowService, + } = createWorkflowAdapterForTests(); + mockAiBuilderTemporaryWorkflowRepository.existsForWorkflow.mockResolvedValue(true); + mockWorkflowFinderService.findWorkflowForUser.mockResolvedValue({ + id: 'wf-archived', + isArchived: true, + }); + + await expect(adapter.archiveIfAiTemporary('wf-archived')).resolves.toBe(false); + + expect(mockWorkflowService.archive).not.toHaveBeenCalled(); + expect(mockAiBuilderTemporaryWorkflowRepository.unmark).toHaveBeenCalledWith('wf-archived'); + }); + describe('instance read-only mode', () => { it('blocks createFromWorkflowJSON when branchReadOnly is true', async () => { const { adapter } = createWorkflowAdapterForTests({ branchReadOnly: true }); @@ -1360,6 +1479,7 @@ function createExecutionAdapterForTests(overrides?: { sharingEnabled?: boolean } {} as unknown as ConstructorParameters[26], mockRoleService as unknown as RoleService, {} as unknown as ConstructorParameters[28], + {} as unknown as ConstructorParameters[29], ); const adapter = service.createContext(mockUser).executionService; diff --git a/packages/cli/src/modules/instance-ai/instance-ai.adapter.service.ts b/packages/cli/src/modules/instance-ai/instance-ai.adapter.service.ts index fb116446d76..36f7f847a41 100644 --- a/packages/cli/src/modules/instance-ai/instance-ai.adapter.service.ts +++ b/packages/cli/src/modules/instance-ai/instance-ai.adapter.service.ts @@ -39,7 +39,7 @@ import { wrapUntrustedData } from '@n8n/instance-ai'; import type { WorkflowJSON } from '@n8n/workflow-sdk'; import { GlobalConfig } from '@n8n/config'; import { Time } from '@n8n/constants'; -import type { User, WorkflowEntity, ExecutionSummaries } from '@n8n/db'; +import type { User, ExecutionSummaries } from '@n8n/db'; import { InstanceAiSettingsService } from './instance-ai-settings.service'; import { @@ -55,9 +55,11 @@ import { LRUCache, } from './web-research'; import { + AiBuilderTemporaryWorkflowRepository, ExecutionRepository, ProjectRepository, SharedWorkflowRepository, + WorkflowEntity, WorkflowRepository, } from '@n8n/db'; import { Logger } from '@n8n/backend-common'; @@ -90,6 +92,7 @@ import { WEBHOOK_NODE_TYPE, SCHEDULE_TRIGGER_NODE_TYPE, TimeoutExecutionCancelledError, + UnexpectedError, jsonParse, } from 'n8n-workflow'; @@ -183,6 +186,7 @@ export class InstanceAiAdapterService { private readonly eventService: EventService, private readonly roleService: RoleService, private readonly telemetry: Telemetry, + private readonly aiBuilderTemporaryWorkflowRepository: AiBuilderTemporaryWorkflowRepository, ) { this.logger = logger.scoped('instance-ai'); this.allowSendingParameterValues = globalConfig.ai.allowSendingParameterValues; @@ -267,6 +271,7 @@ export class InstanceAiAdapterService { workflowFinderService, workflowRepository, sharedWorkflowRepository, + aiBuilderTemporaryWorkflowRepository, workflowHistoryService, enterpriseWorkflowService, license, @@ -323,6 +328,36 @@ export class InstanceAiAdapterService { await workflowService.delete(user, workflowId); }, + async clearAiTemporary(workflowId: string) { + assertNotReadOnly(); + const workflow = await workflowFinderService.findWorkflowForUser(workflowId, user, [ + 'workflow:update', + ]); + if (!workflow) return; + if (!(await aiBuilderTemporaryWorkflowRepository.existsForWorkflow(workflowId))) return; + + await aiBuilderTemporaryWorkflowRepository.unmark(workflowId); + }, + + async archiveIfAiTemporary(workflowId: string) { + assertNotReadOnly(); + const workflow = await workflowFinderService.findWorkflowForUser(workflowId, user, [ + 'workflow:update', + ]); + if (!workflow) return false; + if (!(await aiBuilderTemporaryWorkflowRepository.existsForWorkflow(workflowId))) { + return false; + } + if (workflow.isArchived) { + await aiBuilderTemporaryWorkflowRepository.unmark(workflowId); + return false; + } + + await workflowService.archive(user, workflowId, { skipArchived: true }); + await aiBuilderTemporaryWorkflowRepository.unmark(workflowId); + return true; + }, + async publish( workflowId: string, options?: { versionId?: string; name?: string; description?: string }, @@ -361,7 +396,10 @@ export class InstanceAiAdapterService { return toWorkflowJSON(wf, { redactParameters }); }, - async createFromWorkflowJSON(json: WorkflowJSON, options?: { projectId?: string }) { + async createFromWorkflowJSON( + json: WorkflowJSON, + options?: { projectId?: string; markAsAiTemporary?: boolean }, + ) { assertNotReadOnly(); const projectId = await resolveProjectId(['workflow:create'], options?.projectId); @@ -393,15 +431,23 @@ export class InstanceAiAdapterService { versionId: randomUUID(), } as Partial); - const saved = await workflowRepository.save(newWorkflow); - - await sharedWorkflowRepository.save( - sharedWorkflowRepository.create({ - role: 'workflow:owner', - projectId, - workflow: saved, - }), - ); + const saved = await workflowRepository.manager.transaction(async (transactionManager) => { + const workflow = await transactionManager.save(WorkflowEntity, newWorkflow); + await sharedWorkflowRepository.makeOwner([workflow.id], projectId, transactionManager); + if (options?.markAsAiTemporary) { + if (!threadId) { + throw new UnexpectedError( + 'Cannot mark AI-builder temporary workflow without a thread ID', + ); + } + await aiBuilderTemporaryWorkflowRepository.mark( + workflow.id, + threadId, + transactionManager, + ); + } + return workflow; + }); // Now update with actual nodes — this creates the WorkflowHistory entry // needed for activation and publishing. diff --git a/packages/cli/src/modules/instance-ai/instance-ai.service.ts b/packages/cli/src/modules/instance-ai/instance-ai.service.ts index cc11bd1409a..328e94d7b47 100644 --- a/packages/cli/src/modules/instance-ai/instance-ai.service.ts +++ b/packages/cli/src/modules/instance-ai/instance-ai.service.ts @@ -14,7 +14,7 @@ import { Logger } from '@n8n/backend-common'; import { GlobalConfig } from '@n8n/config'; import { Time } from '@n8n/constants'; import type { InstanceAiConfig } from '@n8n/config'; -import type { User } from '@n8n/db'; +import { AiBuilderTemporaryWorkflowRepository, UserRepository, type User } from '@n8n/db'; import { Service } from '@n8n/di'; import { UrlService } from '@/services/url.service'; import { @@ -64,6 +64,7 @@ import { type SpawnBackgroundTaskOptions, type ServiceProxyConfig, type StreamableAgent, + type SuspendedRunState, WorkflowTaskCoordinator, WorkflowLoopStorage, } from '@n8n/instance-ai'; @@ -216,6 +217,8 @@ export class InstanceAiService { private readonly dbIterationLogStorage: DbIterationLogStorage, private readonly sourceControlPreferencesService: SourceControlPreferencesService, private readonly telemetry: Telemetry, + private readonly userRepository: UserRepository, + private readonly aiBuilderTemporaryWorkflowRepository: AiBuilderTemporaryWorkflowRepository, ) { this.logger = logger.scoped('instance-ai'); this.instanceAiConfig = globalConfig.instanceAi; @@ -520,7 +523,17 @@ export class InstanceAiService { // Fast in-memory check — prevents the read-then-write race within a single process. if (this.creditedThreads.has(threadId)) return; - const thread = await this.threadRepo.findOneBy({ id: threadId }); + let thread: Awaited>; + try { + thread = await this.threadRepo.findOneBy({ id: threadId }); + } catch (error) { + this.logger.warn('Failed to check Instance AI credit status', { + threadId, + runId, + error: getErrorMessage(error), + }); + return; + } if (!thread) return; if (thread.metadata?.creditCounted) { this.creditedThreads.add(threadId); // Sync in-memory with DB state @@ -931,27 +944,7 @@ export class InstanceAiService { if (suspended) { suspended.abortController.abort(); - void this.finalizeRunTracing(suspended.runId, suspended.tracing, { - status: 'cancelled', - reason: 'user_cancelled', - }); - this.eventBus.publish(threadId, { - type: 'run-finish', - runId: suspended.runId, - agentId: ORCHESTRATOR_AGENT_ID, - payload: { status: 'cancelled', reason: 'user_cancelled' }, - }); - // Persist the snapshot so the run-finish event (which clears - // in-flight tool calls) is reflected in the stored tree. - void this.saveAgentTreeSnapshot(threadId, suspended.runId, this.dbSnapshotStorage, true); - if (suspended.mastraRunId) { - void this.cleanupMastraSnapshots(suspended.mastraRunId); - } - void this.maybeFinalizeRunTraceRoot(suspended.runId, { - status: 'cancelled', - reason: 'user_cancelled', - metadata: { completion_source: 'orchestrator' }, - }); + void this.finalizeCancelledSuspendedRun(suspended); } } @@ -1141,6 +1134,7 @@ export class InstanceAiService { this.threadPushRef.delete(threadId); this.deleteTraceContextsForThread(threadId); await this.destroySandbox(threadId); + await this.reapAiTemporaryForThreadCleanup(threadId); this.eventBus.clearThread(threadId); } @@ -1696,6 +1690,7 @@ export class InstanceAiService { let mastraRunId = ''; let tracing: InstanceAiTraceContext | undefined; let messageTraceFinalization: MessageTraceFinalization | undefined; + let aiCreatedWorkflowIds: Set | undefined; try { const messageId = nanoid(); @@ -1733,6 +1728,7 @@ export class InstanceAiService { messageGroupId, executionPushRef, ); + aiCreatedWorkflowIds = context.aiCreatedWorkflowIds ??= new Set(); // Make the current user message available to sub-agents (e.g. planner) // since memory.recall() only returns previously-saved messages. orchestrationContext.currentUserMessage = message; @@ -2060,9 +2056,15 @@ export class InstanceAiService { modelId, metadata: { completion_source: 'orchestrator' }, }; + const archivedWorkflowIds = await this.reapAiTemporaryFromRun( + threadId, + user, + aiCreatedWorkflowIds, + ); await this.finalizeRun(threadId, runId, result.status, snapshotStorage, { userId: user.id, modelId, + archivedWorkflowIds, }); // Count credits on first completed run per thread @@ -2087,7 +2089,12 @@ export class InstanceAiService { reason: 'user_cancelled', metadata: { completion_source: 'orchestrator' }, }; - this.publishRunFinish(threadId, runId, 'cancelled', 'user_cancelled'); + const archivedWorkflowIds = await this.reapAiTemporaryFromRun( + threadId, + user, + aiCreatedWorkflowIds, + ); + this.publishRunFinish(threadId, runId, 'cancelled', 'user_cancelled', archivedWorkflowIds); return; } @@ -2108,6 +2115,11 @@ export class InstanceAiService { metadata: { completion_source: 'orchestrator' }, }; + const archivedWorkflowIds = await this.reapAiTemporaryFromRun( + threadId, + user, + aiCreatedWorkflowIds, + ); this.eventBus.publish(threadId, { type: 'run-finish', runId, @@ -2115,6 +2127,7 @@ export class InstanceAiService { payload: { status: 'error', reason: errorMessage, + ...(archivedWorkflowIds.length > 0 ? { archivedWorkflowIds } : {}), }, }); } finally { @@ -2312,7 +2325,14 @@ export class InstanceAiService { outputText, metadata: { completion_source: 'orchestrator' }, }; - await this.finalizeRun(opts.threadId, opts.runId, result.status, opts.snapshotStorage); + const archivedWorkflowIds = await this.reapAiTemporaryFromRun( + opts.threadId, + opts.user, + undefined, + ); + await this.finalizeRun(opts.threadId, opts.runId, result.status, opts.snapshotStorage, { + archivedWorkflowIds, + }); if (result.status === 'completed') { this.telemetry.track('Builder sent message', { @@ -2334,7 +2354,18 @@ export class InstanceAiService { reason: 'user_cancelled', metadata: { completion_source: 'orchestrator' }, }; - this.publishRunFinish(opts.threadId, opts.runId, 'cancelled', 'user_cancelled'); + const archivedWorkflowIds = await this.reapAiTemporaryFromRun( + opts.threadId, + opts.user, + undefined, + ); + this.publishRunFinish( + opts.threadId, + opts.runId, + 'cancelled', + 'user_cancelled', + archivedWorkflowIds, + ); return; } @@ -2355,6 +2386,11 @@ export class InstanceAiService { metadata: { completion_source: 'orchestrator' }, }; + const archivedWorkflowIds = await this.reapAiTemporaryFromRun( + opts.threadId, + opts.user, + undefined, + ); this.eventBus.publish(opts.threadId, { type: 'run-finish', runId: opts.runId, @@ -2362,6 +2398,7 @@ export class InstanceAiService { payload: { status: 'error', reason: errorMessage, + ...(archivedWorkflowIds.length > 0 ? { archivedWorkflowIds } : {}), }, }); } finally { @@ -2537,21 +2574,178 @@ export class InstanceAiService { }); } + /** + * Archive any workflow the agent created for this thread that still carries + * the AI-builder temporary marker. The orchestrator clears the marker on the + * main deliverable before run-finish, so anything still marked is a + * stepping-stone — chunk, scratch, or sub-workflow the user never sees in + * the workflows list. Soft delete: a mistaken reap is recoverable from the + * archive view. + * + * Best-effort. Individual archive failures are logged but do not block + * the run-finish emit. + */ + private async reapAiTemporaryFromRun( + threadId: string, + user: User, + createdWorkflowIds: Set | undefined, + ): Promise { + let markedWorkflows: Array<{ workflowId: string }> = []; + try { + markedWorkflows = await this.aiBuilderTemporaryWorkflowRepository.findByThread(threadId); + } catch (error) { + this.logger.warn('Failed to inspect AI-builder temporary workflows during run finish', { + threadId, + error: getErrorMessage(error), + }); + } + const workflowIds = new Set([ + ...markedWorkflows.map(({ workflowId }) => workflowId), + ...(createdWorkflowIds ?? []), + ]); + if (workflowIds.size === 0) return []; + + return await this.archiveAiTemporaryWorkflows(threadId, user, workflowIds); + } + + private async archiveAiTemporaryWorkflows( + threadId: string, + user: User, + workflowIds: Set, + ): Promise { + const adapter = this.adapterService.createContext(user, { threadId }); + const archived: string[] = []; + for (const workflowId of workflowIds) { + try { + const didArchive = await adapter.workflowService.archiveIfAiTemporary(workflowId); + if (didArchive) archived.push(workflowId); + } catch (error) { + this.logger.warn('Failed to reap AI-builder temporary workflow', { + threadId, + workflowId, + error: getErrorMessage(error), + }); + } + } + return archived; + } + + private async finalizeCancelledSuspendedRun(suspended: SuspendedRunState): Promise { + await this.finalizeRunTracing(suspended.runId, suspended.tracing, { + status: 'cancelled', + reason: 'user_cancelled', + }); + + const archivedWorkflowIds = await this.reapAiTemporaryFromRun( + suspended.threadId, + suspended.user, + undefined, + ); + this.publishRunFinish( + suspended.threadId, + suspended.runId, + 'cancelled', + 'user_cancelled', + archivedWorkflowIds, + ); + + // Persist the snapshot so the run-finish event (which clears + // in-flight tool calls) is reflected in the stored tree. + await this.saveAgentTreeSnapshot( + suspended.threadId, + suspended.runId, + this.dbSnapshotStorage, + true, + ); + if (suspended.mastraRunId) { + void this.cleanupMastraSnapshots(suspended.mastraRunId); + } + await this.maybeFinalizeRunTraceRoot(suspended.runId, { + status: 'cancelled', + reason: 'user_cancelled', + metadata: { completion_source: 'orchestrator' }, + }); + } + + private async reapAiTemporaryForThreadCleanup(threadId: string): Promise { + let markedWorkflows: Array<{ workflowId: string }>; + try { + markedWorkflows = await this.aiBuilderTemporaryWorkflowRepository.findByThread(threadId); + } catch (error) { + this.logger.warn('Failed to inspect AI-builder temporary workflows during thread cleanup', { + threadId, + error: getErrorMessage(error), + }); + return; + } + + if (markedWorkflows.length === 0) return; + + let thread: Awaited>; + try { + thread = await this.threadRepo.findOneBy({ id: threadId }); + } catch (error) { + this.logger.warn('Failed to load thread owner for AI-builder temporary workflow cleanup', { + threadId, + markedWorkflowCount: markedWorkflows.length, + error: getErrorMessage(error), + }); + return; + } + if (!thread?.resourceId) { + this.logger.warn('Skipping AI-builder temporary workflow cleanup for thread without owner', { + threadId, + markedWorkflowCount: markedWorkflows.length, + }); + return; + } + + let user: User | null; + try { + user = await this.userRepository.findOneBy({ id: thread.resourceId }); + } catch (error) { + this.logger.warn('Failed to load user for AI-builder temporary workflow cleanup', { + threadId, + userId: thread.resourceId, + markedWorkflowCount: markedWorkflows.length, + error: getErrorMessage(error), + }); + return; + } + if (!user) { + this.logger.warn('Skipping AI-builder temporary workflow cleanup for missing thread owner', { + threadId, + userId: thread.resourceId, + markedWorkflowCount: markedWorkflows.length, + }); + return; + } + + await this.archiveAiTemporaryWorkflows( + threadId, + user, + new Set(markedWorkflows.map(({ workflowId }) => workflowId)), + ); + } + private publishRunFinish( threadId: string, runId: string, status: 'completed' | 'cancelled' | 'errored', reason?: string, + archivedWorkflowIds?: string[], ): void { const effectiveStatus = status === 'errored' ? 'error' : status; + const hasArchived = archivedWorkflowIds && archivedWorkflowIds.length > 0; this.eventBus.publish(threadId, { type: 'run-finish', runId, agentId: ORCHESTRATOR_AGENT_ID, - payload: - status === 'cancelled' - ? { status: effectiveStatus, reason: reason ?? 'user_cancelled' } - : { status: effectiveStatus }, + payload: { + status: effectiveStatus, + ...(status === 'cancelled' ? { reason: reason ?? 'user_cancelled' } : {}), + ...(hasArchived ? { archivedWorkflowIds } : {}), + }, }); } @@ -2560,9 +2754,9 @@ export class InstanceAiService { runId: string, status: 'completed' | 'cancelled' | 'errored', snapshotStorage: DbSnapshotStorage, - options?: { userId?: string; modelId?: ModelConfig }, + options?: { userId?: string; modelId?: ModelConfig; archivedWorkflowIds?: string[] }, ): Promise { - this.publishRunFinish(threadId, runId, status); + this.publishRunFinish(threadId, runId, status, undefined, options?.archivedWorkflowIds); if (status === 'completed') { await this.saveAgentTreeSnapshot(threadId, runId, snapshotStorage); if (options?.userId && options?.modelId) { diff --git a/packages/frontend/@n8n/i18n/src/locales/en.json b/packages/frontend/@n8n/i18n/src/locales/en.json index d0fa4604dba..0b9630dc1de 100644 --- a/packages/frontend/@n8n/i18n/src/locales/en.json +++ b/packages/frontend/@n8n/i18n/src/locales/en.json @@ -5143,6 +5143,7 @@ "instanceAi.artifactsPanel.noArtifacts": "No artifacts yet", "instanceAi.artifactsPanel.tasks": "Tasks", "instanceAi.artifactsPanel.openWorkflow": "Open", + "instanceAi.artifactsPanel.archived": "Archived", "instanceAi.previewTabBar.collapse": "Collapse panel", "instanceAi.previewTabBar.openInEditor": "Open in editor", "instanceAi.previewTabBar.copyLink": "Copy link", diff --git a/packages/frontend/editor-ui/src/features/ai/instanceAi/components/AgentActivityTree.vue b/packages/frontend/editor-ui/src/features/ai/instanceAi/components/AgentActivityTree.vue index 382b5d12530..646cd00e72b 100644 --- a/packages/frontend/editor-ui/src/features/ai/instanceAi/components/AgentActivityTree.vue +++ b/packages/frontend/editor-ui/src/features/ai/instanceAi/components/AgentActivityTree.vue @@ -91,6 +91,7 @@ function resolveArtifactName(artifact: ArtifactInfo): string { :name="resolveArtifactName(artifact)" :resource-id="artifact.resourceId" :project-id="artifact.projectId" + :archived="store.producedArtifacts.get(artifact.resourceId)?.archived" :class="$style.artifactCard" /> diff --git a/packages/frontend/editor-ui/src/features/ai/instanceAi/components/AgentTimeline.vue b/packages/frontend/editor-ui/src/features/ai/instanceAi/components/AgentTimeline.vue index 955a7cf42bd..4e99f5ba96c 100644 --- a/packages/frontend/editor-ui/src/features/ai/instanceAi/components/AgentTimeline.vue +++ b/packages/frontend/editor-ui/src/features/ai/instanceAi/components/AgentTimeline.vue @@ -279,6 +279,7 @@ function mapTaskItemsToPlannedTasks(tasks?: TaskList): PlannedTaskArg[] | undefi :name="resolveArtifactName(artifact)" :resource-id="artifact.resourceId" :project-id="artifact.projectId" + :archived="store.producedArtifacts.get(artifact.resourceId)?.archived" :metadata="formatArtifactMetadata(artifact)" /> diff --git a/packages/frontend/editor-ui/src/features/ai/instanceAi/components/ArtifactCard.vue b/packages/frontend/editor-ui/src/features/ai/instanceAi/components/ArtifactCard.vue index 4c12dc5ee75..5074469858c 100644 --- a/packages/frontend/editor-ui/src/features/ai/instanceAi/components/ArtifactCard.vue +++ b/packages/frontend/editor-ui/src/features/ai/instanceAi/components/ArtifactCard.vue @@ -1,13 +1,17 @@