diff --git a/docs/features/site-shell.md b/docs/features/site-shell.md index c0bcc632c..6b2550a96 100644 --- a/docs/features/site-shell.md +++ b/docs/features/site-shell.md @@ -611,6 +611,10 @@ exponential backoff, and each (re)connect re-runs syncStep1 so Yjs state vectors pull exactly the missed delta. `usePersistence` HTTP-loads the document once for first paint, then connects the provider — edits gate on each doc's first sync so an unseeded doc can never receive local ops. +Initial sync projections collect socket tasks over a 16 ms frame and commit +rows plus their roster to the store together. Write gates open after that +projection, so tools never edit a stale HTTP snapshot. Live remote updates +retain microtask projection; undo/redo project synchronously. Because every committed edit is a frame, burst-prone inputs coalesce before they commit: the `ColorInput` primitive throttles picker-drag change events (leading fire for instant clicks, one trailing fire with the final value), diff --git a/src/__tests__/collab/paletteBurstWrite.test.ts b/src/__tests__/collab/paletteBurstWrite.test.ts index 784dfdb21..882228ab2 100644 --- a/src/__tests__/collab/paletteBurstWrite.test.ts +++ b/src/__tests__/collab/paletteBurstWrite.test.ts @@ -24,6 +24,7 @@ import type { CollabResetListener, } from '@site/collab/collabProvider' import { clearCollabBlockNotice } from '@site/store/slices/site/collabNotices' +import { whenCollabWritable } from '@site/store/slices/site/collabWriteGate' import { useEditorStore } from '@site/store/store' import { runSetColorTokens } from '@site/agent/tokenRunners' import '@modules/base/index' @@ -102,8 +103,7 @@ describe('installing a colour palette in one call', () => { const provider = deferredProvider() connectCollabProvider(provider) provider.releaseAll() - // Let the whenSynced promises settle so the gates are open. - await new Promise((resolve) => setTimeout(resolve, 10)) + expect(await whenCollabWritable()).toBe(true) const result = useEditorStore.getState().upsertFrameworkColorTokens(PALETTE) @@ -117,7 +117,7 @@ describe('installing a colour palette in one call', () => { const provider = deferredProvider() connectCollabProvider(provider) provider.releaseAll() - await new Promise((resolve) => setTimeout(resolve, 10)) + expect(await whenCollabWritable()).toBe(true) // Two DIFFERENT entries normalizing to the same slug must not collapse: // the second is a distinct token and gets a suffixed slug, exactly as a @@ -135,7 +135,7 @@ describe('installing a colour palette in one call', () => { const provider = deferredProvider() connectCollabProvider(provider) provider.releaseAll() - await new Promise((resolve) => setTimeout(resolve, 10)) + expect(await whenCollabWritable()).toBe(true) useEditorStore.getState().upsertFrameworkColorTokens(PALETTE) const second = useEditorStore.getState().upsertFrameworkColorTokens(PALETTE) diff --git a/src/__tests__/collab/projectionBatches.test.ts b/src/__tests__/collab/projectionBatches.test.ts new file mode 100644 index 000000000..3b334d113 --- /dev/null +++ b/src/__tests__/collab/projectionBatches.test.ts @@ -0,0 +1,92 @@ +import { afterEach, expect, it } from 'bun:test' +import * as Y from 'yjs' +import { Awareness } from 'y-protocols/awareness' +import '@modules/base' +import { encodeCollabDocId, MAIN_SITE_DOC_ID, metaMap, seedPageDoc, seedSiteDoc } from '@core/collab' +import type { CollabProvider, BoundCollabDoc } from '@site/collab/collabProvider' +import { connectCollabProvider, disconnectCollabProvider } from '@site/store/slices/site/collabBinding' +import { whenCollabWritable } from '@site/store/slices/site/collabWriteGate' +import { clearCollabBlockNotice } from '@site/store/slices/site/collabNotices' +import { useEditorStore } from '@site/store/store' +import { makePage, makeSite } from '../fixtures' + +function seededProvider(site: ReturnType) { + const presence = new Y.Doc() + const awareness = new Awareness(presence) + const bound = new Map() + const releases: Array<() => void> = [] + const provider: CollabProvider = { + bind(id) { + const existing = bound.get(id) + if (existing) return existing + const doc = new Y.Doc() + if (id === MAIN_SITE_DOC_ID) seedSiteDoc(doc, site) + else { + const page = site.pages.find((p) => pageId(p.id) === id) + if (page) seedPageDoc(doc, page) + } + let release = () => {} + const whenSynced = new Promise((resolve) => { release = resolve }) + const entry = { doc, synced: false, whenSynced } + bound.set(id, entry) + releases.push(() => { entry.synced = true; release() }) + return entry + }, + unbind(id) { bound.get(id)?.doc.destroy(); bound.delete(id) }, + awareness, + status: () => 'connected', + canSend: () => true, + reconnectNow() {}, + onStatus: () => () => {}, + onReset: () => () => {}, + destroy() { + for (const entry of bound.values()) entry.doc.destroy() + awareness.destroy() + presence.destroy() + }, + } + return { provider, releaseAll: () => { for (const release of releases) release() } } +} + +function pageId(id: string) { + return encodeCollabDocId({ kind: 'page', branchId: 'main', rowId: id }) +} + +afterEach(() => { + disconnectCollabProvider() + useEditorStore.getState().clearSite() + clearCollabBlockNotice() +}) + +it('commits an initial row burst once and opens writes only after its projection', async () => { + const pages = Array.from({ length: 30 }, (_, i) => makePage({ id: `page-${i}`, slug: `page-${i}` })) + const site = makeSite({ pages }) + useEditorStore.getState().loadSite(site) + const { provider, releaseAll } = seededProvider(makeSite({ + pages: pages.map((p, i) => ({ ...p, title: `Server ${i}` })), + })) + connectCollabProvider(provider) + let changes = 0 + const off = useEditorStore.subscribe((next, previous) => { + if (next.site !== previous.site) changes++ + }) + try { + releaseAll() + await Promise.resolve() + expect(await whenCollabWritable(1)).toBe(false) + useEditorStore.getState().updateSiteName('Premature write') + expect(useEditorStore.getState().site!.name).toBe(site.name) + expect(await whenCollabWritable()).toBe(true) + expect(changes).toBe(1) + expect(useEditorStore.getState().site!.pages.map((p) => p.title)).toEqual( + pages.map((_p, i) => `Server ${i}`), + ) + // After startup a live peer edit still projects in the next microtask. + metaMap(provider.bind(pageId(pages[0].id)).doc).set('title', 'Live edit') + await Promise.resolve() + expect(useEditorStore.getState().site!.pages[0].title).toBe('Live edit') + expect(changes).toBe(2) + } finally { + off() + } +}) diff --git a/src/admin/pages/site/store/slices/site/collabBinding.ts b/src/admin/pages/site/store/slices/site/collabBinding.ts index 7aadd3a81..faf0bafe0 100644 --- a/src/admin/pages/site/store/slices/site/collabBinding.ts +++ b/src/admin/pages/site/store/slices/site/collabBinding.ts @@ -31,9 +31,7 @@ */ import * as Y from 'yjs' import type { Patches } from 'mutative' -import type { Page, SiteDocument, SiteShell } from '@core/page-tree' -import type { VisualComponent } from '@core/visualComponents' -import type { SavedLayout } from '@core/layouts' +import type { SiteDocument } from '@core/page-tree' import { applySitePatchesToDocs, createCollabDocSet, @@ -43,10 +41,6 @@ import { LOCAL_ORIGIN, metaMap, parseCollabDocId, - projectComponentDoc, - projectLayoutDoc, - projectPageDoc, - projectSiteDoc, rostersMap, seedComponentDoc, seedLayoutDoc, @@ -58,14 +52,12 @@ import { treeMap, type CollabDocSet, } from '@core/collab' +import { projectCollabDocument } from './collabProjection' import { allDocIdsForSite, collabBranchId, notifyCollabBranchGone } from './collabBranch' -import { clonePackageJson } from '@core/site-dependencies/manifest' -import { cloneSiteRuntimeConfig } from '@core/site-runtime' -import { validateSite } from '@core/persistence/validate' import type { EditorStoreApi } from '@site/store/types' import { pruneCanvasSelectionDraft } from '../selectionSlice' import type { Awareness } from 'y-protocols/awareness' -import type { CollabProvider } from '@site/collab/collabProvider' +import type { BoundCollabDoc, CollabProvider } from '@site/collab/collabProvider' import { collabBlockToast, clearCollabBlockNotice, @@ -100,6 +92,7 @@ let provider: CollabProvider | null = null let detachProviderReset: (() => void) | null = null let detachProviderStatus: (() => void) | null = null const pendingProjections = new Set() +const pendingSyncCompletions = new Map void>() let projectionFlushScheduled = false /** * The exact store `site` object the doc world currently mirrors. Every path @@ -395,157 +388,53 @@ export function collabClearHistory(): void { function flushProjections(): void { const batch = [...pendingProjections] pendingProjections.clear() + const api = storeApi + const site = api?.getState().site + if (batch.length === 0) return + const completeSync = () => { + for (const id of batch) { + pendingSyncCompletions.get(id)?.() + pendingSyncCompletions.delete(id) + } + } + if (!api || !site) { completeSync(); return } // Site doc last — it assembles rows the row projections just refreshed. batch.sort((a, b) => (isSiteDocId(a) ? 1 : 0) - (isSiteDocId(b) ? 1 : 0)) - for (const id of batch) projectDocIntoStore(id) + let nextSite = site + const context = { + getDoc: (docId: string) => docs.get(docId), + branchId: collabBranchId(), + bindDoc: bindDocThroughProvider, + } + for (const id of batch) nextSite = projectCollabDocument(nextSite, id, context) + if (nextSite === site) { completeSync(); return } + alignedSiteRef = nextSite + api.setState((draft) => { + draft.site = nextSite + if (nextSite.packageJson !== site.packageJson) draft.packageJson = nextSite.packageJson + if (nextSite.runtime !== site.runtime) draft.siteRuntime = nextSite.runtime + if (!nextSite.pages.some((p) => p.id === draft.activePageId)) { + draft.activePageId = nextSite.pages[0]?.id ?? null + } + // Project rows and their roster atomically before pruning selections. + pruneCanvasSelectionDraft(draft) + }) + completeSync() } function scheduleProjection(docId: string): void { pendingProjections.add(docId) if (projectionFlushScheduled) return projectionFlushScheduled = true - queueMicrotask(() => { + const flush = () => { projectionFlushScheduled = false flushProjections() - }) -} - -function rowFromDoc(docId: string): Page | VisualComponent | SavedLayout | null { - const parsed = parseCollabDocId(docId) - if (!parsed || parsed.kind === 'site') return null - const doc = docs.get(docId) - if (!doc) return null - if (parsed.kind === 'page') { - const page = projectPageDoc(doc, parsed.rowId) - return page.rootNodeId ? page : null - } - if (parsed.kind === 'component') { - const vc = projectComponentDoc(doc, parsed.rowId) - return vc.tree.rootNodeId ? vc : null } - const layout = projectLayoutDoc(doc, parsed.rowId) - return layout.rootNodeId ? layout : null -} - -function projectDocIntoStore(docId: string): void { - const api = storeApi - if (!api) return - const state = api.getState() - const site = state.site - if (!site) return - const parsed = parseCollabDocId(docId) - if (!parsed) return - - if (parsed.kind === 'site') { - const doc = docs.get(docId) - if (!doc) return - const projected = projectSiteDoc(doc) - if (Object.keys(projected.shell).length === 0) return - // The projected shell is untyped wire data — validate it before it enters - // the store, exactly like the HTTP load path (validateSite) and the relay's - // persist path both do. `validateSite` is tolerant of individual malformed - // entries (drops bad style rules / conditions / files rather than - // rejecting the whole shell), so one corrupt rule from any source can't - // crash a panel. `id`/`updatedAt` are non-collaborative — inject them like - // the persist path. If the shell is not yet coherent (mid-sync), skip this - // tick; the next projection re-runs once it is. - let shell: SiteShell - try { - shell = validateSite({ - ...projected.shell, - id: 'default', - updatedAt: - typeof projected.shell.updatedAt === 'number' ? projected.shell.updatedAt : Date.now(), - }) - } catch (err) { - console.warn('[collabBinding] projected shell failed validation — projection skipped:', err) - return - } - const byId = { - pages: new Map(site.pages.map((p) => [p.id, p])), - components: new Map(site.visualComponents.map((vc) => [vc.id, vc])), - layouts: new Map(site.layouts.map((l) => [l.id, l])), - } - const assemble = ( - ids: readonly string[], - existing: Map, - kind: 'page' | 'component' | 'layout', - ): T[] => { - const rows: T[] = [] - for (const id of ids) { - const known = existing.get(id) - if (known) { - rows.push(known) - continue - } - const rowDocId = encodeCollabDocId({ kind, branchId: collabBranchId(), rowId: id }) - const fresh = rowFromDoc(rowDocId) as T | null - if (fresh) { - rows.push(fresh) - continue - } - // A peer created this row — its doc isn't bound here yet. Bind it; - // the whenSynced hook re-projects the site once content arrives. - bindDocThroughProvider(rowDocId) - } - return rows - } - const nextSite: SiteDocument = { - ...site, - ...shell, - pages: assemble(projected.rosters.pages, byId.pages, 'page'), - visualComponents: assemble(projected.rosters.components, byId.components, 'component'), - layouts: assemble(projected.rosters.layouts, byId.layouts, 'layout'), - } - if (projected.shell.conditions === undefined) delete nextSite.conditions - const packageJson = clonePackageJson(nextSite.packageJson) - const siteRuntime = cloneSiteRuntimeConfig(nextSite.runtime) - const alignedSite = { ...nextSite, packageJson, runtime: siteRuntime } - alignedSiteRef = alignedSite - api.setState((draft) => { - draft.site = alignedSite - draft.packageJson = packageJson - draft.siteRuntime = siteRuntime - if (!nextSite.pages.some((p) => p.id === draft.activePageId)) { - draft.activePageId = nextSite.pages[0]?.id ?? null - } - // A roster change can drop the whole document the selection lives in (a - // peer deleted the page, or an undo removed it). Prune AFTER site + - // activePageId land, since the pruner resolves the active tree from them. - pruneCanvasSelectionDraft(draft) - }) - return - } - - const row = rowFromDoc(docId) - const collection = - parsed.kind === 'page' ? 'pages' : parsed.kind === 'component' ? 'visualComponents' : 'layouts' - const rows = site[collection] as Array<{ id: string }> - const index = rows.findIndex((r) => r.id === parsed.rowId) - if (!row) { - if (index === -1) return - const nextRows = rows.filter((r) => r.id !== parsed.rowId) - const nextSite = { ...site, [collection]: nextRows } as SiteDocument - alignedSiteRef = nextSite - api.setState((draft) => { - draft.site = nextSite - pruneCanvasSelectionDraft(draft) - }) - return - } - const nextRows = index === -1 ? [...rows, row] : rows.map((r, i) => (i === index ? row : r)) - const nextSite = { ...site, [collection]: nextRows } as SiteDocument - alignedSiteRef = nextSite - api.setState((draft) => { - draft.site = nextSite - // The freshly projected row may have lost nodes — a peer deleted them, or a - // Y.UndoManager undo reverted their creation. Prune by tree-membership, the - // same way a local delete does: survivors keep their selection, dead ids - // (including descendants swept with a subtree) drop out, and an inline-edit - // session on a vanished node is closed. `pruneCanvasSelectionDraft` reads - // the ACTIVE tree, so it self-limits to the doc the user is looking at. - pruneCanvasSelectionDraft(draft) - }) + // Initial sync spans thousands of socket tasks on large sites. Group those + // tasks over a frame while writes are gated; live edits retain microtask + // projection so the next local mutation reads the latest remote content. + if (provider && anyGateUnsynced()) setTimeout(flush, 16) + else queueMicrotask(flush) } // --------------------------------------------------------------------------- @@ -574,6 +463,8 @@ export function resetCollabDocsFromSite(site: SiteDocument | null): void { // Drop any projection still queued for the OLD docs — flushing it against // the fresh doc set would project empty rows into the just-loaded site. pendingProjections.clear() + for (const complete of pendingSyncCompletions.values()) complete() + pendingSyncCompletions.clear() syncUndoFlags() if (!site) return @@ -615,15 +506,25 @@ function bindDocThroughProvider(docId: string): void { const binding = provider.bind(docId) docs.set(docId, binding.doc) ensureManaged(docId, binding.doc) - const gate = registerProviderGate(docId, binding.whenSynced) - gate.synced = binding.synced + registerProjectionGate(docId, binding, true) +} + +/** Writes open only after the server seed has also reached the store. */ +function registerProjectionGate(docId: string, binding: BoundCollabDoc, assembleRoster: boolean): void { + let resolveProjection!: () => void + const projected = new Promise((resolve) => { resolveProjection = resolve }) + const gate = registerProviderGate(docId, projected) void binding.whenSynced.then(() => { - gate.synced = true + if (docs.get(docId) !== binding.doc) { resolveProjection(); return } + pendingSyncCompletions.set(docId, () => { + gate.synced = true + resolveProjection() + }) scheduleProjection(docId) // A row doc bound on demand (a peer created the row) re-assembles the // site once its content arrives — the roster projection skipped it // while it was empty. - if (!isSiteDocId(docId)) scheduleProjection(siteDocId(collabBranchId())) + if (assembleRoster && !isSiteDocId(docId)) scheduleProjection(siteDocId(collabBranchId())) }) } @@ -673,12 +574,7 @@ export function connectCollabProvider(next: CollabProvider): void { const rebound = next.bind(docId) docs.set(docId, rebound.doc) ensureManaged(docId, rebound.doc) - const gate = registerProviderGate(docId, rebound.whenSynced) - gate.synced = rebound.synced - void rebound.whenSynced.then(() => { - gate.synced = true - scheduleProjection(docId) - }) + registerProjectionGate(docId, rebound, false) }) const site = storeApi?.getState().site ?? null resetCollabDocsFromSite(site) diff --git a/src/admin/pages/site/store/slices/site/collabProjection.ts b/src/admin/pages/site/store/slices/site/collabProjection.ts new file mode 100644 index 000000000..43f7291ce --- /dev/null +++ b/src/admin/pages/site/store/slices/site/collabProjection.ts @@ -0,0 +1,116 @@ +import type { Page, SiteDocument, SiteShell } from '@core/page-tree' +import type { VisualComponent } from '@core/visualComponents' +import type { SavedLayout } from '@core/layouts' +import { encodeCollabDocId, parseCollabDocId, projectComponentDoc, projectLayoutDoc, projectPageDoc, projectSiteDoc, type CollabDocSet } from '@core/collab' +import { clonePackageJson } from '@core/site-dependencies/manifest' +import { cloneSiteRuntimeConfig } from '@core/site-runtime' +import { validateSite } from '@core/persistence/validate' + +interface ProjectionContext { + getDoc: CollabDocSet['get'] + branchId: string + bindDoc: (docId: string) => void +} + +function rowFromDoc(docId: string, context: ProjectionContext): Page | VisualComponent | SavedLayout | null { + const parsed = parseCollabDocId(docId) + if (!parsed || parsed.kind === 'site') return null + const doc = context.getDoc(docId) + if (!doc) return null + if (parsed.kind === 'page') { + const page = projectPageDoc(doc, parsed.rowId) + return page.rootNodeId ? page : null + } + if (parsed.kind === 'component') { + const vc = projectComponentDoc(doc, parsed.rowId) + return vc.tree.rootNodeId ? vc : null + } + const layout = projectLayoutDoc(doc, parsed.rowId) + return layout.rootNodeId ? layout : null +} + +export function projectCollabDocument(site: SiteDocument, docId: string, context: ProjectionContext): SiteDocument { + const parsed = parseCollabDocId(docId) + if (!parsed) return site + + if (parsed.kind === 'site') { + const doc = context.getDoc(docId) + if (!doc) return site + const projected = projectSiteDoc(doc) + if (Object.keys(projected.shell).length === 0) return site + // The projected shell is untyped wire data — validate it before it enters + // the store, exactly like the HTTP load path (validateSite) and the relay's + // persist path both do. `validateSite` is tolerant of individual malformed + // entries (drops bad style rules / conditions / files rather than + // rejecting the whole shell), so one corrupt rule from any source can't + // crash a panel. `id`/`updatedAt` are non-collaborative — inject them like + // the persist path. If the shell is not yet coherent (mid-sync), skip this + // tick; the next projection re-runs once it is. + let shell: SiteShell + try { + shell = validateSite({ + ...projected.shell, + id: 'default', + updatedAt: + typeof projected.shell.updatedAt === 'number' ? projected.shell.updatedAt : Date.now(), + }) + } catch (err) { + console.warn('[collabBinding] projected shell failed validation — projection skipped:', err) + return site + } + const byId = { + pages: new Map(site.pages.map((p) => [p.id, p])), + components: new Map(site.visualComponents.map((vc) => [vc.id, vc])), + layouts: new Map(site.layouts.map((l) => [l.id, l])), + } + const assemble = ( + ids: readonly string[], + existing: Map, + kind: 'page' | 'component' | 'layout', + ): T[] => { + const rows: T[] = [] + for (const id of ids) { + const known = existing.get(id) + if (known) { + rows.push(known) + continue + } + const rowDocId = encodeCollabDocId({ kind, branchId: context.branchId, rowId: id }) + const fresh = rowFromDoc(rowDocId, context) as T | null + if (fresh) { + rows.push(fresh) + continue + } + // A peer created this row — its doc isn't bound here yet. Bind it; + // the whenSynced hook re-projects the site once content arrives. + context.bindDoc(rowDocId) + } + return rows + } + const nextSite: SiteDocument = { + ...site, + ...shell, + pages: assemble(projected.rosters.pages, byId.pages, 'page'), + visualComponents: assemble(projected.rosters.components, byId.components, 'component'), + layouts: assemble(projected.rosters.layouts, byId.layouts, 'layout'), + } + if (projected.shell.conditions === undefined) delete nextSite.conditions + const packageJson = clonePackageJson(nextSite.packageJson) + const siteRuntime = cloneSiteRuntimeConfig(nextSite.runtime) + return { ...nextSite, packageJson, runtime: siteRuntime } + } + + const row = rowFromDoc(docId, context) + const collection = + parsed.kind === 'page' ? 'pages' : parsed.kind === 'component' ? 'visualComponents' : 'layouts' + const rows = site[collection] as Array<{ id: string }> + const index = rows.findIndex((r) => r.id === parsed.rowId) + if (!row) { + if (index === -1) return site + const nextRows = rows.filter((r) => r.id !== parsed.rowId) + return { ...site, [collection]: nextRows } as SiteDocument + } + const nextRows = index === -1 ? [...rows, row] : rows.map((r, i) => (i === index ? row : r)) + return { ...site, [collection]: nextRows } as SiteDocument +} + diff --git a/src/admin/pages/site/store/slices/site/collabWriteGate.ts b/src/admin/pages/site/store/slices/site/collabWriteGate.ts index 404e8c783..1b214383f 100644 --- a/src/admin/pages/site/store/slices/site/collabWriteGate.ts +++ b/src/admin/pages/site/store/slices/site/collabWriteGate.ts @@ -18,7 +18,7 @@ export interface ProviderGate { synced: boolean - /** Resolves on this doc's first sync. */ + /** Resolves after this doc's first sync has reached the editor store. */ whenSynced: Promise } @@ -45,7 +45,7 @@ export function clearProviderGates(): void { providerGates.clear() } -/** True when at least one bound doc has not finished its first sync. */ +/** True when at least one bound doc has not finished syncing into the store. */ export function anyGateUnsynced(): boolean { for (const gate of providerGates.values()) { if (!gate.synced) return true