fix(ai-builder): Hide and reap intermediate AI-created workflows (#29066)

This commit is contained in:
Albert Alises
2026-04-27 15:17:54 +00:00
committed by GitHub
parent 90843cf4ba
commit 632ae67de3
30 changed files with 852 additions and 64 deletions
@@ -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({
@@ -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<WorkflowEntity>;
}
+3
View File
@@ -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,
@@ -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');
}
}
@@ -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,
];
@@ -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 };
@@ -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<SelectQueryBuilder<WorkflowEntity>>();
const subQueryBuilder = mock<SelectQueryBuilder<AiBuilderTemporaryWorkflow>>();
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'];
@@ -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<AiBuilderTemporaryWorkflow> {
constructor(dataSource: DataSource) {
super(AiBuilderTemporaryWorkflow, dataSource.manager);
}
async mark(
workflowId: string,
threadId: string,
entityManager: EntityManager = this.manager,
): Promise<void> {
await entityManager.upsert(AiBuilderTemporaryWorkflow, { workflowId, threadId }, [
'workflowId',
]);
}
async unmark(workflowId: string, entityManager: EntityManager = this.manager): Promise<void> {
await entityManager.delete(AiBuilderTemporaryWorkflow, { workflowId });
}
async findByThread(threadId: string): Promise<AiBuilderTemporaryWorkflow[]> {
return await this.find({
where: { threadId },
select: ['workflowId', 'threadId'],
});
}
async existsForWorkflow(workflowId: string): Promise<boolean> {
return await this.existsBy({ workflowId });
}
}
@@ -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';
@@ -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<WorkflowEntity> {
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<WorkflowEntity>): void {
const markerSubquery = qb
.subQuery()
.select('1')
.from(AiBuilderTemporaryWorkflow, 'aitw')
.where('aitw."workflowId" = workflow.id')
.getQuery();
qb.andWhere(`NOT EXISTS ${markerSubquery}`);
}
private applyAvailableInMCPFilter(
@@ -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<string, unknown> }) => {
metadata: Record<string, unknown>;
};
},
) => {
const currentMetadata = metadataByThread.get(opts.threadId) ?? {};
const next = opts.update({ metadata: currentMetadata });
metadataByThread.set(opts.threadId, next.metadata);
},
),
}));
const metadataByThread = new Map<string, Record<string, unknown>>();
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');
});
});
@@ -22,6 +22,8 @@ function createMockContext(overrides?: Partial<InstanceAiContext>): InstanceAiCo
delete: jest.fn(),
publish: jest.fn(),
unpublish: jest.fn(),
clearAiTemporary: jest.fn(),
archiveIfAiTemporary: jest.fn(),
},
executionService: {
list: jest.fn(),
@@ -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();
});
});
@@ -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<void> {
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<string, unknown> {
return typeof value === 'object' && value !== null;
}
type ExecutableTool = Record<string, unknown> & {
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();
@@ -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 };
}
@@ -27,6 +27,8 @@ function createMockContext(overrides?: Partial<InstanceAiContext>): InstanceAiCo
delete: jest.fn(),
publish: jest.fn(),
unpublish: jest.fn(),
clearAiTemporary: jest.fn(),
archiveIfAiTemporary: jest.fn(),
},
executionService: {
list: jest.fn(),
@@ -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<string>()).add(created.id);
return {
success: true,
workflowId: created.id,
@@ -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<string>()).add(created.id);
}
} catch (error) {
const errors = [
+24 -1
View File
@@ -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<WorkflowDetail>;
/** Update a workflow from SDK-produced WorkflowJSON. */
updateFromWorkflowJSON(
@@ -163,6 +163,18 @@ export interface InstanceAiWorkflowService {
): Promise<WorkflowDetail>;
archive(workflowId: string): Promise<void>;
delete(workflowId: string): Promise<void>;
/**
* 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<void>;
/**
* 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<boolean>;
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<string>;
/**
* Attachments from the current user message. Runtime-only — not persisted.
* Used to register `parse-file` and supply data to the parser.
@@ -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<ExecutionPersistence>();
const eventService = mock<EventService>();
const roleService = mock<RoleService>();
const telemetry = mock<Telemetry>();
const aiBuilderTemporaryWorkflowRepository = mock<AiBuilderTemporaryWorkflowRepository>();
const service = new InstanceAiAdapterService(
logger,
@@ -121,6 +123,7 @@ const service = new InstanceAiAdapterService(
eventService,
roleService,
telemetry,
aiBuilderTemporaryWorkflowRepository,
);
const user = mock<User>({
@@ -684,6 +684,7 @@ jest.mock('@/permissions.ee/check-access', () => ({
}));
import type {
AiBuilderTemporaryWorkflowRepository,
User,
ExecutionRepository,
ProjectRepository,
@@ -746,6 +747,7 @@ function createNodeAdapterForTests(nodes: Array<Record<string, unknown>>) {
{} as unknown as ConstructorParameters<typeof InstanceAiAdapterService>[26],
{} as unknown as ConstructorParameters<typeof InstanceAiAdapterService>[27],
{} as unknown as ConstructorParameters<typeof InstanceAiAdapterService>[28],
{} as unknown as ConstructorParameters<typeof InstanceAiAdapterService>[29],
);
(
@@ -874,6 +876,7 @@ function createDataTableAdapterForTests(overrides?: {
{} as unknown as ConstructorParameters<typeof InstanceAiAdapterService>[26],
{} as unknown as ConstructorParameters<typeof InstanceAiAdapterService>[27],
{} as unknown as ConstructorParameters<typeof InstanceAiAdapterService>[28],
{} as unknown as ConstructorParameters<typeof InstanceAiAdapterService>[29],
);
const adapter = service.createContext(mockUser).dataTableService;
@@ -1063,14 +1066,38 @@ function createWorkflowAdapterForTests(overrides?: {
const mockWorkflowRepository = {
create: jest.fn().mockImplementation((data: Record<string, unknown>) => data),
save: jest.fn().mockResolvedValue(savedWorkflow),
update: jest.fn().mockResolvedValue(undefined),
manager: {
transaction: jest.fn(
async (
fn: (transactionManager: { save: jest.Mock }) => Promise<unknown>,
): Promise<unknown> => {
return await fn({
save: jest.fn().mockResolvedValue(savedWorkflow),
});
},
),
},
};
const mockWorkflowFinderService = {
findWorkflowForUser: jest.fn().mockResolvedValue(savedWorkflow),
};
const mockSharedWorkflowRepository = {
create: jest.fn().mockImplementation((data: Record<string, unknown>) => 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<typeof InstanceAiAdapterService>[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<typeof InstanceAiAdapterService>[25],
{} as unknown as ConstructorParameters<typeof InstanceAiAdapterService>[26],
{} as unknown as ConstructorParameters<typeof InstanceAiAdapterService>[27],
{} as unknown as ConstructorParameters<typeof InstanceAiAdapterService>[28],
{ track: jest.fn() } as unknown as ConstructorParameters<typeof InstanceAiAdapterService>[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<typeof InstanceAiAdapterService>[26],
mockRoleService as unknown as RoleService,
{} as unknown as ConstructorParameters<typeof InstanceAiAdapterService>[28],
{} as unknown as ConstructorParameters<typeof InstanceAiAdapterService>[29],
);
const adapter = service.createContext(mockUser).executionService;
@@ -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<WorkflowEntity>);
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.
@@ -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<ReturnType<InstanceAiThreadRepository['findOneBy']>>;
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<string> | undefined;
try {
const messageId = nanoid();
@@ -1733,6 +1728,7 @@ export class InstanceAiService {
messageGroupId,
executionPushRef,
);
aiCreatedWorkflowIds = context.aiCreatedWorkflowIds ??= new Set<string>();
// 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<string> | undefined,
): Promise<string[]> {
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<string>,
): Promise<string[]> {
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<User>): Promise<void> {
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<void> {
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<ReturnType<InstanceAiThreadRepository['findOneBy']>>;
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<void> {
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) {
@@ -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",
@@ -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"
/>
</template>
@@ -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)"
/>
</template>
@@ -1,13 +1,17 @@
<script lang="ts" setup>
import { N8nCard, N8nIcon, N8nText, type IconName } from '@n8n/design-system';
import { useI18n } from '@n8n/i18n';
import { computed, inject } from 'vue';
const i18n = useI18n();
const props = defineProps<{
type: 'workflow' | 'data-table';
name: string;
resourceId: string;
projectId?: string;
metadata?: string;
archived?: boolean;
}>();
const openPreview = inject<((id: string) => void) | undefined>('openWorkflowPreview', undefined);
@@ -46,12 +50,15 @@ function handleClick(e: MouseEvent) {
</script>
<template>
<N8nCard :class="$style.card" @click="handleClick">
<N8nCard :class="[$style.card, props.archived && $style.cardArchived]" @click="handleClick">
<template #prepend>
<N8nIcon :icon="icon" size="large" :class="$style.icon" />
</template>
<template #header>
<N8nText>{{ props.name }}</N8nText>
<span v-if="props.archived" :class="$style.archivedBadge">
{{ i18n.baseText('instanceAi.artifactsPanel.archived') }}
</span>
</template>
<N8nText v-if="props.metadata" :class="$style.metadata">{{ props.metadata }}</N8nText>
</N8nCard>
@@ -68,6 +75,19 @@ function handleClick(e: MouseEvent) {
}
}
.cardArchived {
opacity: 0.55;
}
.archivedBadge {
font-size: var(--font-size--3xs);
color: var(--color--text--tint-1);
background: var(--color--foreground--tint-1);
padding: var(--spacing--5xs) var(--spacing--3xs);
border-radius: var(--radius--sm);
margin-left: var(--spacing--2xs);
}
.icon {
color: var(--icon-color--strong);
flex-shrink: 0;
@@ -82,7 +82,7 @@ const artifactIconMap: Record<string, IconName> = {
<div
v-for="artifact in artifacts"
:key="artifact.id"
:class="$style.artifactRow"
:class="[$style.artifactRow, artifact.archived && $style.artifactRowArchived]"
@click="handleArtifactClick(artifact, $event)"
>
<span :class="$style.artifactIconWrap">
@@ -93,6 +93,9 @@ const artifactIconMap: Record<string, IconName> = {
/>
</span>
<span :class="$style.artifactName">{{ artifact.name }}</span>
<span v-if="artifact.archived" :class="$style.archivedBadge">
{{ i18n.baseText('instanceAi.artifactsPanel.archived') }}
</span>
</div>
</div>
@@ -225,6 +228,25 @@ const artifactIconMap: Record<string, IconName> = {
overflow: hidden;
text-overflow: ellipsis;
white-space: nowrap;
flex: 1;
min-width: 0;
}
.artifactRowArchived {
opacity: 0.55;
.artifactName {
text-decoration: line-through;
}
}
.archivedBadge {
font-size: var(--font-size--3xs);
color: var(--color--text--tint-1);
background: var(--color--foreground--tint-1);
padding: var(--spacing--5xs) var(--spacing--3xs);
border-radius: var(--radius--sm);
flex-shrink: 0;
}
/* Empty state */
@@ -135,6 +135,7 @@ export const useInstanceAiStore = defineStore('instanceAi', () => {
const lastEventIdByThread = ref<Record<string, number>>({});
const activeRunId = ref<string | null>(null);
const messages = ref<InstanceAiMessage[]>([]);
const archivedWorkflowIds = ref<Set<string>>(new Set());
const latestTasks = ref<TaskList | null>(null);
const hydratingThreadId = ref<string | null>(null);
const pendingMessageCount = ref(0);
@@ -169,6 +170,7 @@ export const useInstanceAiStore = defineStore('instanceAi', () => {
const { producedArtifacts, resourceNameIndex } = useResourceRegistry(
() => messages.value,
(id) => workflowsListStore.getWorkflowById(id)?.name,
() => archivedWorkflowIds.value,
);
// Response feedback — rateability selector + submission
@@ -349,6 +351,15 @@ export const useInstanceAiStore = defineStore('instanceAi', () => {
thread.title = parsed.data.payload.title;
}
}
if (parsed.data.type === 'run-finish') {
const ids = parsed.data.payload.archivedWorkflowIds;
if (ids && ids.length > 0) {
// Reassign instead of mutating: Set.add() on a ref doesn't trigger reactivity.
const next = new Set(archivedWorkflowIds.value);
for (const id of ids) next.add(id);
archivedWorkflowIds.value = next;
}
}
// Force Vue reactivity when streaming state changes (run-start can
// re-activate a completed message for auto-follow-up runs, run-finish
// marks it done). In-place mutation of message properties may not
@@ -514,6 +525,7 @@ export const useInstanceAiStore = defineStore('instanceAi', () => {
function resetThreadRuntimeState(nextHydratingThreadId: string | null): void {
hydratingThreadId.value = nextHydratingThreadId;
messages.value = [];
archivedWorkflowIds.value = new Set();
latestTasks.value = null;
activeRunId.value = null;
debugEvents.value = [];
@@ -12,6 +12,13 @@ export interface ResourceEntry {
createdAt?: string;
updatedAt?: string;
projectId?: string;
/**
* Set to true when the run-finish reap archived this workflow — a
* stepping-stone the agent created but never promoted to the main
* deliverable. The artifacts panel renders these as dimmed with an
* "Archived" label.
*/
archived?: boolean;
}
// ---------------------------------------------------------------------------
@@ -234,6 +241,7 @@ function enrichWorkflowNames(
export function useResourceRegistry(
messages: () => InstanceAiMessage[],
workflowNameLookup?: (id: string) => string | undefined,
archivedWorkflowIds?: () => ReadonlySet<string>,
) {
const collections = computed((): Collections => {
const col: Collections = {
@@ -250,6 +258,15 @@ export function useResourceRegistry(
enrichWorkflowNames(col, workflowNameLookup);
}
const archived = archivedWorkflowIds?.();
if (archived && archived.size > 0) {
for (const entry of col.produced.values()) {
if (entry.type === 'workflow' && archived.has(entry.id)) {
entry.archived = true;
}
}
}
return col;
});