feat(workflow): soft-ref core workflow/version (backfill + dual-write) (#22821)
Replaces the shared-UUID model (core row reuses the workspace record id) with a **soft-ref**: the workspace `workflow`/`workflowVersion` records carry a nullable `coreWorkflowId`/`coreWorkflowVersionId` pointing to their **own-id** core rows. This removes the assumption that workspace record ids are globally unique - which is false, since prefilled/seeded workflows share ids across workspaces. Supersedes #22776. ## In this PR **Soft-ref columns (foundation):** - **twenty-shared** `STANDARD_OBJECTS`: `workflowVersion.coreWorkflowVersionId` + `workflow.coreWorkflowId` (+ snapshot test). - **compute utils**: both as system, nullable UUID fields. - **entity classes**: the bare fields. **Version soft-ref sync:** - Core `workflowVersion` rows get their own id, derived deterministically from `workspaceId + record id` (uuidv5). Deterministic so the upsert is idempotent: a failed write-back re-derives the same id and self-heals instead of orphaning rows or colliding on the one-active-per-workflow index. - Sync = find-or-create keyed on the workspace record's `coreWorkflowVersionId`, then write the core id back onto the workspace record. - Migrating over pre-soft-ref data: purges any core row whose id equals the workspace record id before recreating, so old shared-UUID rows aren't orphaned. - Version dual-write listener reworked: delete is keyed by the core id read off `before.coreWorkflowVersionId`. Verified on a fresh `database:reset` (columns materialize, backfill produces deterministic own-id rows linked back, idempotent re-run), a simulated old shared-UUID state (stale rows purged, records re-linked), and a simulated write-back failure (retry re-links to the same id, no orphan, active-version index intact). ## Next steps (follow-up work, not in this PR) 1. Workflow-side soft-ref sync mirroring the version side (service, module, dual-write listener, backfill command). 2. Workspace command to add the two columns to existing workspaces. <!-- This is an auto-generated description by cubic. --> <a href="https://cubic.dev/pr/twentyhq/twenty/pull/22821?utm_source=github" target="_blank" rel="noopener noreferrer" data-no-image-dialog="true"><picture><source media="(prefers-color-scheme: dark)" srcset="https://www.cubic.dev/buttons/review-in-cubic-dark.svg"><source media="(prefers-color-scheme: light)" srcset="https://www.cubic.dev/buttons/review-in-cubic-light.svg"><img alt="Review in cubic" src="https://www.cubic.dev/buttons/review-in-cubic-dark.svg"></picture></a> <!-- End of auto-generated description by cubic. -->
This commit is contained in:
+100
-22
@@ -3,18 +3,26 @@ import { InjectRepository } from '@nestjs/typeorm';
|
||||
|
||||
import { isDefined } from 'twenty-shared/utils';
|
||||
import { In, Repository } from 'typeorm';
|
||||
import { v4 as uuidv4 } from 'uuid';
|
||||
import { v4 as uuidv4, v5 as uuidv5 } from 'uuid';
|
||||
|
||||
import {
|
||||
WorkflowVersionEntity,
|
||||
WorkflowVersionStatus,
|
||||
} from 'src/engine/core-modules/workflow/entities/workflow-version.entity';
|
||||
import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
|
||||
import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager';
|
||||
import { InjectWorkspaceScopedRepository } from 'src/engine/twenty-orm/workspace-scoped-repository/inject-workspace-scoped-repository.decorator';
|
||||
import { WorkspaceScopedRepository } from 'src/engine/twenty-orm/workspace-scoped-repository/workspace-scoped-repository';
|
||||
import { buildSystemAuthContext } from 'src/engine/twenty-orm/utils/build-system-auth-context.util';
|
||||
import { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service';
|
||||
import { type WorkflowVersionWorkspaceEntity } from 'src/modules/workflow/common/standard-objects/workflow-version.workspace-entity';
|
||||
|
||||
// Deriving the core id deterministically from workspaceId + record id makes the
|
||||
// upsert idempotent across retries: a failed write-back re-derives the same id
|
||||
// instead of orphaning a row or hitting the one-active-per-workflow index.
|
||||
const CORE_WORKFLOW_VERSION_ID_NAMESPACE =
|
||||
'f4988927-0a5c-453a-a262-0bd136d7fdaf';
|
||||
|
||||
@Injectable()
|
||||
export class WorkflowVersionCoreSyncService {
|
||||
constructor(
|
||||
@@ -22,6 +30,7 @@ export class WorkflowVersionCoreSyncService {
|
||||
private readonly workflowVersionRepository: WorkspaceScopedRepository<WorkflowVersionEntity>,
|
||||
@InjectRepository(WorkspaceEntity)
|
||||
private readonly workspaceRepository: Repository<WorkspaceEntity>,
|
||||
private readonly globalWorkspaceOrmManager: GlobalWorkspaceOrmManager,
|
||||
private readonly workspaceCacheService: WorkspaceCacheService,
|
||||
) {}
|
||||
|
||||
@@ -35,10 +44,25 @@ export class WorkflowVersionCoreSyncService {
|
||||
|
||||
const applicationId = await this.getCustomApplicationIdOrThrow(workspaceId);
|
||||
|
||||
await this.workflowVersionRepository.upsert(
|
||||
workspaceId,
|
||||
workflowVersions.map((workflowVersion) => ({
|
||||
id: workflowVersion.id,
|
||||
const coreVersionIdByWorkspaceRecordId = new Map<string, string>();
|
||||
|
||||
const coreRows = workflowVersions.map((workflowVersion) => {
|
||||
const coreWorkflowVersionId =
|
||||
workflowVersion.coreWorkflowVersionId ??
|
||||
uuidv5(
|
||||
`${workspaceId}:${workflowVersion.id}`,
|
||||
CORE_WORKFLOW_VERSION_ID_NAMESPACE,
|
||||
);
|
||||
|
||||
if (!isDefined(workflowVersion.coreWorkflowVersionId)) {
|
||||
coreVersionIdByWorkspaceRecordId.set(
|
||||
workflowVersion.id,
|
||||
coreWorkflowVersionId,
|
||||
);
|
||||
}
|
||||
|
||||
return {
|
||||
id: coreWorkflowVersionId,
|
||||
workflowId: workflowVersion.workflowId,
|
||||
triggers: isDefined(workflowVersion.trigger)
|
||||
? [workflowVersion.trigger]
|
||||
@@ -47,13 +71,82 @@ export class WorkflowVersionCoreSyncService {
|
||||
status: workflowVersion.status as unknown as WorkflowVersionStatus,
|
||||
universalIdentifier: uuidv4(),
|
||||
applicationId,
|
||||
})),
|
||||
['id'],
|
||||
};
|
||||
});
|
||||
|
||||
await this.purgeSharedIdCoreRows(
|
||||
workspaceId,
|
||||
Array.from(coreVersionIdByWorkspaceRecordId.keys()),
|
||||
);
|
||||
|
||||
await this.workflowVersionRepository.upsert(workspaceId, coreRows, ['id']);
|
||||
|
||||
await this.writeBackCoreVersionIds(
|
||||
workspaceId,
|
||||
coreVersionIdByWorkspaceRecordId,
|
||||
);
|
||||
|
||||
await this.invalidateAutomatedTriggerMaps(workspaceId);
|
||||
}
|
||||
|
||||
async deleteFromCore(
|
||||
workspaceId: string,
|
||||
coreWorkflowVersionIds: string[],
|
||||
): Promise<void> {
|
||||
if (coreWorkflowVersionIds.length === 0) {
|
||||
return;
|
||||
}
|
||||
|
||||
await this.workflowVersionRepository.delete(workspaceId, {
|
||||
id: In(coreWorkflowVersionIds),
|
||||
});
|
||||
|
||||
await this.invalidateAutomatedTriggerMaps(workspaceId);
|
||||
}
|
||||
|
||||
// The pre-soft-ref model wrote core rows with id === workspace record id.
|
||||
// Delete those before creating own-id rows, otherwise re-migration orphans
|
||||
// them and a second active row collides on the one-active-per-workflow index.
|
||||
private async purgeSharedIdCoreRows(
|
||||
workspaceId: string,
|
||||
workspaceRecordIds: string[],
|
||||
): Promise<void> {
|
||||
if (workspaceRecordIds.length === 0) {
|
||||
return;
|
||||
}
|
||||
|
||||
await this.workflowVersionRepository.delete(workspaceId, {
|
||||
id: In(workspaceRecordIds),
|
||||
});
|
||||
}
|
||||
|
||||
private async writeBackCoreVersionIds(
|
||||
workspaceId: string,
|
||||
coreVersionIdByWorkspaceRecordId: Map<string, string>,
|
||||
): Promise<void> {
|
||||
if (coreVersionIdByWorkspaceRecordId.size === 0) {
|
||||
return;
|
||||
}
|
||||
|
||||
await this.globalWorkspaceOrmManager.executeInWorkspaceContext(async () => {
|
||||
const workflowVersionRepository =
|
||||
await this.globalWorkspaceOrmManager.getRepository<WorkflowVersionWorkspaceEntity>(
|
||||
workspaceId,
|
||||
'workflowVersion',
|
||||
{ shouldBypassPermissionChecks: true },
|
||||
);
|
||||
|
||||
for (const [
|
||||
workspaceRecordId,
|
||||
coreWorkflowVersionId,
|
||||
] of coreVersionIdByWorkspaceRecordId) {
|
||||
await workflowVersionRepository.update(workspaceRecordId, {
|
||||
coreWorkflowVersionId,
|
||||
});
|
||||
}
|
||||
}, buildSystemAuthContext(workspaceId));
|
||||
}
|
||||
|
||||
private async getCustomApplicationIdOrThrow(
|
||||
workspaceId: string,
|
||||
): Promise<string> {
|
||||
@@ -71,21 +164,6 @@ export class WorkflowVersionCoreSyncService {
|
||||
return workspace.workspaceCustomApplicationId;
|
||||
}
|
||||
|
||||
async deleteFromCore(
|
||||
workspaceId: string,
|
||||
workflowVersionIds: string[],
|
||||
): Promise<void> {
|
||||
if (workflowVersionIds.length === 0) {
|
||||
return;
|
||||
}
|
||||
|
||||
await this.workflowVersionRepository.delete(workspaceId, {
|
||||
id: In(workflowVersionIds),
|
||||
});
|
||||
|
||||
await this.invalidateAutomatedTriggerMaps(workspaceId);
|
||||
}
|
||||
|
||||
private async invalidateAutomatedTriggerMaps(
|
||||
workspaceId: string,
|
||||
): Promise<void> {
|
||||
|
||||
Reference in New Issue
Block a user