From d5145cbf17f93e185523e969e062381ef8ba3193 Mon Sep 17 00:00:00 2001 From: Hassan Abdel-Rahman Date: Fri, 24 Jul 2026 09:43:33 -0400 Subject: [PATCH 1/3] Coalesce client loader/store rebuilds during incremental event bursts MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Every incremental index event touching an already-loaded executable module tore down and rebuilt the whole client — resetLoader + store.reset + re-establishing every live card reference — once per event. A write burst therefore cost N full rebuilds and N full-graph re-fetches, arriving at or above the rate they complete. Move the rebuild into a keepLatest task so at most one runs in flight and one stays pending: intermediate events collapse into the pending slot, so N rapid executable invalidations cost at most 2 rebuilds. The final rebuild re-fetches current server state, so the end state still reflects the latest generation. Reactivity is unchanged — an isolated change still resets and re-renders exactly as before, and the full reset stays load-bearing. Once a rebuild is in flight, further executable invalidations fold into it regardless of the isModuleLoaded probe, which reads false against the freshly reset loader and would otherwise drop a mid-burst code change. The per-invalidation reload loop is skipped when a rebuild is scheduled, since the rebuild subsumes it (as the synchronous store.reset already did). The rebuild telemetry now fires once per actual rebuild and carries a coalesced_events count. Co-Authored-By: Claude Opus 4.8 (1M context) --- .../host/app/services/client-telemetry.ts | 5 + packages/host/app/services/store.ts | 208 +++++++++++------- .../host/tests/integration/store-test.gts | 136 ++++++++++++ 3 files changed, 274 insertions(+), 75 deletions(-) diff --git a/packages/host/app/services/client-telemetry.ts b/packages/host/app/services/client-telemetry.ts index 46ed689f7a..b9e7499473 100644 --- a/packages/host/app/services/client-telemetry.ts +++ b/packages/host/app/services/client-telemetry.ts @@ -105,6 +105,11 @@ export interface RebuildEvent extends BaseEvent { trigger_module: string; modules_refetched: number; cards_reloaded: number; + // How many incremental events collapsed into this single rebuild. 1 for an + // isolated change; >1 when a write burst arrived faster than a rebuild + // completes and the events coalesced into one in-flight plus one pending + // rebuild. + coalesced_events: number; } export interface RealmEvent extends BaseEvent { diff --git a/packages/host/app/services/store.ts b/packages/host/app/services/store.ts index 78801ba79a..ae70e9c52b 100644 --- a/packages/host/app/services/store.ts +++ b/packages/host/app/services/store.ts @@ -11,7 +11,7 @@ import { isTesting } from '@embroider/macros'; import { tracked } from '@glimmer/tracking'; import { formatDistanceToNow } from 'date-fns'; -import { task } from 'ember-concurrency'; +import { keepLatestTask, task } from 'ember-concurrency'; import { cloneDeep } from 'lodash-es'; import { isEqual } from 'lodash-es'; @@ -131,7 +131,10 @@ import type { SearchResource } from '../resources/search'; import type * as CardAPI from '@cardstack/base/card-api'; import type { CardDef, BaseDef } from '@cardstack/base/card-api'; import type { FileDef } from '@cardstack/base/file-api'; -import type { RealmEventContent } from '@cardstack/base/matrix-event'; +import type { + IncrementalIndexEventContent, + RealmEventContent, +} from '@cardstack/base/matrix-event'; export { CardErrorJSONAPI, CardSaveSubscriber }; @@ -1884,49 +1887,75 @@ export default class StoreService extends Service implements StoreInterface { return; } let invalidations = event.invalidations as string[]; - // Count reloads triggered while handling this event, shared by the - // realm-event and (when a loader reset fires) rebuild telemetry. - let reloadsTriggered = 0; let ownWrite = event.clientRequestId ? this.cardService.clientRequestIds.has(event.clientRequestId) : false; - let rebuildStart: number | undefined; - let rebuildTriggerModules: string[] = []; - let rebuildModulesRefetched = 0; - let rebuildTask: { then: (...a: unknown[]) => unknown } | undefined; - if ( - invalidations.find( - (i) => - hasExecutableExtension(i) && + // The invalidation triggers a rebuild when it touches an already-loaded + // executable module: the loader must be flushed so the updated code is + // picked up before the open card graph re-runs. Net-new modules that were + // never loaded don't need one. Once a rebuild is in flight, fold every + // further executable invalidation into it regardless of `isModuleLoaded` — + // the freshly reset loader reports those modules as not-loaded, so the + // probe alone would drop a mid-burst code change. + let executableInvalidations = invalidations.filter(hasExecutableExtension); + let needsRebuild = + executableInvalidations.length > 0 && + (this.rebuildForCodeChange.isRunning || + executableInvalidations.some((i) => this.loaderService.loader.isModuleLoaded(i), - ) - ) { - // the invalidation included code changes to modules that are already - // loaded. in this case we need to flush the loader so that we can pick - // up the updated code before re-running the card. net-new modules that - // have never been loaded don't require a loader reset. + )); + + let reloadsTriggered = 0; + if (needsRebuild) { + // Coalesce the rebuild. `keepLatestTask` keeps at most one rebuild in + // flight and one pending, so a burst of executable invalidations arriving + // faster than a rebuild completes collapses into at most 2 rebuilds; the + // final rebuild re-fetches current server state, so the end result + // reflects the latest generation. This is scheduling only and does not + // touch reactivity: an isolated change triggers a single reset and + // re-render. if (telemetry?.isEnabled) { - rebuildStart = performance.now(); - rebuildTriggerModules = invalidations.filter((i) => - hasExecutableExtension(i), - ); - rebuildModulesRefetched = invalidations.filter( - (i) => - hasExecutableExtension(i) && - this.loaderService.loader.isModuleLoaded(i), - ).length; + this.#accumulatePendingRebuild(event.realmURL, executableInvalidations); } - this.loaderService.resetLoader(); - this.store.reset(); - rebuildTask = this.reestablishReferences.perform() as unknown as { - then: (...a: unknown[]) => unknown; - }; + this.rebuildForCodeChange.perform(); + // The rebuild subsumes the per-invalidation reloads below: store.reset + // empties the graph and reestablishReferences re-fetches every live + // reference, so running that loop too would only re-fetch cards the + // rebuild is about to discard. + } else { + reloadsTriggered = this.#reloadInvalidatedInstances(event, invalidations); } + if (telemetry?.isEnabled) { + telemetry.recordEvent({ + event_type: 'realm-event', + realm: event.realmURL, + index_type: 'incremental', + invalidations_count: invalidations.length, + invalidated_ids: invalidations.slice(0, 50), + reloads_triggered: reloadsTriggered, + own_write: ownWrite, + processing_ms: + processingStart !== undefined + ? Math.round(performance.now() - processingStart) + : 0, + }); + } + }; + + // Reload the individual cards / file-meta resources named by an incremental + // invalidation that did not trigger a full rebuild. Returns the number of + // reloads kicked off (for realm-event telemetry). + #reloadInvalidatedInstances( + event: IncrementalIndexEventContent, + invalidations: string[], + ): number { + let reloadsTriggered = 0; for (let invalidation of invalidations) { if (hasExecutableExtension(invalidation)) { - // we already dealt with this + // Executable modules have no card instance to reload here; when an + // already-loaded one changed, the coalesced rebuild handled it. continue; } let fileMetaInstance = @@ -2038,47 +2067,8 @@ export default class StoreService extends Service implements StoreInterface { } } - if (telemetry?.isEnabled) { - telemetry.recordEvent({ - event_type: 'realm-event', - realm: event.realmURL, - index_type: 'incremental', - invalidations_count: invalidations.length, - invalidated_ids: invalidations.slice(0, 50), - reloads_triggered: reloadsTriggered, - own_write: ownWrite, - processing_ms: - processingStart !== undefined - ? Math.round(performance.now() - processingStart) - : 0, - }); - - // When a loader reset was triggered, report the rebuild once the - // reference re-establishment settles so its duration spans the whole - // reset → reload cascade. cards_reloaded reflects the reloads counted - // above (this handler runs to completion before the task settles). - if (rebuildStart !== undefined && rebuildTask) { - let reporter = telemetry; - let startedAt = rebuildStart; - let triggerModules = rebuildTriggerModules.slice(0, 20); - let modulesRefetched = rebuildModulesRefetched; - let realmURL = event.realmURL; - Promise.resolve(rebuildTask as unknown as Promise) - .then(() => { - reporter.recordEvent({ - event_type: 'rebuild', - realm: realmURL, - duration_ms: Math.round(performance.now() - startedAt), - trigger_modules: triggerModules, - trigger_module: triggerModules[0] ?? '', - modules_refetched: modulesRefetched, - cards_reloaded: reloadsTriggered, - }); - }) - .catch(() => {}); - } - } - }; + return reloadsTriggered; + } private loadInstanceTask = task( async (idOrDoc: string | LooseSingleCardDocument) => { @@ -2125,6 +2115,74 @@ export default class StoreService extends Service implements StoreInterface { await Promise.all( [...remoteIds].map((id) => this.getCardInstance({ idOrDoc: id })), ); + return remoteIds.size; + }); + + // Telemetry metadata for the pending/in-flight coalesced rebuild, merged + // across every executable invalidation that collapses into it. Drained when + // the rebuild task begins so the emitted `rebuild` event describes exactly + // the code changes that rebuild picked up. Only populated when telemetry is + // enabled; the coalescing itself does not depend on it. + #pendingRebuild: + | { + realm: string; + triggerModules: Set; + modulesRefetched: Set; + events: number; + } + | undefined = undefined; + + #accumulatePendingRebuild(realm: string, executableInvalidations: string[]) { + let pending = (this.#pendingRebuild ??= { + realm, + triggerModules: new Set(), + modulesRefetched: new Set(), + events: 0, + }); + // The most recent event's realm labels the rebuild; a burst is + // overwhelmingly single-realm, and the final generation wins regardless. + pending.realm = realm; + pending.events++; + for (let module of executableInvalidations) { + pending.triggerModules.add(module); + if (this.loaderService.loader.isModuleLoaded(module)) { + pending.modulesRefetched.add(module); + } + } + } + + // Coalesced client rebuild: flush the loader, reset the store, and re-fetch + // every live card reference. `keepLatest` bounds a write burst to one + // in-flight rebuild plus one pending — intermediate events collapse into the + // pending slot — so a burst of rapid executable invalidations costs at most 2 + // rebuilds regardless of its length. The final rebuild re-fetches current + // server state, so the end state reflects the latest generation. This is + // scheduling only: each rebuild is the full load-bearing reset, not a partial + // one. + private rebuildForCodeChange = keepLatestTask(async () => { + let telemetry = this.#clientTelemetry(); + let pending = this.#pendingRebuild; + this.#pendingRebuild = undefined; + let rebuildStart = + telemetry?.isEnabled && pending ? performance.now() : undefined; + + this.loaderService.resetLoader(); + this.store.reset(); + let cardsReloaded = await this.reestablishReferences.perform(); + + if (telemetry?.isEnabled && pending && rebuildStart !== undefined) { + let triggerModules = [...pending.triggerModules].slice(0, 20); + telemetry.recordEvent({ + event_type: 'rebuild', + realm: pending.realm, + duration_ms: Math.round(performance.now() - rebuildStart), + trigger_modules: triggerModules, + trigger_module: triggerModules[0] ?? '', + modules_refetched: pending.modulesRefetched.size, + cards_reloaded: cardsReloaded ?? 0, + coalesced_events: pending.events, + }); + } }); private reloadTask = task(async (instance: CardDef) => { diff --git a/packages/host/tests/integration/store-test.gts b/packages/host/tests/integration/store-test.gts index 03a597e999..50688e1b6f 100644 --- a/packages/host/tests/integration/store-test.gts +++ b/packages/host/tests/integration/store-test.gts @@ -3251,4 +3251,140 @@ module('Integration | Store', function (hooks) { `reference count for ${jade} is 1`, ); }); + + // Count full rebuilds by the loader flush each one performs: the coalesced + // rebuild calls resetLoader exactly once, and nothing else flushes the loader + // in these tests, so resetLoader-call-count == rebuild-count. + function countRebuilds() { + let count = 0; + let realResetLoader = loaderService.resetLoader.bind(loaderService); + loaderService.resetLoader = (( + options?: Parameters[0], + ) => { + count++; + return realResetLoader(options); + }) as LoaderService['resetLoader']; + return { + get count() { + return count; + }, + restore() { + loaderService.resetLoader = realResetLoader; + }, + }; + } + + async function renderCard(id: string) { + class Driver { + @tracked id: string | undefined; + } + let driver = new Driver(); + class ResourceConsumer extends GlimmerComponent { + resource = getCard(this, () => driver.id); + get renderedCard() { + return this.resource.card?.constructor.getComponent(this.resource.card); + } + + } + await renderComponent( + class TestDriver extends GlimmerComponent { + + }, + ); + driver.id = id; + await waitFor(`[data-test-rendered-card="${id}"]`, { timeout: 5_000 }); + } + + test('a burst of executable invalidations coalesces to at most two rebuilds', async function (assert) { + let hassan = `${testRealmURL}Person/hassan`; + await renderCard(hassan); + + let personModule = `${testRealmURL}person.gts`; + assert.true( + loaderService.loader.isModuleLoaded(personModule), + 'precondition: the person module is loaded', + ); + + let rebuilds = countRebuilds(); + try { + // Deliver executable invalidations faster than a rebuild can complete — + // synchronously, before the first rebuild's re-fetch settles. + let event: RealmEventContent = { + eventName: 'index', + indexType: 'incremental', + realmURL: testRealmURL, + invalidations: [personModule], + }; + for (let i = 0; i < 6; i++) { + (storeService as any).handleInvalidations(event); + } + await settled(); + } finally { + rebuilds.restore(); + } + + let coalesced = rebuilds.count >= 1 && rebuilds.count <= 2; + assert.ok( + coalesced, + `6 rapid executable invalidations coalesce to at most 2 rebuilds (saw ${rebuilds.count})`, + ); + + // End state reflects the latest generation: the open card is re-established + // and rendered against current server state. + await waitFor(`[data-test-rendered-card="${hassan}"]`, { timeout: 5_000 }); + let instance = storeService.peek(hassan); + assert.true( + isCardInstance(instance), + 'the open card is re-established after the burst', + ); + assert.strictEqual( + (instance as any).name, + 'Hassan', + 'the re-established card reflects current server state', + ); + assert.strictEqual( + storeService.getReferenceCount(hassan), + 1, + 'reference count stays balanced across the burst', + ); + }); + + test('an isolated executable invalidation still triggers exactly one rebuild', async function (assert) { + let hassan = `${testRealmURL}Person/hassan`; + await renderCard(hassan); + + let personModule = `${testRealmURL}person.gts`; + assert.true( + loaderService.loader.isModuleLoaded(personModule), + 'precondition: the person module is loaded', + ); + + let rebuilds = countRebuilds(); + try { + (storeService as any).handleInvalidations({ + eventName: 'index', + indexType: 'incremental', + realmURL: testRealmURL, + invalidations: [personModule], + } as RealmEventContent); + await settled(); + } finally { + rebuilds.restore(); + } + + assert.strictEqual( + rebuilds.count, + 1, + 'a single executable invalidation resets exactly once, as before', + ); + await waitFor(`[data-test-rendered-card="${hassan}"]`, { timeout: 5_000 }); + assert.true( + isCardInstance(storeService.peek(hassan)), + 'the card is re-established after the single rebuild', + ); + }); }); From ded06ac765674508ba1cf3be3b0ae851fe80a954 Mon Sep 17 00:00:00 2001 From: Hassan Abdel-Rahman Date: Fri, 24 Jul 2026 10:50:55 -0400 Subject: [PATCH 2/3] Give the patch save test realistic waitUntil slack MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `can skip waiting for the save when patching an instance` waited on the fire-and-forget save with a bare `waitUntil(() => didSave)`, i.e. the 1s default. That save completes a write + incremental-index round-trip before `onSave` fires — measured at ~1.8s under load — so the assertion timed out even though the save always succeeds. Its sibling `can skip waiting for the save when adding to the store` already waits with `{ timeout: 10000 }`; match it. Co-Authored-By: Claude Opus 4.8 (1M context) --- packages/host/tests/integration/store-test.gts | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/packages/host/tests/integration/store-test.gts b/packages/host/tests/integration/store-test.gts index 50688e1b6f..7489541d1b 100644 --- a/packages/host/tests/integration/store-test.gts +++ b/packages/host/tests/integration/store-test.gts @@ -1923,7 +1923,11 @@ module('Integration | Store', function (hooks) { ); assert.false(didSave, 'instance has not been persisted yet'); - await waitUntil(() => didSave); + // The fire-and-forget save (doNotWaitForPersist) completes a write + + // incremental index round-trip before onSave fires — comfortably over the + // 1s waitUntil default under load. Match the save slack the sibling + // "adding to the store" test uses. + await waitUntil(() => didSave, { timeout: 10000 }); let file = await testRealmAdapter.openFile('Person/hassan.json'); assert.strictEqual( From 32f16f2a502e6adce4b6b8c3ffaec20444c39029 Mon Sep 17 00:00:00 2001 From: Hassan Abdel-Rahman Date: Fri, 24 Jul 2026 11:06:55 -0400 Subject: [PATCH 3/3] Restore original resetLoader identity in the rebuild-count helper The rebuild-count spy restored `loaderService.resetLoader` to a bound copy of the method, leaving the service with a different method identity than it started with. Capture the original reference, invoke it via `.call` for the correct `this`, and restore that exact reference. Co-Authored-By: Claude Opus 4.8 (1M context) --- packages/host/tests/integration/store-test.gts | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/packages/host/tests/integration/store-test.gts b/packages/host/tests/integration/store-test.gts index 7489541d1b..660d2ea600 100644 --- a/packages/host/tests/integration/store-test.gts +++ b/packages/host/tests/integration/store-test.gts @@ -3261,19 +3261,19 @@ module('Integration | Store', function (hooks) { // in these tests, so resetLoader-call-count == rebuild-count. function countRebuilds() { let count = 0; - let realResetLoader = loaderService.resetLoader.bind(loaderService); - loaderService.resetLoader = (( + let original = loaderService.resetLoader; + loaderService.resetLoader = function ( options?: Parameters[0], - ) => { + ) { count++; - return realResetLoader(options); - }) as LoaderService['resetLoader']; + return original.call(loaderService, options); + } as LoaderService['resetLoader']; return { get count() { return count; }, restore() { - loaderService.resetLoader = realResetLoader; + loaderService.resetLoader = original; }, }; }