mirror of
https://github.com/n8n-io/n8n.git
synced 2026-09-24 23:22:38 +08:00
feat(core): Publish workflow activations via the outbox (no-changelog) (#31996)
This commit is contained in:
@@ -1,6 +1,7 @@
|
||||
import { GlobalConfig } from '@n8n/config';
|
||||
import { Service } from '@n8n/di';
|
||||
import { DataSource, Repository } from '@n8n/typeorm';
|
||||
import type { EntityManager } from '@n8n/typeorm';
|
||||
import { UnexpectedError } from 'n8n-workflow';
|
||||
|
||||
import {
|
||||
@@ -26,25 +27,33 @@ export class WorkflowPublicationOutboxRepository extends Repository<WorkflowPubl
|
||||
* (`workflowId` where `status = 'pending'`) guarantees at most one pending
|
||||
* record per workflow without an explicit transaction. Callers only need to
|
||||
* know the enqueue succeeded, so no row is returned.
|
||||
*
|
||||
* Pass `trx` to run the UPSERT inside an existing transaction, e.g. to make
|
||||
* the enqueue atomic with a `workflow_entity` update.
|
||||
*/
|
||||
async enqueue(workflowId: string, publishedVersionId: string): Promise<void> {
|
||||
async enqueue(
|
||||
workflowId: string,
|
||||
publishedVersionId: string,
|
||||
trx?: EntityManager,
|
||||
): Promise<void> {
|
||||
if (this.globalConfig.database.type === 'postgresdb') {
|
||||
await this.enqueueWithPostgresUpsert(workflowId, publishedVersionId);
|
||||
await this.enqueueWithPostgresUpsert(workflowId, publishedVersionId, trx ?? this.manager);
|
||||
return;
|
||||
}
|
||||
|
||||
await this.enqueueWithSqliteUpsert(workflowId, publishedVersionId);
|
||||
await this.enqueueWithSqliteUpsert(workflowId, publishedVersionId, trx ?? this.manager);
|
||||
}
|
||||
|
||||
private async enqueueWithPostgresUpsert(
|
||||
workflowId: string,
|
||||
publishedVersionId: string,
|
||||
trx: EntityManager,
|
||||
): Promise<void> {
|
||||
const tableName = this.getTableName('workflow_publication_outbox');
|
||||
|
||||
// `createdAt`/`updatedAt` carry DB-level defaults, so the insert omits
|
||||
// them; the conflict path bumps `updatedAt` explicitly.
|
||||
await this.query(
|
||||
await trx.query(
|
||||
`INSERT INTO ${tableName} ("workflowId", "publishedVersionId", "status")
|
||||
VALUES ($1, $2, '${Status.Pending}')
|
||||
ON CONFLICT ("workflowId", "status") WHERE "status" IN ('${Status.Pending}', '${Status.InProgress}')
|
||||
@@ -56,10 +65,11 @@ export class WorkflowPublicationOutboxRepository extends Repository<WorkflowPubl
|
||||
private async enqueueWithSqliteUpsert(
|
||||
workflowId: string,
|
||||
publishedVersionId: string,
|
||||
trx: EntityManager,
|
||||
): Promise<void> {
|
||||
const tableName = this.getTableName('workflow_publication_outbox');
|
||||
|
||||
await this.query(
|
||||
await trx.query(
|
||||
`INSERT INTO ${tableName} ("workflowId", "publishedVersionId", "status")
|
||||
VALUES (?, ?, '${Status.Pending}')
|
||||
ON CONFLICT ("workflowId", "status") WHERE "status" IN ('${Status.Pending}', '${Status.InProgress}')
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import { Service } from '@n8n/di';
|
||||
import { DataSource, Repository } from '@n8n/typeorm';
|
||||
import type { EntityManager } from '@n8n/typeorm';
|
||||
|
||||
import { WorkflowPublishHistory } from '../entities';
|
||||
|
||||
@@ -9,13 +10,17 @@ export class WorkflowPublishHistoryRepository extends Repository<WorkflowPublish
|
||||
super(WorkflowPublishHistory, dataSource.manager);
|
||||
}
|
||||
|
||||
async addRecord({
|
||||
workflowId,
|
||||
versionId,
|
||||
event,
|
||||
userId,
|
||||
}: Pick<WorkflowPublishHistory, 'event' | 'workflowId' | 'versionId' | 'userId'>) {
|
||||
await this.insert({
|
||||
async addRecord(
|
||||
{
|
||||
workflowId,
|
||||
versionId,
|
||||
event,
|
||||
userId,
|
||||
}: Pick<WorkflowPublishHistory, 'event' | 'workflowId' | 'versionId' | 'userId'>,
|
||||
trx?: EntityManager,
|
||||
) {
|
||||
const repository = trx ? trx.getRepository(WorkflowPublishHistory) : this;
|
||||
await repository.insert({
|
||||
workflowId,
|
||||
versionId,
|
||||
event,
|
||||
|
||||
@@ -153,6 +153,7 @@ describe('Start - AuthRolesService initialization', () => {
|
||||
cache: { backend: 'memory' },
|
||||
taskRunners: {},
|
||||
expressionEngine: { engine: 'legacy', poolSize: 1, maxCodeCacheSize: 1024 },
|
||||
workflows: { useWorkflowPublicationService: false },
|
||||
};
|
||||
// @ts-expect-error - Accessing protected method for testing
|
||||
start.initCrashJournal = jest.fn().mockResolvedValue(undefined);
|
||||
@@ -206,6 +207,7 @@ describe('Start - AuthRolesService initialization', () => {
|
||||
cache: { backend: 'memory' },
|
||||
taskRunners: {},
|
||||
expressionEngine: { engine: 'legacy', poolSize: 1, maxCodeCacheSize: 1024 },
|
||||
workflows: { useWorkflowPublicationService: false },
|
||||
};
|
||||
|
||||
await start.init();
|
||||
@@ -240,6 +242,7 @@ describe('Start - AuthRolesService initialization', () => {
|
||||
cache: { backend: 'memory' },
|
||||
taskRunners: {},
|
||||
expressionEngine: { engine: 'legacy', poolSize: 1, maxCodeCacheSize: 1024 },
|
||||
workflows: { useWorkflowPublicationService: false },
|
||||
};
|
||||
|
||||
await start.init();
|
||||
@@ -285,6 +288,7 @@ describe('Start - AuthRolesService initialization', () => {
|
||||
cache: { backend: 'memory' },
|
||||
taskRunners: {},
|
||||
expressionEngine: { engine: 'legacy' as const, poolSize: 1, maxCodeCacheSize: 1024 },
|
||||
workflows: { useWorkflowPublicationService: false },
|
||||
};
|
||||
|
||||
beforeEach(() => {
|
||||
|
||||
@@ -232,6 +232,13 @@ export class Start extends BaseCommand<z.infer<typeof flagsSchema>> {
|
||||
await this.initOrchestration();
|
||||
}
|
||||
|
||||
if (this.globalConfig.workflows.useWorkflowPublicationService) {
|
||||
const { WorkflowPublicationOutboxConsumer } = await import(
|
||||
'@/workflows/workflow-publication-outbox-consumer'
|
||||
);
|
||||
Container.get(WorkflowPublicationOutboxConsumer).init();
|
||||
}
|
||||
|
||||
await this.instanceSettings.initialize(Container.get(DeploymentKeyRepository));
|
||||
await Container.get(JwtService).initialize(Container.get(DeploymentKeyRepository));
|
||||
await Container.get(BinaryDataConfig).initialize(Container.get(DeploymentKeyRepository));
|
||||
|
||||
@@ -13,7 +13,7 @@ import type {
|
||||
} from '@n8n/db';
|
||||
import type { EntityManager } from '@n8n/typeorm';
|
||||
import { mock } from 'jest-mock-extended';
|
||||
import type { ErrorReporter } from 'n8n-core';
|
||||
import type { ErrorReporter, InstanceSettings } from 'n8n-core';
|
||||
import type { INode } from 'n8n-workflow';
|
||||
|
||||
import type { ActivationErrorsService } from '@/activation-errors.service';
|
||||
@@ -31,6 +31,7 @@ describe('WorkflowPublicationOutboxConsumer', () => {
|
||||
const workflowPublishedVersionRepository = mock<WorkflowPublishedVersionRepository>();
|
||||
const activeWorkflowManager = mock<ActiveWorkflowManager>();
|
||||
const activationErrorsService = mock<ActivationErrorsService>();
|
||||
const instanceSettings = mock<InstanceSettings>();
|
||||
|
||||
let consumer: WorkflowPublicationOutboxConsumer;
|
||||
|
||||
@@ -51,6 +52,7 @@ describe('WorkflowPublicationOutboxConsumer', () => {
|
||||
workflowPublishedVersionRepository,
|
||||
activeWorkflowManager,
|
||||
activationErrorsService,
|
||||
instanceSettings,
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
@@ -1,7 +1,15 @@
|
||||
import type { LicenseState } from '@n8n/backend-common';
|
||||
import type { Project, User, WorkflowRepository, WorkflowPublishHistoryRepository } from '@n8n/db';
|
||||
import type { GlobalConfig, WorkflowsConfig } from '@n8n/config';
|
||||
import type {
|
||||
Project,
|
||||
User,
|
||||
WorkflowRepository,
|
||||
WorkflowPublishHistoryRepository,
|
||||
WorkflowPublicationOutboxRepository,
|
||||
} from '@n8n/db';
|
||||
import { WorkflowEntity, WorkflowHistory } from '@n8n/db';
|
||||
import type { Scope } from '@n8n/permissions';
|
||||
import type { EntityManager } from '@n8n/typeorm';
|
||||
import type { MockProxy } from 'jest-mock-extended';
|
||||
import { mock } from 'jest-mock-extended';
|
||||
import type { IConnections, INode } from 'n8n-workflow';
|
||||
@@ -10,6 +18,7 @@ import type { ActiveWorkflowManager } from '@/active-workflow-manager';
|
||||
import { BadRequestError } from '@/errors/response-errors/bad-request.error';
|
||||
import { UnprocessableRequestError } from '@/errors/response-errors/unprocessable.error';
|
||||
import { WorkflowActivationBadRequestError } from '@/errors/response-errors/workflow-activation-bad-request.error';
|
||||
import type { EventService } from '@/events/event.service';
|
||||
import type { ExternalHooks } from '@/external-hooks';
|
||||
import type { RedactionEnforcementService } from '@/modules/redaction/redaction-enforcement.service';
|
||||
import { userHasScopes } from '@/permissions.ee/check-access';
|
||||
@@ -67,6 +76,7 @@ describe('WorkflowService', () => {
|
||||
mock(), // workflowFinderService
|
||||
mock(), // workflowPublishedVersionRepository
|
||||
mock(), // workflowPublishHistoryRepository
|
||||
mock(), // outboxRepository
|
||||
Object.assign(mock<WorkflowValidationService>(), {
|
||||
validateCredentialNodeRestrictions: () => ({ isValid: true }),
|
||||
}), // workflowValidationService
|
||||
@@ -218,6 +228,7 @@ describe('WorkflowService', () => {
|
||||
workflowFinderServiceMock, // workflowFinderService
|
||||
mock(), // workflowPublishedVersionRepository
|
||||
mock(), // workflowPublishHistoryRepository
|
||||
mock(), // outboxRepository
|
||||
Object.assign(mock<WorkflowValidationService>(), {
|
||||
validateCredentialNodeRestrictions: () => ({ isValid: true }),
|
||||
}), // workflowValidationService
|
||||
@@ -785,8 +796,11 @@ describe('WorkflowService', () => {
|
||||
let workflowHistoryServiceMock: MockProxy<WorkflowHistoryService>;
|
||||
let workflowRepositoryMock: MockProxy<WorkflowRepository>;
|
||||
let workflowPublishHistoryRepositoryMock: MockProxy<WorkflowPublishHistoryRepository>;
|
||||
let outboxRepositoryMock: MockProxy<WorkflowPublicationOutboxRepository>;
|
||||
let globalConfigMock: MockProxy<GlobalConfig>;
|
||||
let activeWorkflowManagerMock: MockProxy<ActiveWorkflowManager>;
|
||||
let externalHooksMock: MockProxy<ExternalHooks>;
|
||||
let eventServiceMock: MockProxy<EventService>;
|
||||
|
||||
const WORKFLOW_ID = 'workflow-1';
|
||||
const PREVIOUS_VERSION_ID = 'v1';
|
||||
@@ -821,8 +835,13 @@ describe('WorkflowService', () => {
|
||||
workflowHistoryServiceMock = mock<WorkflowHistoryService>();
|
||||
workflowRepositoryMock = mock();
|
||||
workflowPublishHistoryRepositoryMock = mock();
|
||||
outboxRepositoryMock = mock();
|
||||
globalConfigMock = mock<GlobalConfig>({
|
||||
workflows: mock<WorkflowsConfig>({ useWorkflowPublicationService: false }),
|
||||
});
|
||||
activeWorkflowManagerMock = mock();
|
||||
externalHooksMock = mock<ExternalHooks>();
|
||||
eventServiceMock = mock<EventService>();
|
||||
|
||||
workflowRepositoryMock.create.mockImplementation(
|
||||
(data) => Object.assign(new WorkflowEntity(), data) as WorkflowEntity,
|
||||
@@ -842,12 +861,13 @@ describe('WorkflowService', () => {
|
||||
mock(), // roleService
|
||||
mock(), // projectService
|
||||
mock(), // executionRepository
|
||||
mock(), // eventService
|
||||
mock(), // globalConfig
|
||||
eventServiceMock, // eventService
|
||||
globalConfigMock, // globalConfig
|
||||
mock(), // folderRepository
|
||||
workflowFinderServiceMock, // workflowFinderService
|
||||
mock(), // workflowPublishedVersionRepository
|
||||
workflowPublishHistoryRepositoryMock, // workflowPublishHistoryRepository
|
||||
outboxRepositoryMock, // outboxRepository
|
||||
Object.assign(mock<WorkflowValidationService>(), {
|
||||
validateCredentialNodeRestrictions: () => ({ isValid: true }),
|
||||
}), // workflowValidationService
|
||||
@@ -941,5 +961,87 @@ describe('WorkflowService', () => {
|
||||
expect(candidate.nodes).toBe(workflow.nodes);
|
||||
expect(candidate.connections).toBe(workflow.connections);
|
||||
});
|
||||
|
||||
test('with the publication outbox enabled, updates the version, writes history, enqueues and emits events without touching the active workflow manager', async () => {
|
||||
globalConfigMock.workflows.useWorkflowPublicationService = true;
|
||||
|
||||
const workflow = makeWorkflowEntity({ activeVersionId: PREVIOUS_VERSION_ID });
|
||||
const versionToActivate = makeVersionToActivate();
|
||||
workflowFinderServiceMock.findWorkflowForUser.mockResolvedValue(workflow);
|
||||
workflowHistoryServiceMock.getVersion.mockResolvedValue(versionToActivate);
|
||||
workflowRepositoryMock.findOne.mockResolvedValue(workflow);
|
||||
externalHooksMock.run.mockResolvedValue(undefined);
|
||||
|
||||
const trx = mock<EntityManager>();
|
||||
const managerMock = mock<EntityManager>();
|
||||
(managerMock.transaction as unknown as jest.Mock).mockImplementation(
|
||||
async (runInTransaction: (entityManager: EntityManager) => Promise<unknown>) =>
|
||||
await runInTransaction(trx),
|
||||
);
|
||||
Object.defineProperty(workflowRepositoryMock, 'manager', {
|
||||
value: managerMock,
|
||||
configurable: true,
|
||||
});
|
||||
|
||||
const addToActiveWorkflowManagerSpy = jest.spyOn(
|
||||
workflowService as never,
|
||||
'_addToActiveWorkflowManager',
|
||||
);
|
||||
|
||||
const user = mock<User>({ id: 'user-1' });
|
||||
|
||||
await workflowService.activateWorkflow(user, WORKFLOW_ID, {
|
||||
versionId: TARGET_VERSION_ID,
|
||||
});
|
||||
|
||||
// activeVersionId + active are updated inside the transaction
|
||||
expect(trx.update).toHaveBeenCalledWith(
|
||||
WorkflowEntity,
|
||||
{ id: WORKFLOW_ID },
|
||||
expect.objectContaining({ active: true, activeVersionId: TARGET_VERSION_ID }),
|
||||
);
|
||||
// the outbox record is enqueued in the same transaction
|
||||
expect(outboxRepositoryMock.enqueue).toHaveBeenCalledWith(
|
||||
WORKFLOW_ID,
|
||||
TARGET_VERSION_ID,
|
||||
trx,
|
||||
);
|
||||
// publish-history records (deactivated for the previous version, activated for the
|
||||
// target) are written in the same transaction
|
||||
expect(workflowPublishHistoryRepositoryMock.addRecord).toHaveBeenCalledWith(
|
||||
expect.objectContaining({ event: 'deactivated', versionId: PREVIOUS_VERSION_ID }),
|
||||
trx,
|
||||
);
|
||||
expect(workflowPublishHistoryRepositoryMock.addRecord).toHaveBeenCalledWith(
|
||||
expect.objectContaining({ event: 'activated', versionId: TARGET_VERSION_ID }),
|
||||
trx,
|
||||
);
|
||||
expect(eventServiceMock.emit).toHaveBeenNthCalledWith(1, 'workflow-deactivated', {
|
||||
user,
|
||||
workflowId: WORKFLOW_ID,
|
||||
workflow,
|
||||
publicApi: false,
|
||||
deactivatedVersionId: PREVIOUS_VERSION_ID,
|
||||
source: 'ui',
|
||||
});
|
||||
expect(eventServiceMock.emit).toHaveBeenNthCalledWith(2, 'workflow-activated', {
|
||||
user,
|
||||
workflowId: WORKFLOW_ID,
|
||||
workflow: expect.objectContaining({
|
||||
active: true,
|
||||
activeVersionId: TARGET_VERSION_ID,
|
||||
activeVersion: versionToActivate,
|
||||
nodes: versionToActivate.nodes,
|
||||
connections: versionToActivate.connections,
|
||||
}),
|
||||
publicApi: false,
|
||||
source: 'ui',
|
||||
});
|
||||
// trigger reapplication is deferred to the consumer
|
||||
expect(addToActiveWorkflowManagerSpy).not.toHaveBeenCalled();
|
||||
expect(activeWorkflowManagerMock.add).not.toHaveBeenCalled();
|
||||
expect(activeWorkflowManagerMock.remove).not.toHaveBeenCalled();
|
||||
expect(workflowRepositoryMock.update).not.toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
@@ -12,7 +12,7 @@ import {
|
||||
} from '@n8n/db';
|
||||
import { OnLeaderStepdown, OnLeaderTakeover, OnShutdown } from '@n8n/decorators';
|
||||
import { Service } from '@n8n/di';
|
||||
import { ErrorReporter } from 'n8n-core';
|
||||
import { ErrorReporter, InstanceSettings } from 'n8n-core';
|
||||
import { UnexpectedError, ensureError } from 'n8n-workflow';
|
||||
|
||||
import { ActivationErrorsService } from '@/activation-errors.service';
|
||||
@@ -48,10 +48,17 @@ export class WorkflowPublicationOutboxConsumer {
|
||||
private readonly workflowPublishedVersionRepository: WorkflowPublishedVersionRepository,
|
||||
private readonly activeWorkflowManager: ActiveWorkflowManager,
|
||||
private readonly activationErrorsService: ActivationErrorsService,
|
||||
private readonly instanceSettings: InstanceSettings,
|
||||
) {
|
||||
this.logger = this.logger.scoped('workflow-publication');
|
||||
}
|
||||
|
||||
init() {
|
||||
if (this.instanceSettings.isLeader) {
|
||||
this.startPolling();
|
||||
}
|
||||
}
|
||||
|
||||
@OnLeaderTakeover()
|
||||
startPolling() {
|
||||
if (!this.workflowsConfig.useWorkflowPublicationService || this.isShuttingDown) return;
|
||||
|
||||
@@ -1,15 +1,10 @@
|
||||
import { UpdateWorkflowHistoryVersionDto } from '@n8n/api-types';
|
||||
import { LicenseState, Logger } from '@n8n/backend-common';
|
||||
import { GlobalConfig } from '@n8n/config';
|
||||
import type {
|
||||
User,
|
||||
ListQueryDb,
|
||||
WorkflowFolderUnionFull,
|
||||
WorkflowHistory,
|
||||
WorkflowEntity,
|
||||
} from '@n8n/db';
|
||||
import type { User, ListQueryDb, WorkflowFolderUnionFull, WorkflowHistory } from '@n8n/db';
|
||||
import {
|
||||
SharedWorkflow,
|
||||
WorkflowEntity,
|
||||
ExecutionRepository,
|
||||
FolderRepository,
|
||||
WorkflowTagMappingRepository,
|
||||
@@ -17,6 +12,7 @@ import {
|
||||
WorkflowRepository,
|
||||
WorkflowPublishedVersionRepository,
|
||||
WorkflowPublishHistoryRepository,
|
||||
WorkflowPublicationOutboxRepository,
|
||||
ProjectRepository,
|
||||
} from '@n8n/db';
|
||||
import { Container, Service } from '@n8n/di';
|
||||
@@ -94,6 +90,7 @@ export class WorkflowService {
|
||||
private readonly workflowFinderService: WorkflowFinderService,
|
||||
private readonly workflowPublishedVersionRepository: WorkflowPublishedVersionRepository,
|
||||
private readonly workflowPublishHistoryRepository: WorkflowPublishHistoryRepository,
|
||||
private readonly outboxRepository: WorkflowPublicationOutboxRepository,
|
||||
private readonly workflowValidationService: WorkflowValidationService,
|
||||
private readonly nodeTypes: NodeTypes,
|
||||
private readonly webhookService: WebhookService,
|
||||
@@ -646,16 +643,6 @@ export class WorkflowService {
|
||||
if (didPublish) {
|
||||
assert(workflow.activeVersionId !== null);
|
||||
|
||||
// Temporary: In the future, the workflow publication service will
|
||||
// manage this mapping. We set it here for now to support incremental
|
||||
// development and testing.
|
||||
if (this.globalConfig.workflows.useWorkflowPublicationService) {
|
||||
await this.workflowPublishedVersionRepository.setPublishedVersion(
|
||||
workflowId,
|
||||
workflow.activeVersionId,
|
||||
);
|
||||
}
|
||||
|
||||
await this.workflowPublishHistoryRepository.addRecord({
|
||||
workflowId,
|
||||
versionId: workflow.activeVersionId,
|
||||
@@ -806,51 +793,89 @@ export class WorkflowService {
|
||||
});
|
||||
}
|
||||
|
||||
if (previousActiveVersionId) {
|
||||
await this.activeWorkflowManager.remove(workflowId);
|
||||
|
||||
this.eventService.emit('workflow-deactivated', {
|
||||
if (this.globalConfig.workflows.useWorkflowPublicationService) {
|
||||
await this._publishViaOutbox(
|
||||
user,
|
||||
workflowId,
|
||||
workflow,
|
||||
versionIdToActivate,
|
||||
previousActiveVersionId,
|
||||
workflow.updatedAt,
|
||||
);
|
||||
|
||||
if (previousActiveVersionId) {
|
||||
this.eventService.emit('workflow-deactivated', {
|
||||
user,
|
||||
workflowId,
|
||||
workflow,
|
||||
publicApi,
|
||||
deactivatedVersionId: previousActiveVersionId,
|
||||
source,
|
||||
});
|
||||
}
|
||||
|
||||
const activatedWorkflow = this.workflowRepository.create({
|
||||
...workflow,
|
||||
active: true,
|
||||
activeVersionId: versionIdToActivate,
|
||||
activeVersion: versionToActivate,
|
||||
nodes: versionToActivate.nodes,
|
||||
connections: versionToActivate.connections,
|
||||
});
|
||||
|
||||
this.eventService.emit('workflow-activated', {
|
||||
user,
|
||||
workflowId,
|
||||
workflow: activatedWorkflow,
|
||||
publicApi,
|
||||
deactivatedVersionId: previousActiveVersionId,
|
||||
source,
|
||||
});
|
||||
await this.workflowPublishHistoryRepository.addRecord({
|
||||
workflowId,
|
||||
versionId: previousActiveVersionId,
|
||||
event: 'deactivated',
|
||||
userId: user.id,
|
||||
} else {
|
||||
if (previousActiveVersionId) {
|
||||
await this.activeWorkflowManager.remove(workflowId);
|
||||
|
||||
this.eventService.emit('workflow-deactivated', {
|
||||
user,
|
||||
workflowId,
|
||||
workflow,
|
||||
publicApi,
|
||||
deactivatedVersionId: previousActiveVersionId,
|
||||
source,
|
||||
});
|
||||
await this.workflowPublishHistoryRepository.addRecord({
|
||||
workflowId,
|
||||
versionId: previousActiveVersionId,
|
||||
event: 'deactivated',
|
||||
userId: user.id,
|
||||
});
|
||||
}
|
||||
|
||||
const activationMode = previousActiveVersionId ? 'update' : 'activate';
|
||||
|
||||
await this.workflowRepository.update(workflowId, {
|
||||
activeVersionId: versionIdToActivate,
|
||||
active: true,
|
||||
// workflow content did not change, so we keep updatedAt as is
|
||||
updatedAt: workflow.updatedAt,
|
||||
});
|
||||
|
||||
const workflowForActivation = await this.workflowRepository.findOne({
|
||||
where: { id: workflowId },
|
||||
relations: ['activeVersion'],
|
||||
});
|
||||
|
||||
if (!workflowForActivation) {
|
||||
throw new NotFoundError(`Workflow with ID "${workflowId}" could not be found.`);
|
||||
}
|
||||
|
||||
await this._addToActiveWorkflowManager(
|
||||
user,
|
||||
workflowId,
|
||||
workflowForActivation,
|
||||
activationMode,
|
||||
{ source },
|
||||
);
|
||||
}
|
||||
|
||||
const activationMode = previousActiveVersionId ? 'update' : 'activate';
|
||||
|
||||
await this.workflowRepository.update(workflowId, {
|
||||
activeVersionId: versionIdToActivate,
|
||||
active: true,
|
||||
// workflow content did not change, so we keep updatedAt as is
|
||||
updatedAt: workflow.updatedAt,
|
||||
});
|
||||
|
||||
const workflowForActivation = await this.workflowRepository.findOne({
|
||||
where: { id: workflowId },
|
||||
relations: ['activeVersion'],
|
||||
});
|
||||
|
||||
if (!workflowForActivation) {
|
||||
throw new NotFoundError(`Workflow with ID "${workflowId}" could not be found.`);
|
||||
}
|
||||
|
||||
await this._addToActiveWorkflowManager(
|
||||
user,
|
||||
workflowId,
|
||||
workflowForActivation,
|
||||
activationMode,
|
||||
{ source },
|
||||
);
|
||||
|
||||
if (options?.name !== undefined || options?.description !== undefined) {
|
||||
const updateFields: UpdateWorkflowHistoryVersionDto = {};
|
||||
if (options.name !== undefined) updateFields.name = options.name;
|
||||
@@ -1318,4 +1343,55 @@ export class WorkflowService {
|
||||
throw new WorkflowValidationError(validation.error ?? 'Sub-workflow validation failed');
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Atomically records the requested version and enqueues an outbox record.
|
||||
* The publication outbox consumer reapplies the triggers and advances the
|
||||
* published version asynchronously, so we do not touch the active workflow
|
||||
* manager here.
|
||||
*/
|
||||
private async _publishViaOutbox(
|
||||
user: User,
|
||||
workflowId: string,
|
||||
versionIdToActivate: string,
|
||||
previousActiveVersionId: string | null,
|
||||
updatedAt: Date,
|
||||
): Promise<void> {
|
||||
await this.workflowRepository.manager.transaction(async (trx) => {
|
||||
await trx.update(
|
||||
WorkflowEntity,
|
||||
{ id: workflowId },
|
||||
{
|
||||
activeVersionId: versionIdToActivate,
|
||||
active: true,
|
||||
// workflow content did not change, so we keep updatedAt as is
|
||||
updatedAt,
|
||||
},
|
||||
);
|
||||
|
||||
if (previousActiveVersionId) {
|
||||
await this.workflowPublishHistoryRepository.addRecord(
|
||||
{
|
||||
workflowId,
|
||||
versionId: previousActiveVersionId,
|
||||
event: 'deactivated',
|
||||
userId: user.id,
|
||||
},
|
||||
trx,
|
||||
);
|
||||
}
|
||||
|
||||
await this.workflowPublishHistoryRepository.addRecord(
|
||||
{
|
||||
workflowId,
|
||||
versionId: versionIdToActivate,
|
||||
event: 'activated',
|
||||
userId: user.id,
|
||||
},
|
||||
trx,
|
||||
);
|
||||
|
||||
await this.outboxRepository.enqueue(workflowId, versionIdToActivate, trx);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
+21
@@ -45,6 +45,27 @@ describe('WorkflowPublicationOutboxRepository', () => {
|
||||
expect(claimedAgain).toBeNull();
|
||||
});
|
||||
|
||||
it('enqueues within a provided transaction and is visible once it commits', async () => {
|
||||
await repository.manager.transaction(async (trx) => {
|
||||
await repository.enqueue('wf-1', 'v-1', trx);
|
||||
});
|
||||
|
||||
const claimed = await repository.claimNextPendingRecord();
|
||||
expect(claimed?.workflowId).toBe('wf-1');
|
||||
expect(claimed?.publishedVersionId).toBe('v-1');
|
||||
});
|
||||
|
||||
it('discards the enqueued record when the surrounding transaction rolls back', async () => {
|
||||
await expect(
|
||||
repository.manager.transaction(async (trx) => {
|
||||
await repository.enqueue('wf-1', 'v-1', trx);
|
||||
throw new Error('rollback');
|
||||
}),
|
||||
).rejects.toThrow('rollback');
|
||||
|
||||
expect(await repository.claimNextPendingRecord()).toBeNull();
|
||||
});
|
||||
|
||||
it('claims pending records in FIFO order', async () => {
|
||||
await repository.enqueue('wf-1', 'v-1');
|
||||
await repository.enqueue('wf-2', 'v-1');
|
||||
|
||||
@@ -14,6 +14,7 @@ import {
|
||||
type WorkflowEntity,
|
||||
WorkflowPublishedVersionRepository,
|
||||
WorkflowPublishHistoryRepository,
|
||||
WorkflowPublicationOutboxRepository,
|
||||
WorkflowRepository,
|
||||
ProjectRepository,
|
||||
} from '@n8n/db';
|
||||
@@ -44,6 +45,7 @@ let workflowRepository: WorkflowRepository;
|
||||
let workflowService: WorkflowService;
|
||||
let workflowPublishedVersionRepository: WorkflowPublishedVersionRepository;
|
||||
let workflowPublishHistoryRepository: WorkflowPublishHistoryRepository;
|
||||
let outboxRepository: WorkflowPublicationOutboxRepository;
|
||||
let workflowHistoryService: WorkflowHistoryService;
|
||||
const loggerMock = mock<Logger>();
|
||||
const activeWorkflowManager = mockInstance(ActiveWorkflowManager);
|
||||
@@ -60,6 +62,7 @@ beforeAll(async () => {
|
||||
workflowRepository = Container.get(WorkflowRepository);
|
||||
workflowPublishedVersionRepository = Container.get(WorkflowPublishedVersionRepository);
|
||||
workflowPublishHistoryRepository = Container.get(WorkflowPublishHistoryRepository);
|
||||
outboxRepository = Container.get(WorkflowPublicationOutboxRepository);
|
||||
workflowHistoryService = Container.get(WorkflowHistoryService);
|
||||
workflowService = new WorkflowService(
|
||||
loggerMock,
|
||||
@@ -81,6 +84,7 @@ beforeAll(async () => {
|
||||
Container.get(WorkflowFinderService),
|
||||
workflowPublishedVersionRepository,
|
||||
workflowPublishHistoryRepository,
|
||||
outboxRepository,
|
||||
workflowValidationService,
|
||||
nodeTypes,
|
||||
webhookServiceMock,
|
||||
@@ -104,6 +108,7 @@ afterEach(async () => {
|
||||
'SharedWorkflow',
|
||||
'ProjectRelation',
|
||||
'WorkflowPublishedVersion',
|
||||
'WorkflowPublicationOutbox',
|
||||
'WorkflowEntity',
|
||||
'WorkflowHistory',
|
||||
'WorkflowPublishHistory',
|
||||
@@ -515,7 +520,7 @@ describe('deactivateWorkflow()', () => {
|
||||
});
|
||||
});
|
||||
|
||||
describe('workflow_published_version table population', () => {
|
||||
describe('workflow publication outbox', () => {
|
||||
describe('when feature flag is enabled', () => {
|
||||
beforeEach(() => {
|
||||
globalConfig.workflows.useWorkflowPublicationService = true;
|
||||
@@ -525,19 +530,31 @@ describe('workflow_published_version table population', () => {
|
||||
globalConfig.workflows.useWorkflowPublicationService = false;
|
||||
});
|
||||
|
||||
test('should write to workflow_published_version on activation', async () => {
|
||||
test('should set the active version and enqueue a pending outbox record on activation', async () => {
|
||||
const owner = await createOwner();
|
||||
const workflow = await createWorkflowWithHistory({}, owner);
|
||||
|
||||
await workflowService.activateWorkflow(owner, workflow.id);
|
||||
|
||||
const updated = await workflowRepository.findOne({ where: { id: workflow.id } });
|
||||
expect(updated?.active).toBe(true);
|
||||
expect(updated?.activeVersionId).toBe(workflow.versionId);
|
||||
|
||||
const outboxRecord = await outboxRepository.findOne({
|
||||
where: { workflowId: workflow.id },
|
||||
});
|
||||
expect(outboxRecord?.publishedVersionId).toBe(workflow.versionId);
|
||||
expect(outboxRecord?.status).toBe('pending');
|
||||
|
||||
// Publication is async: the published version is advanced by the
|
||||
// outbox consumer, not synchronously by the service.
|
||||
const publishedVersion = await workflowPublishedVersionRepository.findOne({
|
||||
where: { workflowId: workflow.id },
|
||||
});
|
||||
expect(publishedVersion?.publishedVersionId).toBe(workflow.versionId);
|
||||
expect(publishedVersion).toBeNull();
|
||||
});
|
||||
|
||||
test('should update workflow_published_version when activating a new version', async () => {
|
||||
test('should supersede the pending outbox record when activating a new version', async () => {
|
||||
const owner = await createOwner();
|
||||
const workflow = await createWorkflowWithHistory({}, owner);
|
||||
|
||||
@@ -550,10 +567,13 @@ describe('workflow_published_version table population', () => {
|
||||
versionId: newVersionId,
|
||||
});
|
||||
|
||||
const publishedVersion = await workflowPublishedVersionRepository.findOne({
|
||||
where: { workflowId: workflow.id },
|
||||
});
|
||||
expect(publishedVersion?.publishedVersionId).toBe(newVersionId);
|
||||
const outboxRecords = await outboxRepository.find({ where: { workflowId: workflow.id } });
|
||||
expect(outboxRecords).toHaveLength(1);
|
||||
expect(outboxRecords[0].publishedVersionId).toBe(newVersionId);
|
||||
expect(outboxRecords[0].status).toBe('pending');
|
||||
|
||||
const updated = await workflowRepository.findOne({ where: { id: workflow.id } });
|
||||
expect(updated?.activeVersionId).toBe(newVersionId);
|
||||
});
|
||||
|
||||
test('should remove workflow_published_version on deactivation', async () => {
|
||||
@@ -561,6 +581,8 @@ describe('workflow_published_version table population', () => {
|
||||
const workflow = await createWorkflowWithHistory({}, owner);
|
||||
|
||||
await workflowService.activateWorkflow(owner, workflow.id);
|
||||
// Simulate the outbox consumer having advanced the published version.
|
||||
await workflowPublishedVersionRepository.setPublishedVersion(workflow.id, workflow.versionId);
|
||||
|
||||
await workflowService.deactivateWorkflow(owner, workflow.id);
|
||||
|
||||
@@ -575,11 +597,8 @@ describe('workflow_published_version table population', () => {
|
||||
const workflow = await createWorkflowWithHistory({}, owner);
|
||||
|
||||
await workflowService.activateWorkflow(owner, workflow.id);
|
||||
|
||||
const publishedVersionBefore = await workflowPublishedVersionRepository.findOne({
|
||||
where: { workflowId: workflow.id },
|
||||
});
|
||||
expect(publishedVersionBefore).not.toBeNull();
|
||||
// Simulate the outbox consumer having advanced the published version.
|
||||
await workflowPublishedVersionRepository.setPublishedVersion(workflow.id, workflow.versionId);
|
||||
|
||||
await workflowService.archive(owner, workflow.id);
|
||||
|
||||
@@ -591,12 +610,15 @@ describe('workflow_published_version table population', () => {
|
||||
});
|
||||
|
||||
describe('when feature flag is disabled', () => {
|
||||
test('should not write to workflow_published_version on activation', async () => {
|
||||
test('should not enqueue an outbox record on activation', async () => {
|
||||
const owner = await createOwner();
|
||||
const workflow = await createWorkflowWithHistory({}, owner);
|
||||
|
||||
await workflowService.activateWorkflow(owner, workflow.id);
|
||||
|
||||
const outboxCount = await outboxRepository.count();
|
||||
expect(outboxCount).toBe(0);
|
||||
|
||||
const publishedVersion = await workflowPublishedVersionRepository.findOne({
|
||||
where: { workflowId: workflow.id },
|
||||
});
|
||||
|
||||
@@ -5,6 +5,9 @@ test.use({
|
||||
env: {
|
||||
TEST_ISOLATION: 'workflow-publication-service',
|
||||
N8N_USE_WORKFLOW_PUBLICATION_SERVICE: 'true',
|
||||
// Activation is applied asynchronously by the publication outbox
|
||||
// consumer, so poll frequently to keep the test fast.
|
||||
N8N_WORKFLOW_PUBLICATION_OUTBOX_POLL_INTERVAL_MS: '250',
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user