fix(plugin-workflow): fix duplicated executing (#7533)

This commit is contained in:
Junyi
2025-09-29 19:50:18 +08:00
committed by GitHub
parent fd4f7fee1c
commit 4102ffb5c1
@@ -41,7 +41,7 @@ import WorkflowRepository from './repositories/WorkflowRepository';
type ID = number | string;
type Pending = [ExecutionModel, JobModel?];
type Pending = { execution: ExecutionModel; job?: JobModel; force?: boolean };
export type EventOptions = {
eventKey?: string;
@@ -77,7 +77,7 @@ export default class PluginWorkflowServer extends Plugin {
private onQueueExecution: QueueEventOptions['process'] = async (event) => {
const ExecutionRepo = this.db.getRepository('executions');
const execution = await ExecutionRepo.findOne({
const execution: ExecutionModel = await ExecutionRepo.findOne({
filterByTk: event.executionId,
});
if (!execution || execution.status !== EXECUTION_STATUS.QUEUEING) {
@@ -89,7 +89,7 @@ export default class PluginWorkflowServer extends Plugin {
this.getLogger(execution.workflowId).info(
`execution (${execution.id}) received from queue, adding to pending list`,
);
this.run(execution);
this.run({ execution });
};
private onBeforeSave = async (instance: WorkflowModel, { transaction, cycling }) => {
@@ -529,16 +529,8 @@ export default class PluginWorkflowServer extends Plugin {
return null;
}
async run(execution: ExecutionModel, job?: JobModel): Promise<void> {
while (this.executing) {
await this.executing;
}
this.executing = this.process(execution, job);
await this.executing;
this.executing = null;
async run(pending: Pending): Promise<void> {
this.pending.push(pending);
this.dispatch();
}
@@ -552,7 +544,7 @@ export default class PluginWorkflowServer extends Plugin {
`execution (${execution.id}) resuming from job (${job.id}) added to pending list`,
);
this.run(execution, job);
this.run({ execution, job, force: true });
}
/**
@@ -565,7 +557,7 @@ export default class PluginWorkflowServer extends Plugin {
}
this.getLogger(execution.workflowId).info(`starting deferred execution (${execution.id})`);
this.run(execution);
this.run({ execution, force: true });
}
private async validateEvent(workflow: WorkflowModel, context: any, options: EventOptions) {
@@ -680,7 +672,7 @@ export default class PluginWorkflowServer extends Plugin {
if (execution?.status === EXECUTION_STATUS.QUEUEING) {
if (!this.executing && !this.pending.length) {
logger.info(`local pending list is empty, adding execution (${execution.id}) to pending list`);
this.pending.push([execution]);
this.pending.push({ execution });
} else {
logger.info(`local pending list is not empty, sending execution (${execution.id}) to queue`);
if (this.ready) {
@@ -727,46 +719,20 @@ export default class PluginWorkflowServer extends Plugin {
}
this.executing = (async () => {
let next: Pending | null = null;
let next: [ExecutionModel, JobModel?] | null = null;
let execution: ExecutionModel | null = null;
// resuming has high priority
if (this.pending.length) {
next = this.pending.shift() as Pending;
this.getLogger(next[0].workflowId).info(`pending execution (${next[0].id}) ready to process`);
const pending = this.pending.shift() as Pending;
execution = pending.force ? pending.execution : await this.acquirePendingExecution(pending.execution);
if (execution) {
next = [execution, pending.job];
this.getLogger(next[0].workflowId).info(`pending execution (${next[0].id}) ready to process`);
}
} else {
try {
await this.db.sequelize.transaction(
{
isolationLevel:
this.db.options.dialect === 'sqlite' ? [][0] : Transaction.ISOLATION_LEVELS.REPEATABLE_READ,
},
async (transaction) => {
const execution = (await this.db.getRepository('executions').findOne({
filter: {
status: {
[Op.is]: EXECUTION_STATUS.QUEUEING,
},
'workflow.enabled': true,
},
sort: 'id',
transaction,
})) as ExecutionModel;
if (execution) {
this.getLogger(execution.workflowId).info(`execution (${execution.id}) fetched from db`);
await execution.update(
{
status: EXECUTION_STATUS.STARTED,
},
{ transaction },
);
execution.workflow = this.enabledCache.get(execution.workflowId);
next = [execution];
} else {
this.getLogger('dispatcher').debug(`no execution in db queued to process`);
}
},
);
} catch (error) {
this.getLogger('dispatcher').error(`fetching execution from db failed: ${error.message}`, { error });
execution = await this.acquireQueueingExecution();
if (execution) {
next = [execution];
}
}
if (next) {
@@ -781,6 +747,77 @@ export default class PluginWorkflowServer extends Plugin {
})();
}
private async acquirePendingExecution(execution: ExecutionModel): Promise<ExecutionModel | null> {
const logger = this.getLogger(execution.workflowId);
const isolationLevel = this.db.options.dialect === 'sqlite' ? [][0] : Transaction.ISOLATION_LEVELS.REPEATABLE_READ;
let fetched = execution;
try {
await this.db.sequelize.transaction({ isolationLevel }, async (transaction) => {
const ExecutionModelClass = this.db.getModel('executions');
const [affected] = await ExecutionModelClass.update(
{ status: EXECUTION_STATUS.STARTED },
{
where: {
id: execution.id,
status: {
[Op.is]: EXECUTION_STATUS.QUEUEING,
},
},
transaction,
},
);
if (!affected) {
fetched = null;
return;
}
await execution.reload({ transaction });
});
} catch (error) {
logger.error(`acquiring pending execution failed: ${error.message}`, { error });
}
return fetched;
}
private async acquireQueueingExecution(): Promise<ExecutionModel | null> {
const isolationLevel = this.db.options.dialect === 'sqlite' ? [][0] : Transaction.ISOLATION_LEVELS.REPEATABLE_READ;
let fetched: ExecutionModel | null = null;
try {
await this.db.sequelize.transaction(
{
isolationLevel,
},
async (transaction) => {
const execution = (await this.db.getRepository('executions').findOne({
filter: {
status: {
[Op.is]: EXECUTION_STATUS.QUEUEING,
},
'workflow.enabled': true,
},
sort: 'id',
transaction,
})) as ExecutionModel;
if (execution) {
this.getLogger(execution.workflowId).info(`execution (${execution.id}) fetched from db`);
await execution.update(
{
status: EXECUTION_STATUS.STARTED,
},
{ transaction },
);
execution.workflow = this.enabledCache.get(execution.workflowId);
fetched = execution;
} else {
this.getLogger('dispatcher').debug(`no execution in db queued to process`);
}
},
);
} catch (error) {
this.getLogger('dispatcher').error(`fetching execution from db failed: ${error.message}`, { error });
}
return fetched;
}
public createProcessor(execution: ExecutionModel, options = {}): Processor {
return new Processor(execution, { ...options, plugin: this });
}