diff --git a/packages/editor/src/utils/dep/collect-worker-client.ts b/packages/editor/src/utils/dep/collect-worker-client.ts index 464de3cd..843c40b9 100644 --- a/packages/editor/src/utils/dep/collect-worker-client.ts +++ b/packages/editor/src/utils/dep/collect-worker-client.ts @@ -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 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 中断收集时调用。这里不能销毁 worker:worker 重建要重新加载并启动整份收集逻辑, + * 而 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 { - 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((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); + } } } } diff --git a/packages/editor/tests/unit/services/dep.spec.ts b/packages/editor/tests/unit/services/dep.spec.ts index 3e4ded54..c6c3e2cc 100644 --- a/packages/editor/tests/unit/services/dep.spec.ts +++ b/packages/editor/tests/unit/services/dep.spec.ts @@ -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[] = []; + + 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 Worker(instances 不增长);请求有实际投递 + 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 () => { diff --git a/packages/editor/tests/unit/utils/dep/collect-worker-client.spec.ts b/packages/editor/tests/unit/utils/dep/collect-worker-client.spec.ts index b3bc2aba..cc76c091 100644 --- a/packages/editor/tests/unit/utils/dep/collect-worker-client.spec.ts +++ b/packages/editor/tests/unit/utils/dep/collect-worker-client.spec.ts @@ -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 = []; + + // 首轮投递后,在 worker 忙时连续模拟 root 更新:abort + 新全量收集 + const first = client.collectDsl(dsl).then((result) => results.push(result)); + const overlapping: Promise[] = []; + 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); }); });