fix(editor): abort 时保留常驻依赖收集 worker,避免反复重建

clearIdleTasks/reset 频繁触发时,销毁 worker 会导致加载中被反复打断、依赖收集无法完成;改为排队投递并丢弃未投递/在途请求,复用常驻 worker。
This commit is contained in:
roymondchen 2026-08-06 11:40:33 +08:00
parent 192c16a7be
commit efb637e72e
3 changed files with 268 additions and 39 deletions

View File

@ -12,6 +12,14 @@ export interface CollectWorkerResult {
nodeIds: Id[];
}
interface CollectTask {
id: number;
createRequest: (id: number) => CollectWorkerRequest;
resolve: (response: CollectWorkerResponse | null) => void;
/** 已被 abort 丢弃:调用方已按 null 结算,结果回来后不再回传 */
discarded: boolean;
}
/**
* worker 线
*
@ -20,13 +28,20 @@ export interface CollectWorkerResult {
*
* worker
* worker 线
*
* 线worker abort
* 线worker
*/
export class CollectWorkerClient {
private worker: Worker | null = null;
private seed = 0;
private pending = new Map<number, (response: CollectWorkerResponse | null) => void>();
/** 等待投递给 worker 的请求 */
private queue: CollectTask[] = [];
/** 已投递、等待 worker 响应的请求worker 串行处理,同一时刻最多一个 */
private inflight: CollectTask | null = null;
public get isSupported() {
return typeof Worker !== 'undefined';
@ -53,24 +68,37 @@ export class CollectWorkerClient {
public terminate() {
this.worker?.terminate();
this.worker = null;
this.settleAllPending();
this.settleAll();
}
/**
* worker
* clearIdleTasks / reset abort worker
* worker
*
* clearIdleTasks / reset workerworker
* clearIdleTasks / reset root abort
* worker
* worker
*
* abort
* 线worker
*/
public abort() {
this.terminate();
const { queue } = this;
this.queue = [];
for (const task of queue) {
task.resolve(null);
}
// 已投递的请求无法从外部打断worker 忙时收不到取消消息),只能丢弃其结果
if (this.inflight && !this.inflight.discarded) {
this.inflight.discarded = true;
this.inflight.resolve(null);
}
}
/**
* worker id
*/
private request(createRequest: (id: number) => CollectWorkerRequest): Promise<CollectWorkerResponse | null> {
const worker = this.getWorker();
if (!worker) {
if (!this.isSupported) {
return Promise.resolve(null);
}
@ -78,19 +106,44 @@ export class CollectWorkerClient {
const id = this.seed;
return new Promise<CollectWorkerResponse | null>((resolve) => {
this.pending.set(id, resolve);
try {
// 节点配置中可能存在函数等无法结构化克隆的值,需要先序列化再传递
worker.postMessage(createRequest(id));
} catch (e) {
error('magic editor: 依赖收集 worker 通信失败', e);
this.pending.delete(id);
resolve(null);
}
this.queue.push({ id, createRequest, resolve, discarded: false });
this.flush();
});
}
/**
* worker id
*/
private flush() {
if (this.inflight || !this.queue.length) {
return;
}
const worker = this.getWorker();
if (!worker) {
this.settleAll();
return;
}
const task = this.queue.shift()!;
this.inflight = task;
try {
// 节点配置中可能存在函数等无法结构化克隆的值,需要先序列化再传递
worker.postMessage(task.createRequest(task.id));
} catch (e) {
error('magic editor: 依赖收集 worker 通信失败', e);
this.inflight = null;
if (!task.discarded) {
task.resolve(null);
}
this.flush();
}
}
private getWorker() {
if (!this.isSupported) {
return null;
@ -122,28 +175,41 @@ export class CollectWorkerClient {
}
private handleResponse(response: CollectWorkerResponse) {
const resolve = this.pending.get(response?.id);
const task = this.inflight;
if (!resolve) {
if (task?.id !== response?.id) {
return;
}
this.pending.delete(response.id);
resolve(response.failed ? null : response);
this.inflight = null;
// 被 abort 丢弃的请求已按 null 结算,结果不能再写回主线程
if (!task.discarded) {
task.resolve(response.failed ? null : response);
}
this.flush();
}
private handleFatalError() {
this.worker?.terminate();
this.worker = null;
this.settleAllPending();
this.settleAll();
}
private settleAllPending() {
const resolvers = [...this.pending.values()];
this.pending.clear();
private settleAll() {
const tasks = this.queue;
this.queue = [];
for (const resolve of resolvers) {
resolve(null);
if (this.inflight) {
tasks.push(this.inflight);
this.inflight = null;
}
for (const task of tasks) {
if (!task.discarded) {
task.resolve(null);
}
}
}
}

View File

@ -545,21 +545,101 @@ describe('Dep service', () => {
expect(depService.get('collecting')).toBe(false);
});
test('clearIdleTasks 会 abort 常驻 worker避免在途任务继续占用', async () => {
test('clearIdleTasks 会丢弃在途 worker 结果,但不会销毁重建常驻 worker', async () => {
const fakeCollectWorker = await enableCollectWorker();
fakeCollectWorker.delay = 50;
depService.addTarget(createDataSourceTarget({ id: 'ds_1', fields: [] }, reactive({})));
const promise = depService.collectIdle([{ id: 'n1', type: 'text' }] as any, {}, true, DepTargetType.DATA_SOURCE);
expect(fakeCollectWorker.instances).toHaveLength(1);
const instanceLength = fakeCollectWorker.instances.length;
depService.clearIdleTasks();
await expect(promise).resolves.toBe(false);
// abort 后再次收集会新建 worker而不是排队在已 abort 的旧任务后面
// 中断只丢结果worker 复用,避免「加载中被销毁 → 重建 → 又被销毁」的反复加载
fakeCollectWorker.delay = 0;
await depService.collectIdle([{ id: 'n2', type: 'text' }] as any, {}, true, DepTargetType.DATA_SOURCE);
expect(fakeCollectWorker.instances.length).toBeGreaterThan(1);
expect(fakeCollectWorker.instances.length).toBe(instanceLength);
});
test('root 反复更新时复用常驻 worker不会重复创建', async () => {
const fakeCollectWorker = await enableCollectWorker();
const dsl = { id: 'app', type: 'app', items: [] } as any;
depService.clearIdleTasks();
await depService.collectByWorker(dsl);
const instanceLength = fakeCollectWorker.instances.length;
// root 更新流程clearIdleTasks + 全量收集,多次更新应始终复用同一个 worker
for (let i = 0; i < 3; i++) {
depService.clearIdleTasks();
await depService.collectByWorker(dsl);
}
expect(fakeCollectWorker.requests).toHaveLength(4);
expect(fakeCollectWorker.instances.length).toBe(instanceLength);
});
test('大页面 + 连续 root/数据源变更:不反复创建 worker最终收集可完成', async () => {
const fakeCollectWorker = await enableCollectWorker();
// 模拟大页面worker 单次收集较慢,后续变更会打在「上一次还没结束」的窗口
fakeCollectWorker.nextDelay = 80;
fakeCollectWorker.delay = 80;
fakeCollectWorker.nextData = {
[DepTargetType.DATA_SOURCE]: { ds_1: { n1: { name: 'n1', keys: ['text'] } } },
};
fakeCollectWorker.response = {
deps: { [DepTargetType.DATA_SOURCE]: { ds_1: { n1: { name: 'n1', keys: ['text'] } } } },
nodeIds: ['n1'],
};
const dsl = {
id: 'app',
type: 'app',
items: Array.from({ length: 200 }, (_, index) => ({ id: `n${index}`, type: 'text', text: '${ds_1.name}' })),
dataSources: [{ id: 'ds_1', fields: [{ name: 'name', type: 'string' }] }],
dataSourceDeps: {},
} as any;
depService.addTarget(
createDataSourceTarget({ id: 'ds_1', fields: [{ name: 'name', type: 'string' }] }, reactive({})),
);
// 先暖机一次depService 单例可能复用上个用例的 worker以当前 instances 为基线
await depService.collectByWorker(dsl);
const instanceLength = fakeCollectWorker.instances.length;
const requestLength = fakeCollectWorker.requests.length;
const settled: boolean[] = [];
// 首轮 root 全量收集尚未完成时,交错触发 root 更新与数据源重收
const firstRoot = depService.collectByWorker(dsl).then(() => settled.push(true));
const overlapping: Promise<unknown>[] = [];
for (let i = 0; i < 8; i++) {
depService.clearIdleTasks();
overlapping.push(depService.collectByWorker(dsl).then(() => settled.push(true)));
depService.clearIdleTasks();
overlapping.push(
depService
.collectIdle(
[{ id: `n${i}`, type: 'text', text: '${ds_1.name}' }] as any,
{},
true,
DepTargetType.DATA_SOURCE,
)
.then((completed) => settled.push(completed)),
);
}
await Promise.all([firstRoot, ...overlapping]);
// 连续变更不应再 new Workerinstances 不增长);请求有实际投递
expect(fakeCollectWorker.instances.length).toBe(instanceLength);
expect(fakeCollectWorker.requests.length).toBeGreaterThan(requestLength);
// 最终至少有一轮完整收集完成;中间被 clearIdleTasks 打断的 idle 可能是 false
expect(settled.some(Boolean)).toBe(true);
expect(depService.get('collecting')).toBe(false);
});
test('collectIdle 节点为空时不会启动 worker', async () => {

View File

@ -23,6 +23,8 @@ vi.mock('@editor/utils/dep/worker.ts?worker&inline', () => ({
public static throwOnCreate = false;
/** postMessage 抛错,如结构化克隆失败 */
public static throwOnPost = false;
/** 模拟大页面收集耗时 */
public static delay = 0;
public onmessage: ((e: any) => void) | null = null;
public onerror: (() => void) | null = null;
@ -52,7 +54,7 @@ vi.mock('@editor/utils/dep/worker.ts?worker&inline', () => ({
this.onmessage?.({
data: { id: request.id, failed: FakeCollectWorker.failed, ...FakeCollectWorker.response },
});
});
}, FakeCollectWorker.delay);
}
public terminate() {
@ -81,6 +83,7 @@ beforeEach(async () => {
fakeCollectWorker.fatal = false;
fakeCollectWorker.throwOnCreate = false;
fakeCollectWorker.throwOnPost = false;
fakeCollectWorker.delay = 0;
(globalThis as any).Worker = class {};
});
@ -241,17 +244,97 @@ describe('dep/collect-worker-client', () => {
expect(fakeCollectWorker.instances[0].terminated).toBe(true);
});
test('abort 与 terminate 一样丢弃在途请求,后续收集会重建 worker', async () => {
test('abort 丢弃在途请求,但保留常驻 worker不重建、不重复加载', async () => {
const fakeCollectWorker = await getFakeWorker();
const client = new CollectWorkerClient();
const promise = client.collect(payload);
client.abort();
// 在途请求按 null 结算,调用方回退主线程收集
await expect(promise).resolves.toBeNull();
expect(fakeCollectWorker.instances[0].terminated).toBe(true);
expect(fakeCollectWorker.instances[0].terminated).toBe(false);
// worker 复用,不会因为 abort 重建
await expect(client.collect(payload)).resolves.not.toBeNull();
expect(fakeCollectWorker.instances).toHaveLength(2);
expect(fakeCollectWorker.instances).toHaveLength(1);
});
test('被 abort 丢弃的请求结果不会回传,且不阻塞后续请求', async () => {
const fakeCollectWorker = await getFakeWorker();
const client = new CollectWorkerClient();
const settled: unknown[] = [];
const aborted = client.collect(payload).then((result) => settled.push(result));
client.abort();
// abort 后紧接着的新请求(如 root 更新触发的全量收集)排在被丢弃请求之后,仍能拿到结果
const next = client.collectDsl(dsl);
await Promise.all([aborted, next]);
expect(settled).toEqual([null]);
await expect(next).resolves.not.toBeNull();
expect(fakeCollectWorker.instances).toHaveLength(1);
});
test('无在途请求时 abort 保留常驻 worker不会反复加载 worker 产物', async () => {
const fakeCollectWorker = await getFakeWorker();
const client = new CollectWorkerClient();
await client.collectDsl(dsl);
expect(fakeCollectWorker.instances).toHaveLength(1);
// 模拟 root 更新clearIdleTasks 先 abort紧接着发起全量收集
client.abort();
await client.collectDsl(dsl);
expect(fakeCollectWorker.instances).toHaveLength(1);
expect(fakeCollectWorker.instances[0].terminated).toBe(false);
});
test('大页面长耗时 + 连续 abort/收集:不重建 worker最终请求仍能完成', async () => {
const fakeCollectWorker = await getFakeWorker();
// 模拟大页面收集耗时;连续 abort 落在「上一次还没跑完」的窗口内
fakeCollectWorker.delay = 80;
const client = new CollectWorkerClient();
const results: Array<unknown> = [];
// 首轮投递后,在 worker 忙时连续模拟 root 更新abort + 新全量收集
const first = client.collectDsl(dsl).then((result) => results.push(result));
const overlapping: Promise<unknown>[] = [];
for (let i = 0; i < 12; i++) {
client.abort();
overlapping.push(client.collectDsl(dsl).then((result) => results.push(result)));
}
const all = await Promise.all([first, ...overlapping]);
// 不应陷入「加载中被销毁 → 重建」:始终只有 1 个 worker 实例
expect(fakeCollectWorker.instances).toHaveLength(1);
expect(fakeCollectWorker.instances[0].terminated).toBe(false);
// 被 abort 的立即 null最后一次保留的请求能拿到结果
expect(results.filter((item) => item === null).length).toBeGreaterThan(0);
expect(all[all.length - 1]).not.toBeNull();
// discarded 长任务仍占 worker真正投递给 worker 的次数远少于发起次数
expect(fakeCollectWorker.requests.length).toBeLessThan(all.length);
expect(fakeCollectWorker.requests.length).toBeGreaterThanOrEqual(1);
});
test('abort 后新请求要等 discarded 长任务结束才能投递', async () => {
const fakeCollectWorker = await getFakeWorker();
fakeCollectWorker.delay = 60;
const client = new CollectWorkerClient();
const first = client.collectDsl(dsl);
client.abort();
const next = client.collectDsl(dsl);
// abort 后立刻结算旧请求,但新请求尚未投递(仍被 discarded inflight 堵住)
await expect(first).resolves.toBeNull();
expect(fakeCollectWorker.requests).toHaveLength(1);
await expect(next).resolves.not.toBeNull();
expect(fakeCollectWorker.requests).toHaveLength(2);
expect(fakeCollectWorker.instances).toHaveLength(1);
});
});