feat(workflow): mirror the workflow row to core synchronously on activation

This commit is contained in:
Thomas Trompette committed 2026-09-15 18:39:19 +02:00
1 parent 80875d7ec4
commit 4aa910f625
8 files changed
+181

No files matched your search

@@ -167,6 +167,28 @@ export class WorkflowCoreSyncService {
);
}
async mirrorWorkflowById(
workspaceId: string,
workflowId: string,
): Promise<void> {
const workflow = await this.workspaceOrmManager.executeInWorkspaceContext(
async () => {
return await this.workspaceOrmManager
.getRepository<WorkflowWorkspaceEntity>('workflow', {
shouldBypassPermissionChecks: true,
})
.findOne({ where: { id: workflowId } });
},
buildSystemAuthContext(workspaceId),
);
if (!isDefined(workflow)) {
return;
}
await this.upsertToCore(workspaceId, [workflow]);
}
async deleteFromCore(
workspaceId: string,
coreWorkflowIds: string[],
@@ -3,6 +3,7 @@ import { Injectable } from '@nestjs/common';
import { type ToolSet } from 'ai';
import { RecordPositionService } from 'src/engine/core-modules/record-position/services/record-position.service';
import { WorkflowCoreSyncService } from 'src/engine/core-modules/workflow/services/workflow-core-sync.service';
import { WorkflowVersionCoreSyncService } from 'src/engine/core-modules/workflow/services/workflow-version-core-sync.service';
import { AgentService } from 'src/engine/metadata-modules/ai/ai-agent/agent.service';
import { WorkspaceManyOrAllFlatEntityMapsCacheService } from 'src/engine/metadata-modules/flat-entity/services/workspace-many-or-all-flat-entity-maps-cache.service';
@@ -60,6 +61,7 @@ export class WorkflowToolWorkspaceService {
agentService: AgentService,
workflowCommonService: WorkflowCommonWorkspaceService,
workflowVersionCoreSyncService: WorkflowVersionCoreSyncService,
workflowCoreSyncService: WorkflowCoreSyncService,
) {
this.deps = {
workflowVersionStepService,
@@ -76,6 +78,7 @@ export class WorkflowToolWorkspaceService {
agentService,
workflowCommonService,
workflowVersionCoreSyncService,
workflowCoreSyncService,
};
}
@@ -60,6 +60,7 @@ type CreateCompleteWorkflowToolDeps = Pick<
| 'workspaceOrmManager'
| 'recordPositionService'
| 'workflowVersionCoreSyncService'
| 'workflowCoreSyncService'
>;
type CreateCompleteWorkflowToolContext = WorkflowToolContext & {
@@ -324,4 +325,9 @@ const updateWorkflowStatus = async ({
lastPublishedVersionId: workflowVersionId,
});
}, authContext);
await deps.workflowCoreSyncService.mirrorWorkflowById(
context.workspaceId,
workflowId,
);
};
@@ -1,4 +1,5 @@
import type { RecordPositionService } from 'src/engine/core-modules/record-position/services/record-position.service';
import type { WorkflowCoreSyncService } from 'src/engine/core-modules/workflow/services/workflow-core-sync.service';
import type { WorkflowVersionCoreSyncService } from 'src/engine/core-modules/workflow/services/workflow-version-core-sync.service';
import type { AgentService } from 'src/engine/metadata-modules/ai/ai-agent/agent.service';
import type { LogicFunctionFromSourceService } from 'src/engine/metadata-modules/logic-function/services/logic-function-from-source.service';
@@ -28,6 +29,7 @@ export type WorkflowToolDependencies = {
agentService: AgentService;
workflowCommonService: WorkflowCommonWorkspaceService;
workflowVersionCoreSyncService: WorkflowVersionCoreSyncService;
workflowCoreSyncService: WorkflowCoreSyncService;
};
export type WorkflowToolContext = {
@@ -5,6 +5,7 @@ import { WORKFLOW_TOOL_SERVICE_TOKEN } from 'src/engine/core-modules/tool-provid
import { AiAgentModule } from 'src/engine/metadata-modules/ai/ai-agent/ai-agent.module';
import { WorkspaceManyOrAllFlatEntityMapsCacheModule } from 'src/engine/metadata-modules/flat-entity/services/workspace-many-or-all-flat-entity-maps-cache.module';
import { LogicFunctionModule } from 'src/engine/metadata-modules/logic-function/logic-function.module';
import { WorkflowCoreModule } from 'src/engine/core-modules/workflow/workflow-core.module';
import { WorkflowVersionCoreModule } from 'src/engine/core-modules/workflow/workflow-version-core.module';
import { WorkflowCommonModule } from 'src/modules/workflow/common/workflow-common.module';
import { WorkflowSchemaModule } from 'src/modules/workflow/workflow-builder/workflow-schema/workflow-schema.module';
@@ -33,6 +34,7 @@ import { WorkflowToolWorkspaceService } from './services/workflow-tool.workspace
WorkspaceManyOrAllFlatEntityMapsCacheModule,
AiAgentModule,
WorkflowVersionCoreModule,
WorkflowCoreModule,
],
providers: [
WorkflowToolWorkspaceService,
@@ -6,6 +6,7 @@ import { CacheStorageModule } from 'src/engine/core-modules/cache-storage/cache-
import { FeatureFlagModule } from 'src/engine/core-modules/feature-flag/feature-flag.module';
import { CommandMenuItemModule } from 'src/engine/metadata-modules/command-menu-item/command-menu-item.module';
import { LogicFunctionModule } from 'src/engine/metadata-modules/logic-function/logic-function.module';
import { WorkflowCoreModule } from 'src/engine/core-modules/workflow/workflow-core.module';
import { WorkflowVersionCoreModule } from 'src/engine/core-modules/workflow/workflow-version-core.module';
import { WorkflowCommonModule } from 'src/modules/workflow/common/workflow-common.module';
import { CodeStepBuildModule } from 'src/modules/workflow/workflow-builder/workflow-version-step/code-step/code-step-build.module';
@@ -27,6 +28,7 @@ import { WorkflowTriggerWorkspaceService } from 'src/modules/workflow/workflow-t
FeatureFlagModule,
LogicFunctionModule,
WorkflowVersionCoreModule,
WorkflowCoreModule,
WorkflowVersionValidationModule,
],
providers: [WorkflowTriggerWorkspaceService, WorkflowTriggerJob],
@@ -11,6 +11,7 @@ import { InjectCacheStorage } from 'src/engine/core-modules/cache-storage/decora
import { CacheStorageService } from 'src/engine/core-modules/cache-storage/services/cache-storage.service';
import { CacheStorageNamespace } from 'src/engine/core-modules/cache-storage/types/cache-storage-namespace.enum';
import { buildCoreDispatchIds } from 'src/engine/core-modules/workflow/utils/build-core-dispatch-ids.util';
import { WorkflowCoreSyncService } from 'src/engine/core-modules/workflow/services/workflow-core-sync.service';
import { WorkflowVersionCoreSyncService } from 'src/engine/core-modules/workflow/services/workflow-version-core-sync.service';
import { CommandMenuItemService } from 'src/engine/metadata-modules/command-menu-item/command-menu-item.service';
import { EngineComponentKey } from 'src/engine/metadata-modules/command-menu-item/enums/engine-component-key.enum';
@@ -68,6 +69,7 @@ export class WorkflowTriggerWorkspaceService {
private readonly workspaceEventEmitter: WorkspaceEventEmitter,
private readonly commandMenuItemService: CommandMenuItemService,
private readonly workflowVersionCoreSyncService: WorkflowVersionCoreSyncService,
private readonly workflowCoreSyncService: WorkflowCoreSyncService,
@InjectCacheStorage(CacheStorageNamespace.ModuleWorkflow)
private readonly cacheStorageService: CacheStorageService,
) {}
@@ -368,6 +370,11 @@ export class WorkflowTriggerWorkspaceService {
},
);
await this.workflowCoreSyncService.mirrorWorkflowById(
workspaceId,
workflow.id,
);
await this.createOrUpdateCommandMenuItem(
workflow,
workflowVersion,
@@ -0,0 +1,137 @@
import { updateWorkflowVersionTrigger } from 'test/integration/graphql/suites/workflow/utils/update-workflow-version-trigger.util';
import { workflowGraphqlRequest } from 'test/integration/graphql/suites/workflow/utils/workflow-graphql-request.util';
const GET_CORE_WORKFLOW_QUERY = `
query GetCoreWorkflow($workspaceWorkflowId: UUID!) {
coreWorkflow(workspaceWorkflowId: $workspaceWorkflowId) {
id
statuses
lastPublishedVersionId
}
}
`;
describe('workflow activation mirrors the core workflow row synchronously (e2e)', () => {
let workspaceWorkflowId: string;
let workflowVersionId: string;
beforeAll(async () => {
const createResponse = await workflowGraphqlRequest(`
mutation {
createCoreWorkflow(input: { name: "Activation Core Mirror" }) {
id
workspaceWorkflowId
}
}
`);
expect(createResponse.body.errors).toBeUndefined();
workspaceWorkflowId =
createResponse.body.data.createCoreWorkflow.workspaceWorkflowId;
const versionsResponse = await workflowGraphqlRequest(
`
query GetWorkflow($id: UUID!) {
workflow(filter: { id: { eq: $id } }) {
versions {
edges {
node {
id
}
}
}
}
}
`,
{ id: workspaceWorkflowId },
);
expect(versionsResponse.body.errors).toBeUndefined();
workflowVersionId =
versionsResponse.body.data.workflow.versions.edges[0].node.id;
await updateWorkflowVersionTrigger({
workflowVersionId,
trigger: {
name: 'Webhook Trigger',
type: 'WEBHOOK',
settings: {
outputSchema: {},
httpMethod: 'GET',
authentication: null,
},
nextStepIds: [],
position: { x: 0, y: 0 },
},
});
const createStepResponse = await workflowGraphqlRequest(
`
mutation CreateWorkflowVersionStep(
$input: CreateWorkflowVersionStepInput!
) {
createWorkflowVersionStep(input: $input) {
stepsDiff
}
}
`,
{
input: {
workflowVersionId,
stepType: 'CODE',
parentStepId: 'trigger',
position: { x: 200, y: 0 },
},
},
);
expect(createStepResponse.body.errors).toBeUndefined();
});
afterAll(async () => {
await workflowGraphqlRequest(
`
mutation DestroyWorkflow($id: ID!) {
destroyWorkflow(id: $id) {
id
}
}
`,
{ id: workspaceWorkflowId },
);
});
it('starts with no published version on the core row', async () => {
const response = await workflowGraphqlRequest(GET_CORE_WORKFLOW_QUERY, {
workspaceWorkflowId,
});
expect(response.body.errors).toBeUndefined();
expect(response.body.data.coreWorkflow.lastPublishedVersionId).toBeNull();
});
it('exposes the published version on the core row as soon as activation returns', async () => {
const activateResponse = await workflowGraphqlRequest(
`
mutation ActivateWorkflowVersion($workflowVersionId: UUID!) {
activateWorkflowVersion(workflowVersionId: $workflowVersionId)
}
`,
{ workflowVersionId },
);
expect(activateResponse.body.errors).toBeUndefined();
const response = await workflowGraphqlRequest(GET_CORE_WORKFLOW_QUERY, {
workspaceWorkflowId,
});
expect(response.body.errors).toBeUndefined();
expect(response.body.data.coreWorkflow.lastPublishedVersionId).toBe(
workflowVersionId,
);
expect(response.body.data.coreWorkflow.statuses).toContain('ACTIVE');
});
});