fix(editor): 修复 collectIdle 批次结算与中断时 Promise 挂起问题

引入批次级 Promise 结算,避免快速连续触发或 clearIdleTasks 后 collecting 卡死。
collectIdle 返回 boolean 表示是否完整完成,initService 据此跳过被中断的 stage 更新。
This commit is contained in:
roymondchen 2026-07-24 17:17:55 +08:00
parent bf8df74864
commit b5abf31066
5 changed files with 253 additions and 43 deletions

View File

@ -333,11 +333,11 @@ export const initServiceEvents = (
Promise.all( Promise.all(
nodes.map((node) => { nodes.map((node) => {
if (node.type === NodeType.ROOT) { if (node.type === NodeType.ROOT) {
return Promise.resolve(); return Promise.resolve(true);
} }
return depService.collectIdle([node], { pageId: getPageIdByNode(node) }, deep, type); return depService.collectIdle([node], { pageId: getPageIdByNode(node) }, deep, type);
}), }),
); ).then((results) => results.every(Boolean));
watch( watch(
() => editorService.get('stage'), () => editorService.get('stage'),
@ -351,7 +351,7 @@ export const initServiceEvents = (
if (!node) return; if (!node) return;
await collectIdle([node], true, DepTargetType.DATA_SOURCE); if (!(await collectIdle([node], true, DepTargetType.DATA_SOURCE))) return;
updateStageNode(node); updateStageNode(node);
}); });
}, },
@ -494,8 +494,10 @@ export const initServiceEvents = (
// 新增节点,收集依赖 // 新增节点,收集依赖
const nodeAddHandler = (nodes: MComponent[]) => { const nodeAddHandler = (nodes: MComponent[]) => {
collectIdle(nodes, true).then(() => { collectIdle(nodes, true).then((completed) => {
updateStageNodes(nodes); if (completed) {
updateStageNodes(nodes);
}
}); });
}; };
@ -549,8 +551,8 @@ export const initServiceEvents = (
if (needRecollectNodes.length) { if (needRecollectNodes.length) {
// 有数据源依赖需要等依赖重新收集完才更新stage // 有数据源依赖需要等依赖重新收集完才更新stage
const handler = async () => { const handler = async () => {
await collectIdle(needRecollectNodes, true, DepTargetType.DATA_SOURCE); if (!(await collectIdle(needRecollectNodes, true, DepTargetType.DATA_SOURCE))) return;
await collectIdle(needRecollectNodes, true, DepTargetType.DATA_SOURCE_COND); if (!(await collectIdle(needRecollectNodes, true, DepTargetType.DATA_SOURCE_COND))) return;
updateStageNodes(needRecollectNodes); updateStageNodes(needRecollectNodes);
}; };
handler(); handler();
@ -572,8 +574,10 @@ export const initServiceEvents = (
// 历史记录变化时,需要重新收集依赖 // 历史记录变化时,需要重新收集依赖
const historyChangeHandler = (page: MPage | MPageFragment) => { const historyChangeHandler = (page: MPage | MPageFragment) => {
collectIdle([page], true).then(() => { collectIdle([page], true).then((completed) => {
updateStageNode(page); if (completed) {
updateStageNode(page);
}
}); });
}; };
@ -654,7 +658,7 @@ export const initServiceEvents = (
if (Array.isArray(root?.items)) { if (Array.isArray(root?.items)) {
depService.clearIdleTasks(); depService.clearIdleTasks();
let collectIdlePromises: Promise<void[]>[] = []; let collectIdlePromises: Promise<boolean>[] = [];
if (isModifyField) { if (isModifyField) {
depService.removeTarget(config.id, DepTargetType.DATA_SOURCE); depService.removeTarget(config.id, DepTargetType.DATA_SOURCE);
depService.removeTarget(config.id, DepTargetType.DATA_SOURCE_COND); depService.removeTarget(config.id, DepTargetType.DATA_SOURCE_COND);
@ -679,10 +683,15 @@ export const initServiceEvents = (
collectIdlePromises = [collectIdle(root.items, true, DepTargetType.DATA_SOURCE_METHOD)]; collectIdlePromises = [collectIdle(root.items, true, DepTargetType.DATA_SOURCE_METHOD)];
} }
Promise.all(collectIdlePromises) const handler = async () => {
.then(() => updateDataSourceSchema()) const results = await Promise.all(collectIdlePromises);
.then(() => updateDsData()) if (!results.every(Boolean)) return;
.then(() => updateStageNodes(root.items));
updateDataSourceSchema();
await updateDsData();
updateStageNodes(root.items);
};
handler();
} }
} else if (root?.dataSources) { } else if (root?.dataSources) {
updateDsData(); updateDsData();
@ -706,11 +715,12 @@ export const initServiceEvents = (
const nodeIds = Object.keys(root.dataSourceDeps?.[id] || {}); const nodeIds = Object.keys(root.dataSourceDeps?.[id] || {});
const nodes = getNodes(nodeIds, root.items); const nodes = getNodes(nodeIds, root.items);
await Promise.all([ const results = await Promise.all([
collectIdle(nodes, false, DepTargetType.DATA_SOURCE), collectIdle(nodes, false, DepTargetType.DATA_SOURCE),
collectIdle(nodes, false, DepTargetType.DATA_SOURCE_COND), collectIdle(nodes, false, DepTargetType.DATA_SOURCE_COND),
collectIdle(nodes, false, DepTargetType.DATA_SOURCE_METHOD), collectIdle(nodes, false, DepTargetType.DATA_SOURCE_METHOD),
]); ]);
if (!results.every(Boolean)) return;
updateDataSourceSchema(); updateDataSourceSchema();

View File

@ -42,6 +42,21 @@ interface State {
type StateKey = keyof State; type StateKey = keyof State;
/**
* collectIdle /
* Promise idleTask finish resolve
* clearTasks Promise resolvecollecting
*/
interface CollectBatch {
nodes: MNode[];
deep: boolean;
pending: number;
dsPending: number;
collectedEmitted: boolean;
dsSettled: boolean;
resolve: (completed: boolean) => void;
}
class Dep extends BaseService { class Dep extends BaseService {
private state = shallowReactive<State>({ private state = shallowReactive<State>({
collecting: false, collecting: false,
@ -54,6 +69,12 @@ class Dep extends BaseService {
private waitingWorker?: Promise<void>; private waitingWorker?: Promise<void>;
private resolveWaitingWorker?: () => void;
private workerGeneration = 0;
private activeBatches = new Set<CollectBatch>();
constructor() { constructor() {
super(); super();
@ -137,38 +158,44 @@ class Dep extends BaseService {
await this.waitingWorker; await this.waitingWorker;
} }
this.set('collecting', true); const batch: CollectBatch = {
let startTask = false; nodes,
this.watcher.collectByCallback(nodes, type, ({ node, target }) => { deep,
startTask = true; pending: 0,
dsPending: 0,
collectedEmitted: false,
dsSettled: false,
resolve: () => {},
};
this.enqueueTask(node, target, depExtendedData, deep); this.watcher.collectByCallback(nodes, type, ({ node, target }) => {
this.enqueueTask(node, target, depExtendedData, deep, batch);
}); });
return new Promise<void>((resolve) => { // 没有命中任何 target无需收集直接完成
if (!startTask) { if (batch.pending === 0) {
this.emit('collected', nodes, deep); this.emit('collected', nodes, deep);
this.set('collecting', false); this.updateCollectingState();
resolve(); return true;
return; }
}
this.idleTask.once('finish', () => { this.activeBatches.add(batch);
this.emit('collected', nodes, deep); this.set('collecting', true);
this.set('collecting', false);
}); return new Promise<boolean>((resolve) => {
this.idleTask.once('hight-level-finish', () => { batch.resolve = resolve;
this.emit('ds-collected', nodes, deep);
resolve();
});
}); });
} }
public collectByWorker(dsl: MApp) { public collectByWorker(dsl: MApp) {
this.set('collecting', true); this.set('collecting', true);
this.workerGeneration += 1;
const generation = this.workerGeneration;
const { promise, resolve: waitingResolve } = Promise.withResolvers<void>(); const { promise, resolve: waitingResolve } = Promise.withResolvers<void>();
this.waitingWorker = promise; this.waitingWorker = promise;
this.resolveWaitingWorker = waitingResolve;
return new Promise<Record<string, Record<string, DepData>>>((resolve) => { return new Promise<Record<string, Record<string, DepData>>>((resolve) => {
const worker = new Work(); const worker = new Work();
@ -182,6 +209,11 @@ class Dep extends BaseService {
resolve({}); resolve({});
}; };
}).then((depsData) => { }).then((depsData) => {
if (generation !== this.workerGeneration) {
waitingResolve();
return depsData;
}
traverseTarget(this.watcher.getTargetsList(), (target) => { traverseTarget(this.watcher.getTargetsList(), (target) => {
if (depsData[target.type]?.[target.id]) { if (depsData[target.type]?.[target.id]) {
target.deps = reactive(depsData[target.type][target.id]); target.deps = reactive(depsData[target.type][target.id]);
@ -200,6 +232,10 @@ class Dep extends BaseService {
this.emit('collected', dsl.items, true); this.emit('collected', dsl.items, true);
this.emit('ds-collected', dsl.items, true); this.emit('ds-collected', dsl.items, true);
waitingResolve(); waitingResolve();
if (this.waitingWorker === promise) {
this.waitingWorker = undefined;
this.resolveWaitingWorker = undefined;
}
return depsData; return depsData;
}); });
@ -233,6 +269,7 @@ class Dep extends BaseService {
} }
public clearIdleTasks() { public clearIdleTasks() {
this.abortActiveBatches();
this.idleTask.clearTasks(); this.idleTask.clearTasks();
} }
@ -251,6 +288,11 @@ class Dep extends BaseService {
} }
public reset() { public reset() {
this.abortActiveBatches();
this.workerGeneration += 1;
this.resolveWaitingWorker?.();
this.resolveWaitingWorker = undefined;
this.waitingWorker = undefined;
this.idleTask.clearTasks(); this.idleTask.clearTasks();
for (const type of Object.keys(this.watcher.getTargetsList())) { for (const type of Object.keys(this.watcher.getTargetsList())) {
@ -286,25 +328,99 @@ class Dep extends BaseService {
} }
} }
private enqueueTask(node: MNode, target: Target, depExtendedData: DepExtendedData, deep: boolean) { private enqueueTask(
node: MNode,
target: Target,
depExtendedData: DepExtendedData,
deep: boolean,
batch: CollectBatch,
) {
const isDataSource = target.type === DepTargetType.DATA_SOURCE;
batch.pending += 1;
if (isDataSource) {
batch.dsPending += 1;
}
this.idleTask.enqueueTask( this.idleTask.enqueueTask(
({ node, deep, target }) => { ({ node, deep, target }) => {
this.collectNode(node, target, depExtendedData, deep); try {
this.collectNode(node, target, depExtendedData, deep);
} finally {
this.onBatchTaskDone(batch, isDataSource);
}
}, },
{ {
node, node,
deep: false, deep: false,
target, target,
}, },
target.type === DepTargetType.DATA_SOURCE, isDataSource,
); );
if (deep && Array.isArray(node.items) && node.items.length) { if (deep && Array.isArray(node.items) && node.items.length) {
node.items.forEach((item) => { node.items.forEach((item) => {
this.enqueueTask(item, target, depExtendedData, deep); this.enqueueTask(item, target, depExtendedData, deep, batch);
}); });
} }
} }
private onBatchTaskDone(batch: CollectBatch, isDataSource: boolean) {
if (isDataSource) {
batch.dsPending -= 1;
// 数据源依赖收集完先 resolve让 stage 尽快更新,其余依赖继续在后台收集
if (batch.dsPending === 0) {
this.settleBatchDs(batch);
}
}
batch.pending -= 1;
if (batch.pending === 0) {
this.finishBatch(batch);
}
}
private settleBatchDs(batch: CollectBatch) {
if (batch.dsSettled) return;
batch.dsSettled = true;
this.emit('ds-collected', batch.nodes, batch.deep);
batch.resolve(true);
}
private finishBatch(batch: CollectBatch) {
if (!batch.collectedEmitted) {
batch.collectedEmitted = true;
this.emit('collected', batch.nodes, batch.deep);
}
// 没有数据源任务的批次在此结算 Promise
this.settleBatchDs(batch);
this.activeBatches.delete(batch);
this.updateCollectingState();
}
/**
* idleTask
* collectIdle Promise resolvecollecting trueonce
* emit collected/ds-collected Promise
*/
private abortActiveBatches() {
if (!this.activeBatches.size) {
return;
}
const batches = [...this.activeBatches];
this.activeBatches.clear();
for (const batch of batches) {
batch.resolve(false);
}
this.updateCollectingState();
}
private updateCollectingState() {
this.set('collecting', this.activeBatches.size > 0);
}
} }
export type DepService = Dep; export type DepService = Dep;

View File

@ -25,6 +25,12 @@ globalThis.requestIdleCallback =
}, 1); }, 1);
}; };
globalThis.cancelIdleCallback =
globalThis.cancelIdleCallback ||
function (handle) {
clearTimeout(handle as unknown as ReturnType<typeof setTimeout>);
};
export class IdleTask<T = any> extends EventEmitter { export class IdleTask<T = any> extends EventEmitter {
private taskList: TaskList<T> = []; private taskList: TaskList<T> = [];

View File

@ -91,7 +91,7 @@ const mkServices = () => {
clear: vi.fn(), clear: vi.fn(),
clearTargets: vi.fn(), clearTargets: vi.fn(),
clearIdleTasks: vi.fn(), clearIdleTasks: vi.fn(),
collectIdle: vi.fn(async () => undefined), collectIdle: vi.fn(async () => true),
collectByWorker: vi.fn(async () => undefined), collectByWorker: vi.fn(async () => undefined),
reset: vi.fn(), reset: vi.fn(),
}; };

View File

@ -13,16 +13,17 @@ vi.mock('@editor/utils/dep/worker.ts?worker&inline', () => ({
default: class FakeWorker { default: class FakeWorker {
public static nextData: Record<string, any> = {}; public static nextData: Record<string, any> = {};
public static nextError = false; public static nextError = false;
public static nextDelay = 0;
public onmessage: ((e: any) => void) | null = null; public onmessage: ((e: any) => void) | null = null;
public onerror: (() => void) | null = null; public onerror: (() => void) | null = null;
public postMessage() { public postMessage() {
setTimeout(() => { setTimeout(() => {
if (FakeWorker.nextError) { if (FakeWorker.nextError) {
this.onerror?.(new Event('error')); this.onerror?.();
return; return;
} }
this.onmessage?.({ data: FakeWorker.nextData }); this.onmessage?.({ data: FakeWorker.nextData });
}, 0); }, FakeWorker.nextDelay);
} }
}, },
})); }));
@ -121,7 +122,7 @@ describe('Dep service', () => {
test('collectIdle - 没有命中时立即 resolve 并 emit collected', async () => { test('collectIdle - 没有命中时立即 resolve 并 emit collected', async () => {
const fn = vi.fn(); const fn = vi.fn();
depService.on('collected', fn); depService.on('collected', fn);
await depService.collectIdle([{ id: 'n1', type: 'text' }] as any); await expect(depService.collectIdle([{ id: 'n1', type: 'text' }] as any)).resolves.toBe(true);
expect(fn).toHaveBeenCalled(); expect(fn).toHaveBeenCalled();
depService.off('collected', fn); depService.off('collected', fn);
}); });
@ -219,6 +220,83 @@ describe('Dep service', () => {
fakeWorker.nextData = {}; fakeWorker.nextData = {};
}); });
test('collectIdle 命中 target 时最终 resolve 并按批次 emit collected/ds-collected', async () => {
depService.addTarget(makeTarget('ds1', DepTargetType.DATA_SOURCE));
const collected = vi.fn();
const dsCollected = vi.fn();
depService.on('collected', collected);
depService.on('ds-collected', dsCollected);
const nodes = [{ id: 'n1', type: 'text' }] as any;
await expect(depService.collectIdle(nodes, {}, false, DepTargetType.DATA_SOURCE)).resolves.toBe(true);
expect(dsCollected).toHaveBeenCalledWith(nodes, false);
expect(collected).toHaveBeenCalledWith(nodes, false);
expect(depService.get('collecting')).toBe(false);
depService.off('collected', collected);
depService.off('ds-collected', dsCollected);
});
test('clearIdleTasks 会结算在途 collectIdle避免 Promise 永久挂起且 collecting 复位', async () => {
depService.addTarget(makeTarget('ds1', DepTargetType.DATA_SOURCE));
const promise = depService.collectIdle([{ id: 'n1', type: 'text' }] as any, {}, false, DepTargetType.DATA_SOURCE);
expect(depService.get('collecting')).toBe(true);
// 快速触发:任务尚未执行就清空队列,批次应被主动结算而不是永久挂起
depService.clearIdleTasks();
await expect(promise).resolves.toBe(false);
expect(depService.get('collecting')).toBe(false);
});
test('reset 会结算在途 collectIdle', async () => {
depService.addTarget(makeTarget('ds1', DepTargetType.DATA_SOURCE));
const promise = depService.collectIdle([{ id: 'n1', type: 'text' }] as any, {}, false, DepTargetType.DATA_SOURCE);
depService.reset();
await expect(promise).resolves.toBe(false);
expect(depService.get('collecting')).toBe(false);
});
test('reset 会忽略在途 worker 的过期结果,避免覆盖新依赖', async () => {
const fakeWorker = (await import('@editor/utils/dep/worker.ts?worker&inline')).default as any;
fakeWorker.nextDelay = 20;
fakeWorker.nextData = {
[DepTargetType.DATA_SOURCE]: { ds1: { n1: { data: {} } } },
};
const workerPromise = depService.collectByWorker({ items: [], id: 'app', type: 'app' } as any);
depService.reset();
const target = makeTarget('ds1', DepTargetType.DATA_SOURCE);
depService.addTarget(target);
const idlePromise = depService.collectIdle(
[{ id: 'n1', type: 'text' }] as any,
{},
false,
DepTargetType.DATA_SOURCE,
);
await Promise.all([workerPromise, idlePromise]);
expect(target.deps.n1).toBeUndefined();
fakeWorker.nextDelay = 0;
fakeWorker.nextData = {};
});
test('多个批次并发时各自独立 resolve全部完成后 collecting 复位', async () => {
depService.addTarget(makeTarget('ds1', DepTargetType.DATA_SOURCE));
const p1 = depService.collectIdle([{ id: 'n1', type: 'text' }] as any, {}, false, DepTargetType.DATA_SOURCE);
const p2 = depService.collectIdle([{ id: 'n2', type: 'text' }] as any, {}, false, DepTargetType.DATA_SOURCE);
await Promise.all([p1, p2]);
expect(depService.get('collecting')).toBe(false);
});
test('destroy 会 reset 并移除监听', () => { test('destroy 会 reset 并移除监听', () => {
depService.addTarget(makeTarget('destroy-me')); depService.addTarget(makeTarget('destroy-me'));
expect(() => depService.destroy()).not.toThrow(); expect(() => depService.destroy()).not.toThrow();