diff --git a/.dev.vars.example b/.dev.vars.example new file mode 100644 index 0000000..581e9e5 --- /dev/null +++ b/.dev.vars.example @@ -0,0 +1,13 @@ +# IS_HELIX is a manual local-testing override for URL-based backend selection. +# Leave it unset for normal deployments. +# It accepts three values: +# false = always go to da-admin (service binding), regardless of the document hostname +# true = always go to Helix, rewriting every request to https://api.aem.live +# local = always go to Helix, rewriting every request to http://localhost:3000 +# When unset, the backend is chosen by the document URL: api.aem.live -> Helix, +# everything else -> da-admin. +# +# Copy this file to .dev.vars and uncomment a line to force the backend locally. +# Note: values set in .dev.vars are masked as "(hidden)" in wrangler's startup log; +# declare IS_HELIX in wrangler.toml [env.dev.vars] instead if you want it printed. +# IS_HELIX=local diff --git a/README.md b/README.md index ece91ca..ce80773 100644 --- a/README.md +++ b/README.md @@ -15,6 +15,13 @@ To run da-admin locally see https://github.com/adobe/da-admin/blob/main/README.m 1. In a terminal, run `npm run dev` this repo's folder. 1. The da-collab service API is available via http://localhost:4711 +#### Local backend testing override + +For manual local testing, `IS_HELIX` in `.dev.vars` can force Helix at `api.aem.live` (`true`), +Helix at `localhost:3000` (`local`), or da-admin (`false`). Leave it unset for normal URL-based +backend selection. Forced URL rewriting and the `*** Calling` log are isolated in the +`getLocalTestBackendOverride()` helper. + #### Access via da-live To access the locally running da-collab via da-live also running locally, first run da-live on your local machine @@ -32,6 +39,32 @@ To access da-collab and da-admin running on stage, open this URL in a browser: h 1. When passing in `?da-collab=local&da-admin=local` each service will set a localStorage value and will not clear until you use `?name-of-service=reset`. It is recommended to use an incognito browser window to ensure you don't forget about this setting. ## Additional details +### Helix backend polling +For Helix-backed documents, collab checks the backend ETag with a HEAD request every five seconds +while the session is connected. A changed ETag or a 404/410 response closes the session, cancels +pending local saves, and removes the stored restore anchor. Invalidation therefore does not +flush stale edits back to the backend, and reconnecting sessions reload the source document. +Ordinary connection closure and da-admin invalidation retain their existing save-flushing behavior. + +Helix ETags are compared without the weak `W/` prefix. If GET omits its ETag, collab brackets a +fresh read with HEAD requests to establish a matching content/version baseline. Without a +verified baseline, reads remain available but saves report an error instead of overwriting +an unknown version. A successful HEAD can establish a polling-only baseline without disconnecting editors. + +Helix saves use the known version in `If-Match` rather than `*`. A 412 is treated as an external +version conflict: pending saves are cancelled and the session is reloaded, not reported as a +deletion. The Helix API must return the written ETag on successful POST; if it does not, collab +logs the missing validator and pauses further protected saves until reload. da-admin retains +its existing multipart PUT and `If-Match: *` behavior. Complete lost-update protection additionally +requires the API to enforce the precondition atomically at the storage write. + +Failed HEAD requests and successful responses without an ETag are logged and retried on the next +interval without disconnecting editors. Only one HEAD request is in flight per document, and +responses overlapping a local save or belonging to an old session are ignored. Polling starts +only if initialization finishes for the current, connected document. Destroying the document +clears its polling timer and aborts any pending HEAD request; late initialization cannot restart +polling after disconnect. The idle-but-connected session policy is unchanged. + ### Recommendations 1. We recommend running `npm run lint` for linting. diff --git a/src/edge.js b/src/edge.js index ad3d6b1..50c3234 100644 --- a/src/edge.js +++ b/src/edge.js @@ -228,7 +228,7 @@ export async function handleApiRequest(request, env) { } const timingBeforeDaAdminHead = Date.now(); - const initialReq = await getBackend(docName, env.daadmin).fetch(docName, opts); + const initialReq = await getBackend(docName, env).fetch(docName, opts); timingDaAdminHeadDuration = Date.now() - timingBeforeDaAdminHead; @@ -429,7 +429,7 @@ export class DocRoom extends DurableObject { // Helix does not yet report auth actions, so grant collaborators // read,write; otherwise honour what da-admin reported. // TODO: remove the isHelixDoc branch once Helix reports auth actions. - const authActions = isHelixDoc(docName) + const authActions = isHelixDoc(docName, this.env) ? 'read,write' : request.headers.get('X-auth-actions') ?? ''; diff --git a/src/shareddoc.js b/src/shareddoc.js index a8ce353..e4cdfa2 100644 --- a/src/shareddoc.js +++ b/src/shareddoc.js @@ -23,10 +23,71 @@ const wsReadyStateConnecting = 0; const wsReadyStateOpen = 1; /** - * True for documents that live in Helix (api.aem.live) rather than da-admin. + * Convert a da-admin source path to the Helix path format. + * + * da-admin paths look like `/source/{org}/{site}/{rest}`; the equivalent Helix + * path is `/{org}/sites/{site}/source/{rest}`. Paths already in Helix format + * (`/{org}/sites/{site}/source/...`) and paths that don't match the da-admin + * source shape are returned unchanged. + * @param {string} pathname - the URL pathname (leading slash included) + */ +export const toHelixPath = (pathname) => { + const segs = pathname.split('/').filter(Boolean); + // Already Helix format: {org}/sites/{site}/source/... + if (segs[1] === 'sites' && segs[3] === 'source') { + return pathname; + } + // da-admin format: source/{org}/{site}/{rest} + if (segs[0] === 'source' && segs.length >= 3) { + const [, org, site, ...rest] = segs; + return `/${org}/sites/${site}/source/${rest.join('/')}`; + } + return pathname; +}; + +/** + * Manual local-testing override. Normal deployments leave IS_HELIX unset. + * 'false' forces da-admin without rewriting URLs; 'true' and 'local' force + * Helix at api.aem.live and localhost:3000, respectively. + * @param {object} [env] - the worker environment + * @returns {{ + * isHelix: boolean, + * rewriteUrl: (url: string) => string, + * } | undefined} + */ +function getLocalTestBackendOverride(env) { + if (env?.IS_HELIX === undefined) { + return undefined; + } + const mode = String(env.IS_HELIX); + let substituteOrigin; + if (mode === 'true') { + substituteOrigin = 'https://api.aem.live'; + } else if (mode === 'local') { + substituteOrigin = 'http://localhost:3000'; + } + + return { + isHelix: substituteOrigin !== undefined, + rewriteUrl: (url) => { + if (!substituteOrigin) { + return url; + } + const u = new URL(url); + const t = `${substituteOrigin}${toHelixPath(u.pathname)}${u.search}`; + console.log('*** Calling', t); + return t; + }, + }; +} + +/** + * True for Helix documents, using the URL unless a local-testing override is set. * @param {string} docName - the document URL + * @param {object} [env] - the worker environment (may carry `IS_HELIX`) */ -export const isHelixDoc = (docName) => docName.startsWith('https://api.aem.live/'); +export const isHelixDoc = (docName, env) => getLocalTestBackendOverride(env)?.isHelix + ?? docName.startsWith('https://api.aem.live/'); /** * Resolve the content backend for a document. @@ -36,18 +97,21 @@ export const isHelixDoc = (docName) => docName.startsWith('https://api.aem.live/ * through the da-admin service binding. * * @param {string} docName - the document URL - * @param {Fetcher} daadmin - the da-admin service binding + * @param {object} env - the worker environment; provides the da-admin service + * binding (`env.daadmin`) and optional local-testing backend configuration. * @returns {{ * fetch: (url: string, opts?: object) => Promise, * putReqData: (content: string, mimeType: string) => { body: *, size: number, headers: object }, * }} */ -export function getBackend(docName, daadmin) { - const isHelix = isHelixDoc(docName); +export function getBackend(docName, env) { + const localTestOverride = getLocalTestBackendOverride(env); + const isHelix = isHelixDoc(docName, env); return { // A fetch that already knows where to go. - fetch: (url, opts) => (isHelix ? globalThis : daadmin).fetch(url, opts), + fetch: (url, opts) => (isHelix ? globalThis : env.daadmin) + .fetch(localTestOverride?.rewriteUrl(url) ?? url, opts), // Build the body (and any body-specific headers) for a content PUT. // Helix takes the raw content with an explicit Content-Type; da-admin takes @@ -138,8 +202,9 @@ function getDocType(docName) { * the ydoc from the local cache map. * @param {ydoc} doc - the ydoc to close the connection for. * @param {WebSocket} conn - the websocket connection to close. + * @param {boolean} skipSave - skip flushing stale or re-entrant saves. */ -export const closeConn = async (doc, conn, isReentrant = false) => { +export const closeConn = async (doc, conn, skipSave = false) => { try { if (doc.conns.has(conn)) { const controlledIds = doc.conns.get(conn); @@ -154,10 +219,9 @@ export const closeConn = async (doc, conn, isReentrant = false) => { } if (doc.conns.size === 0) { - // Skip flushSave when called re-entrantly from persistence.update's closeAll - // loop — the in-flight save owns persistence; awaiting savingPromise here - // would deadlock because persistence.update hasn't returned yet. - if (doc.flushSave && !isReentrant) { + // External changes must discard pending edits. Re-entrant closes from + // persistence.update must also skip flushing to avoid awaiting their own save. + if (doc.flushSave && !skipSave && !doc.discardPendingChanges) { // eslint-disable-next-line no-console console.log('[docroom] Flushing pending save on last connection close', doc.name); await doc.flushSave(); @@ -317,6 +381,100 @@ export const showError = (ydoc, err) => { } }; +/** + * Invalidate the worker storage for the document, which will ensure that when accessed + * the worker will fetch the latest version of the document from the da-admin. + * Invalidation is implemented by closing all client connections to the doc, which will + * cause it to be reinitialised when accessed. + * @param {string} docName - The name of the document + * @param {boolean} discardPendingChanges - cancel pending saves instead of flushing them. + * @returns true if the document was found and invalidated, false otherwise. + */ +export const invalidateFromAdmin = async (docName, discardPendingChanges = false) => { + // eslint-disable-next-line no-console + console.log('[worker] Invalidate document cache', docName); + const ydoc = docs.get(docName); + if (ydoc) { + if (discardPendingChanges) { + ydoc.discardPendingChanges = true; + ydoc.cancelSave?.(); + } + // Remove the restore anchor before reconnecting, including when an external + // replacement has identical content but must discard stored pending edits. + const discardAnchor = discardPendingChanges && ydoc.storage?.delete + ? ydoc.storage.delete('lastsync') + : undefined; + // Await all closeConn calls so that flushSave + docs.delete complete before + // this function returns. Without this, an editor that reconnects during the + // flushSave window finds the stale ydoc still in `docs` and reuses it + // instead of creating a fresh one that loads the new da-admin content. + await Promise.all([ + discardAnchor, + ...Array.from(ydoc.conns.keys()).map((c) => closeConn(ydoc, c, discardPendingChanges)), + ]); + return true; + } else { + // eslint-disable-next-line no-console + console.log('[worker] Document not found', docName); + } + return false; +}; + +const normalizeEtag = (etag) => (etag ? etag.replace(/^W\//, '') : null); + +async function getHelixEtagBaseline(backend, docName, opts, initialResponse) { + let response = initialResponse; + try { + for (let attempt = 0; attempt < 2; attempt += 1) { + const etag = normalizeEtag(response.headers.get('etag')); + if (!response.ok || etag) { + return { response, etag }; + } + + // Without a GET validator, bracket a fresh read with matching HEAD tags. + // Adopting a HEAD tag for an earlier body could authorize a stale overwrite. + // eslint-disable-next-line no-await-in-loop + const before = await backend.fetch(docName, { ...opts, method: 'HEAD' }); + const beforeEtag = normalizeEtag(before.headers.get('etag')); + before.body?.cancel(); + if (!before.ok) { + response.body?.cancel(); + return { response: before, etag: null }; + } + if (!beforeEtag) { + break; + } + + // eslint-disable-next-line no-await-in-loop + const reloaded = await backend.fetch(docName, opts); + response.body?.cancel(); + response = reloaded; + const reloadedEtag = normalizeEtag(response.headers.get('etag')); + if (!response.ok || reloadedEtag) { + return { response, etag: reloadedEtag }; + } + + // eslint-disable-next-line no-await-in-loop + const after = await backend.fetch(docName, { ...opts, method: 'HEAD' }); + const afterEtag = normalizeEtag(after.headers.get('etag')); + after.body?.cancel(); + if (!after.ok) { + response.body?.cancel(); + return { response: after, etag: null }; + } + if (afterEtag === beforeEtag) { + return { response, etag: afterEtag }; + } + } + // eslint-disable-next-line no-console + console.warn('[docroom] Helix ETag baseline unavailable or unstable; protected saves require a reload', docName); + return { response, etag: null }; + } catch (error) { + response.body?.cancel(); + throw error; + } +} + export const persistence = { closeConn, @@ -324,19 +482,32 @@ export const persistence = { * Get the document from da-admin. * @param {string} docName - The document name * @param {string} auth - The authorization header - * @param {object} daadmin - The da-admin worker service binding + * @param {WSSharedDoc} ydoc - the ydoc; provides the da-admin service binding + * (`ydoc.daadmin`) and receives the document's ETag (`ydoc.etag`) read from + * the response for later HEAD-based change detection. * @returns {Promise} - The content of the document * @throws {Error} - If the document cannot be retrieved (including 404) */ - get: async (docName, auth, daadmin) => { + get: async (docName, auth, ydoc) => { const docType = getDocType(docName); const initalOpts = {}; if (auth) { initalOpts.headers = new Headers({ Authorization: auth }); } - const initialReq = await getBackend(docName, daadmin).fetch(docName, initalOpts); + const backend = getBackend(docName, ydoc); + const isHelix = isHelixDoc(docName, ydoc); + let initialReq = await backend.fetch(docName, initalOpts); + let etag; + if (isHelix && initialReq.ok) { + const baseline = await getHelixEtagBaseline(backend, docName, initalOpts, initialReq); + ({ response: initialReq, etag } = baseline); + // eslint-disable-next-line no-param-reassign + ydoc.etagUnverified = !etag; + } if (initialReq.ok) { + // eslint-disable-next-line no-param-reassign + ydoc.etag = isHelix ? etag : initialReq.headers.get('etag'); return docType === 'json' ? initialReq.json() : initialReq.text(); } else { const msg = `[docroom] Unable to get resource from da-admin: ${initialReq.status} - ${initialReq.statusText}`; @@ -367,12 +538,13 @@ export const persistence = { * @returns {Promise} The response from da-admin. */ put: async (ydoc, content) => { - const backend = getBackend(ydoc.name, ydoc.daadmin); + const backend = getBackend(ydoc.name, ydoc); + const isHelix = isHelixDoc(ydoc.name, ydoc); const mimeType = getDocType(ydoc.name) === 'json' ? 'application/json' : 'text/html'; const { body: putBody, size: bodySize, headers: bodyHeaders } = backend .putReqData(content, mimeType); - const opts = { method: isHelixDoc(ydoc.name) ? 'POST' : 'PUT', body: putBody }; + const opts = { method: isHelix ? 'POST' : 'PUT', body: putBody }; const keys = Array.from(ydoc.conns.keys()); const allReadOnly = keys.length > 0 && keys.every((con) => con.readOnly === true); if (allReadOnly) { @@ -381,8 +553,12 @@ export const persistence = { return { ok: true }; } + const storedEtag = normalizeEtag(ydoc.etag); + if (isHelix && (!storedEtag || storedEtag === '*' || ydoc.etagUnverified)) { + throw new Error('Cannot save Helix document without a verified ETag; reload the document'); + } const headers = { - 'If-Match': '*', + 'If-Match': isHelix ? storedEtag : '*', 'X-DA-Initiator': 'collab', ...bodyHeaders, }; @@ -392,7 +568,7 @@ export const persistence = { .map((con) => con.auth); if (auth.length > 0) { - if (isHelixDoc(ydoc.name)) { + if (isHelix) { // eslint-disable-next-line prefer-destructuring headers.Authorization = auth[0]; } else { @@ -407,9 +583,41 @@ export const persistence = { console.warn('[docroom] Writing back an empty document', ydoc.name, bodySize); } + // Track writes even when the backend omits an ETag or returns the same one. + // Any HEAD overlapping a write is too old to invalidate this session. + // eslint-disable-next-line no-param-reassign + ydoc.saveGeneration = (ydoc.saveGeneration ?? 0) + 1; + // eslint-disable-next-line no-param-reassign + ydoc.pendingPuts = (ydoc.pendingPuts ?? 0) + 1; + let response; + try { + response = await backend.fetch(ydoc.name, opts); + } finally { + // eslint-disable-next-line no-param-reassign + ydoc.pendingPuts -= 1; + } const { - ok, status, statusText, body, - } = await backend.fetch(ydoc.name, opts); + ok, status, statusText, body, headers: respHeaders, + } = response; + + if (ok) { + const etag = respHeaders.get('etag'); + if (isHelix) { + // eslint-disable-next-line no-param-reassign + ydoc.etag = normalizeEtag(etag); + // eslint-disable-next-line no-param-reassign + ydoc.etagUnverified = !etag; + if (!etag) { + // eslint-disable-next-line no-console + console.warn('[docroom] Helix save response missing ETag; protected saves require a reload', ydoc.name); + } + } else if (etag) { + // Keep the current document ETag so the periodic HEAD check can detect + // out-of-band changes to da-admin. + // eslint-disable-next-line no-param-reassign + ydoc.etag = etag; + } + } if (body) { // tell CloudFlare to consider the request as completed @@ -423,6 +631,88 @@ export const persistence = { }; }, + /** + * Issue a HEAD request to the backend and compare the document's current ETag + * against the one stored on the ydoc (`ydoc.etag`). Ignore obsolete responses + * and retry failed requests on the next polling interval. + * @param {WSSharedDoc} ydoc - the ydoc holding the last known ETag. + */ + checkEtag: async (ydoc) => { + if (ydoc.isDestroyed || ydoc !== docs.get(ydoc.name) || ydoc.conns.size === 0 + || ydoc.etagCheckAbortController || ydoc.pendingPuts > 0) { + return; + } + const { etag, saveGeneration } = ydoc; + const controller = new AbortController(); + // eslint-disable-next-line no-param-reassign + ydoc.etagCheckAbortController = controller; + try { + const opts = { method: 'HEAD', signal: controller.signal }; + const auth = Array.from(ydoc.conns.keys()) + .map((con) => con.auth) + .filter(Boolean); + if (auth.length > 0) { + const authorization = isHelixDoc(ydoc.name, ydoc) ? auth[0] : [...new Set(auth)].join(','); + opts.headers = new Headers({ Authorization: authorization }); + } + const res = await getBackend(ydoc.name, ydoc).fetch(ydoc.name, opts); + if (res.body) { + res.body.cancel(); + } + if (ydoc.isDestroyed || controller.signal.aborted + || ydoc !== docs.get(ydoc.name) || ydoc.conns.size === 0 + || ydoc.etag !== etag || ydoc.saveGeneration !== saveGeneration) { + return; + } + if (res.status === 404 || res.status === 410) { + await invalidateFromAdmin(ydoc.name, true); + return; + } + if (!res.ok) { + const msg = `[docroom] Etag HEAD failed: ${res.status} - ${res.statusText}`; + if (res.status === 401) { + // eslint-disable-next-line no-console + console.warn(msg, ydoc.name); + } else if (res.status === 403) { + // eslint-disable-next-line no-console + console.log(msg, ydoc.name); + } else { + // eslint-disable-next-line no-console + console.error(msg, ydoc.name); + } + return; + } + const currentEtag = normalizeEtag(res.headers.get('etag')); + if (!currentEtag) { + // eslint-disable-next-line no-console + console.warn('[docroom] Etag HEAD response missing ETag', ydoc.name); + return; + } + const storedEtag = normalizeEtag(etag); + if (!storedEtag) { + // A HEAD-only baseline is useful for polling, but must not authorize + // writes to content loaded or saved without a matching validator. + // eslint-disable-next-line no-param-reassign + ydoc.etag = currentEtag; + // eslint-disable-next-line no-param-reassign + ydoc.etagUnverified = true; + return; + } + if (currentEtag !== storedEtag) { + // eslint-disable-next-line no-console + console.log('[docroom] Etag check', ydoc.name, `stored=${etag}`, `current=${currentEtag}`, 'same=false'); + await invalidateFromAdmin(ydoc.name, true); + } + } catch (err) { + if (!controller.signal.aborted || err !== controller.signal.reason) { + logError(err, '[docroom] Etag check failed', ydoc.name, err); + } + } finally { + // eslint-disable-next-line no-param-reassign + delete ydoc.etagCheckAbortController; + } + }, + /** * An update to the document has been received. Store it in da-admin. * @param {WSSharedDoc} ydoc - the ydoc that has been updated. @@ -452,6 +742,14 @@ export const persistence = { const { ok, status, statusText } = await persistence.put(ydoc, content); if (!ok) { + if (status === 412 && isHelixDoc(docName, ydoc)) { + const conflict = new Error('412 - Helix document changed externally; reloading the document'); + // eslint-disable-next-line no-console + console.warn('[docroom] Helix version conflict', docName, conflict.message); + showError(ydoc, conflict); + await invalidateFromAdmin(docName, true); + return current; + } if (status === 412) { // Document doesn't exist - clean up cached state if (ydoc.storage) { @@ -492,7 +790,7 @@ export const persistence = { } if (closeAll) { // We had an unauthorized from da-admin - lets reset the connections. - // Pass isReentrant=true so closeConn skips flushSave here; the outer + // Pass skipSave=true so closeConn skips flushSave here; the outer // save already handled (or failed to handle) persistence. for (const con of Array.from(ydoc.conns.keys())) { // eslint-disable-next-line no-await-in-loop @@ -522,7 +820,7 @@ export const persistence = { // Get document from da-admin (throws on error including 404) const timingBeforeDaAdminGet = Date.now(); - current = await persistence.get(docName, conn.auth, ydoc.daadmin); + current = await persistence.get(docName, conn.auth, ydoc); const timingDaAdminGetDuration = Date.now() - timingBeforeDaAdminGet; // Read the stored state from internal worker storage (errors are non-fatal) @@ -656,7 +954,7 @@ export const persistence = { if (saving) { return; } - if (!current || ydoc !== docs.get(docName)) { + if (!current || ydoc.discardPendingChanges || ydoc !== docs.get(docName)) { return; } saving = true; @@ -690,6 +988,19 @@ export const persistence = { debouncedSave.cancel(); }; + // Periodically (every 5s) verify the document's ETag still matches the one we + // last saw on get/put, logging the result. Cleared in WSSharedDoc.destroy. + // Only Helix-backed documents carry an ETag, so skip scheduling otherwise. + if (!ydoc.etagCheckInterval && isHelixDoc(docName, ydoc) + && !ydoc.isDestroyed && ydoc === docs.get(docName) && ydoc.conns.size > 0) { + // eslint-disable-next-line no-param-reassign + ydoc.etagCheckInterval = setInterval(() => { + if (ydoc === docs.get(docName)) { + persistence.checkEtag(ydoc); + } + }, 5000); + } + const timingMap = new Map(); timingMap.set('timingReadStateDuration', timingReadStateDuration); timingMap.set('timingDaAdminGetDuration', timingDaAdminGetDuration); @@ -753,6 +1064,9 @@ export class WSSharedDoc extends Y.Doc { } destroy() { + this.cancelSave?.(); + clearInterval(this.etagCheckInterval); + this.etagCheckAbortController?.abort(); super.destroy(); this.awareness.destroy(); } @@ -788,6 +1102,10 @@ export const getYDoc = async (docname, conn, env, storage, timingData, ctx, gc = // Store the service binding to da-admin which we receive through the environment in the doc doc.daadmin = env.daadmin; + // Carry the IS_HELIX backend override too, so backend resolution can honour it + // without holding a reference to the whole environment. Together with daadmin, + // this makes the ydoc usable directly as the `env` argument to getBackend. + doc.IS_HELIX = env.IS_HELIX; if (!doc.promise) { // The doc is not yet bound to the persistence layer, do so now. The promise will be resolved // when bound. @@ -895,32 +1213,6 @@ export const messageListener = async (conn, doc, message) => { } }; -/** - * Invalidate the worker storage for the document, which will ensure that when accessed - * the worker will fetch the latest version of the document from the da-admin. - * Invalidation is implemented by closing all client connections to the doc, which will - * cause it to be reinitialised when accessed. - * @param {string} docName - The name of the document - * @returns true if the document was found and invalidated, false otherwise. - */ -export const invalidateFromAdmin = async (docName) => { - // eslint-disable-next-line no-console - console.log('[worker] Invalidate from Admin received', docName); - const ydoc = docs.get(docName); - if (ydoc) { - // Await all closeConn calls so that flushSave + docs.delete complete before - // this function returns. Without this, an editor that reconnects during the - // flushSave window finds the stale ydoc still in `docs` and reuses it - // instead of creating a fresh one that loads the new da-admin content. - await Promise.all(Array.from(ydoc.conns.keys()).map((c) => closeConn(ydoc, c))); - return true; - } else { - // eslint-disable-next-line no-console - console.log('[worker] Document not found', docName); - } - return false; -}; - /** * Called when a new (Yjs) WebSocket connection is being established. * @param {WebSocket} conn - The WebSocket connection diff --git a/test/shareddoc.test.js b/test/shareddoc.test.js index 2e21cd0..0177609 100644 --- a/test/shareddoc.test.js +++ b/test/shareddoc.test.js @@ -25,7 +25,7 @@ import { invalidateFromAdmin, isExpectedPlatformEvent, messageFlushRequest, messageFlushResponse, messageListener, persistence, readState, safePutLastsync, setupWSConnection, setYDoc, - showError, storeState, updateHandler, WSSharedDoc, + showError, storeState, toHelixPath, updateHandler, WSSharedDoc, } from '../src/shareddoc.js'; function isSubArray(full, sub) { @@ -157,11 +157,13 @@ describe('Collab Test Suite', () => { assert.equal(opts.method, undefined); assert(opts.headers === undefined); return { - ok: true, text: async () => 'content', status: 200, statusText: 'OK', + ok: true, text: async () => 'content', status: 200, statusText: 'OK', headers: new Headers({ etag: 'W/"abc"' }), }; }; - const result = await persistence.get('foo', undefined, daadmin); + const ydoc = { daadmin }; + const result = await persistence.get('foo', undefined, ydoc); assert.equal(result, 'content'); + assert.equal(ydoc.etag, 'W/"abc"', 'ETag from the GET response must be stored on the ydoc'); }); it('Test persistence get auth', async () => { @@ -171,10 +173,10 @@ describe('Collab Test Suite', () => { assert.equal(opts.method, undefined); assert.equal(opts.headers.get('authorization'), 'auth'); return { - ok: true, text: async () => 'content', status: 200, statusText: 'OK', + ok: true, text: async () => 'content', status: 200, statusText: 'OK', headers: new Headers(), }; }; - const result = await persistence.get('foo', 'auth', daadmin); + const result = await persistence.get('foo', 'auth', { daadmin }); assert.equal(result, 'content'); }); @@ -189,7 +191,7 @@ describe('Collab Test Suite', () => { }; }; try { - await persistence.get('foo', 'auth', daadmin); + await persistence.get('foo', 'auth', { daadmin }); assert.fail('Should have thrown an error'); } catch (error) { assert(error.toString().includes('unable to get resource - status: 404')); @@ -208,7 +210,7 @@ describe('Collab Test Suite', () => { }; }; try { - await persistence.get('foo', 'auth', daadmin); + await persistence.get('foo', 'auth', { daadmin }); assert.fail('Expected get to throw'); } catch (error) { // expected @@ -222,7 +224,7 @@ describe('Collab Test Suite', () => { fetch: async () => ({ ok: false, status: 401, statusText: 'Unauthorized' }), }; try { - await persistence.get('foo', 'auth', daadmin); + await persistence.get('foo', 'auth', { daadmin }); assert.fail('Should have thrown an error'); } catch (error) { assert.equal(401, error.status, 'Error must carry 401 status'); @@ -235,7 +237,7 @@ describe('Collab Test Suite', () => { fetch: async () => ({ ok: false, status: 403, statusText: 'Forbidden' }), }; try { - await persistence.get('foo', 'auth', daadmin); + await persistence.get('foo', 'auth', { daadmin }); assert.fail('Should have thrown an error'); } catch (error) { assert.equal(403, error.status, 'Error must carry 403 status'); @@ -255,7 +257,7 @@ describe('Collab Test Suite', () => { console.log = (...a) => logged.push(['log', ...a]); console.error = (...a) => logged.push(['error', ...a]); try { - await persistence.get('foo', 'auth', daadmin); + await persistence.get('foo', 'auth', { daadmin }); assert.fail('Should have thrown'); } catch (_) { // expected @@ -288,7 +290,9 @@ describe('Collab Test Suite', () => { assert.equal(opts.method, 'PUT'); assert.equal(opts.headers.get('If-Match'), '*', 'Should include If-Match: * header'); assert.equal(await opts.body.get('data').text(), 'test'); - return { ok: true, status: 200, statusText: 'OK - Stored' }; + return { + ok: true, status: 200, statusText: 'OK - Stored', headers: new Headers(), + }; }; const conns = new Map(); // conns.set({}, new Set()); @@ -307,7 +311,9 @@ describe('Collab Test Suite', () => { assert.equal('collab', opts.headers.get('X-DA-Initiator')); assert.equal('*', opts.headers.get('If-Match')); assert.equal(await opts.body.get('data').text(), 'test'); - return { ok: true, status: 200, statusText: 'OK - Stored too' }; + return { + ok: true, status: 200, statusText: 'OK - Stored too', headers: new Headers(), + }; }; const conns = new Map(); conns.set({ auth: 'myauth' }, new Set()); @@ -324,7 +330,9 @@ describe('Collab Test Suite', () => { assert.equal(opts.method, 'PUT'); assert.equal(opts.headers.get('If-Match'), '*'); assert.equal(await opts.body.get('data').text(), 'test'); - return { ok: true, status: 200, statusText: 'OK' }; + return { + ok: true, status: 200, statusText: 'OK', headers: new Headers(), + }; }; const result = await persistence.put({ name: 'foo', conns: new Map(), daadmin }, 'test'); assert(result.ok); @@ -339,7 +347,9 @@ describe('Collab Test Suite', () => { assert.equal(opts.headers.get('X-DA-Initiator'), 'collab'); assert.equal(opts.headers.get('If-Match'), '*'); assert.equal(await opts.body.get('data').text(), 'test'); - return { ok: true, status: 200, statusText: 'okidoki' }; + return { + ok: true, status: 200, statusText: 'okidoki', headers: new Headers(), + }; }; const result = await persistence.put({ name: 'foo', @@ -410,7 +420,9 @@ describe('Collab Test Suite', () => { assert.equal(ydoc, mockYDoc); assert.equal(content, 'Svr content update'); called = true; - return { ok: true, status: 201, statusText: 'Created' }; + return { + ok: true, status: 201, statusText: 'Created', headers: new Headers(), + }; }; let calledCloseCon = false; @@ -579,7 +591,7 @@ describe('Collab Test Suite', () => { let putCalled = false; pss.persistence.put = async () => { putCalled = true; - return { ok: true, status: 200 }; + return { ok: true, status: 200, headers: new Headers() }; }; const real = '

real customer content

'; @@ -608,7 +620,7 @@ describe('Collab Test Suite', () => { const putCalls = []; pss.persistence.put = async (yd, c) => { putCalls.push(c); - return { ok: true, status: 200 }; + return { ok: true, status: 200, headers: new Headers() }; }; const result = await pss.persistence.update(mockYDoc, '

old

', 'test.html'); @@ -827,7 +839,7 @@ describe('Collab Test Suite', () => { const putCalls = []; pss.persistence.put = async () => { putCalls.push('put'); - return { ok: true }; + return { ok: true, headers: new Headers() }; }; // Simulate another update after 412 @@ -1025,7 +1037,7 @@ describe('Collab Test Suite', () => { pss.persistence.get = async () => unchangedContent; pss.persistence.put = async () => { putCalls.push('put'); - return { ok: true, status: 200 }; + return { ok: true, status: 200, headers: new Headers() }; }; const conn = { auth: undefined, close() {} }; @@ -1066,7 +1078,7 @@ describe('Collab Test Suite', () => { pss.persistence.get = async () => '

initial content

'; pss.persistence.put = async (doc, content) => { putCalls.push(content); - return { ok: true, status: 200 }; + return { ok: true, status: 200, headers: new Headers() }; }; const conn = { auth: undefined, close() {} }; @@ -1143,7 +1155,7 @@ describe('Collab Test Suite', () => { await new Promise((res) => { resolvePut = res; }); - return { ok: true, status: 200 }; + return { ok: true, status: 200, headers: new Headers() }; }; const conn = { auth: undefined, close() {} }; @@ -1232,7 +1244,7 @@ describe('Collab Test Suite', () => { const mockStorage = { list: () => new Map(), get: async () => undefined }; - pss.persistence.get = async (nm, au, ad) => `Get: ${nm}-${au}-${ad}`; + pss.persistence.get = async (nm, au, yd) => `Get: ${nm}-${au}-${yd.daadmin}`; const updated = new Map(); pss.persistence.update = async (d, v) => updated.set(d, v); @@ -1269,7 +1281,7 @@ describe('Collab Test Suite', () => { const mockStorage = { list: () => new Map(), get: async () => undefined }; - pss.persistence.get = async (nm, au, ad) => `Get: ${nm}-${au}-${ad}`; + pss.persistence.get = async (nm, au, yd) => `Get: ${nm}-${au}-${yd.daadmin}`; const updated = new Map(); pss.persistence.update = async (d, v) => updated.set(d, v); @@ -1305,7 +1317,7 @@ describe('Collab Test Suite', () => { pss.setYDoc(docName, testYDoc); const mockStorage = { list: () => new Map(), get: async () => undefined }; - pss.persistence.get = async (nm, au, ad) => `Get: ${nm}-${au}-${ad}`; + pss.persistence.get = async (nm, au, yd) => `Get: ${nm}-${au}-${yd.daadmin}`; pss.persistence.update = async () => {}; await pss.persistence.bindState(docName, testYDoc, mockConn, mockStorage); @@ -1343,7 +1355,7 @@ describe('Collab Test Suite', () => { pss.setYDoc(docName, testYDoc); const mockStorage = { list: () => new Map(), get: async () => undefined }; - pss.persistence.get = async (nm, au, ad) => `Get: ${nm}-${au}-${ad}`; + pss.persistence.get = async (nm, au, yd) => `Get: ${nm}-${au}-${yd.daadmin}`; pss.persistence.update = async () => {}; await pss.persistence.bindState(docName, testYDoc, mockConn, mockStorage); @@ -1378,7 +1390,7 @@ describe('Collab Test Suite', () => { pss.setYDoc(docName, testYDoc); const mockStorage = { list: () => new Map(), get: async () => undefined }; - pss.persistence.get = async (nm, au, ad) => `Get: ${nm}-${au}-${ad}`; + pss.persistence.get = async (nm, au, yd) => `Get: ${nm}-${au}-${yd.daadmin}`; pss.persistence.update = async () => {}; const logged = []; @@ -1569,7 +1581,7 @@ describe('Collab Test Suite', () => { pss.persistence.put = async (yd, c) => { if (yd === ydoc && c.includes('newcontent')) { putCalls.push(c); - return { ok: true, status: 200 }; + return { ok: true, status: 200, headers: new Headers() }; } }; @@ -1634,7 +1646,7 @@ describe('Collab Test Suite', () => { savedSetTimeout(resolve, 30); }); concurrentPuts -= 1; - return { ok: true, status: 200 }; + return { ok: true, status: 200, headers: new Headers() }; }; await pss.persistence.bindState(docName, ydoc, {}, storage); @@ -2219,7 +2231,13 @@ describe('Collab Test Suite', () => { } return map.get(keyOrKeys); }, - async put(d) { Object.entries(d).forEach(([k, v]) => map.set(k, v)); }, + async put(d, value) { + if (typeof d === 'string') { + map.set(d, value); + } else { + Object.entries(d).forEach(([k, v]) => map.set(k, v)); + } + }, async delete(keys) { (Array.isArray(keys) ? keys : [keys]).forEach((k) => map.delete(k)); }, @@ -3037,7 +3055,9 @@ describe('Collab Test Suite', () => { await new Promise((resolve) => { savedSetTimeout(resolve, 50); }); - return { ok: true, status: 200, statusText: 'OK' }; + return { + ok: true, status: 200, statusText: 'OK', headers: new Headers(), + }; }; await pss.persistence.bindState(docName, ydoc, {}, storage); @@ -3056,9 +3076,8 @@ describe('Collab Test Suite', () => { // --------------------------------------------------------------------------- // Backend resolution (api-live-switch branch) // - // The storage backend is determined entirely by the doc URL: docs under - // https://api.aem.live live in Helix (global fetch); everything else goes - // through the da-admin service binding. There is no isHelix flag to thread. + // Normal routing follows the doc URL. The manual local-testing override can + // force backend selection and, for Helix, rewrite the request URL. // --------------------------------------------------------------------------- it('getBackend routes da-admin docs through the daadmin binding', async () => { @@ -3069,7 +3088,7 @@ describe('Collab Test Suite', () => { return 'da-resp'; }, }; - const backend = getBackend('https://admin.da.live/x.html', daadmin); + const backend = getBackend('https://admin.da.live/x.html', { daadmin }); const resp = await backend.fetch('https://admin.da.live/x.html', { method: 'HEAD' }); assert.equal('da-resp', resp); @@ -3088,7 +3107,7 @@ describe('Collab Test Suite', () => { const daadmin = { fetch: async () => { assert.fail('daadmin.fetch must not be called for Helix docs'); }, }; - const backend = getBackend('https://api.aem.live/o/r/p.html', daadmin); + const backend = getBackend('https://api.aem.live/o/r/p.html', { daadmin }); const resp = await backend.fetch('https://api.aem.live/o/r/p.html', { method: 'HEAD' }); assert.equal('helix-resp', resp); @@ -3099,6 +3118,112 @@ describe('Collab Test Suite', () => { } }); + it('toHelixPath converts da-admin source paths and leaves Helix paths untouched', () => { + // da-admin /source/{org}/{site}/{rest} -> /{org}/sites/{site}/source/{rest} + assert.equal(toHelixPath('/source/o/r/p.html'), '/o/sites/r/source/p.html'); + assert.equal(toHelixPath('/source/o/r/a/b/c.html'), '/o/sites/r/source/a/b/c.html'); + // Already in Helix format -> unchanged. + assert.equal(toHelixPath('/o/sites/r/source/p.html'), '/o/sites/r/source/p.html'); + // Not a da-admin source path -> unchanged. + assert.equal(toHelixPath('/something/else.html'), '/something/else.html'); + }); + + it('getBackend rewrites the URL to api.aem.live when IS_HELIX=true', async () => { + const savedFetch = globalThis.fetch; + const calls = []; + globalThis.fetch = async (url, opts) => { + calls.push({ url, opts }); + return 'helix-resp'; + }; + try { + const backend = getBackend('https://admin.da.live/source/o/r/p.html', { IS_HELIX: 'true' }); + const resp = await backend.fetch('https://admin.da.live/source/o/r/p.html?x=1', { method: 'HEAD' }); + + assert.equal('helix-resp', resp); + assert.equal(1, calls.length); + assert.equal('https://api.aem.live/o/sites/r/source/p.html?x=1', calls[0].url, 'host swapped to api.aem.live and path converted to Helix format, query preserved'); + } finally { + globalThis.fetch = savedFetch; + } + }); + + it('getBackend rewrites the URL to localhost:3000 when IS_HELIX=local', async () => { + const savedFetch = globalThis.fetch; + const calls = []; + globalThis.fetch = async (url, opts) => { + calls.push({ url, opts }); + return 'local-resp'; + }; + try { + const backend = getBackend('https://admin.da.live/source/o/r/p.html', { IS_HELIX: 'local' }); + const resp = await backend.fetch('https://admin.da.live/source/o/r/p.html?x=1', { method: 'HEAD' }); + + assert.equal('local-resp', resp); + assert.equal(1, calls.length); + assert.equal('http://localhost:3000/o/sites/r/source/p.html?x=1', calls[0].url, 'host swapped to localhost:3000 and path converted to Helix format, query preserved'); + } finally { + globalThis.fetch = savedFetch; + } + }); + + it('getBackend with IS_HELIX=false routes to da-admin without rewriting the URL', async () => { + const calls = []; + const daadmin = { + fetch: async (url, opts) => { + calls.push({ url, opts }); + return 'da-resp'; + }, + }; + const savedFetch = globalThis.fetch; + globalThis.fetch = async () => { + assert.fail('global fetch must not be used when IS_HELIX=false'); + }; + try { + const backend = getBackend('https://api.aem.live/o/r/p.html', { IS_HELIX: 'false', daadmin }); + const resp = await backend.fetch('https://api.aem.live/o/r/p.html?x=1', { method: 'HEAD' }); + + assert.equal('da-resp', resp); + assert.equal(1, calls.length); + assert.equal('https://api.aem.live/o/r/p.html?x=1', calls[0].url, 'URL left untouched for da-admin'); + } finally { + globalThis.fetch = savedFetch; + } + }); + + for (const mode of [undefined, 'false', 'true', 'local']) { + it(`getBackend emits the local-testing rewrite log only for forced Helix (IS_HELIX=${mode})`, async () => { + const savedFetch = globalThis.fetch; + const savedLog = console.log; + const calls = []; + const logged = []; + const fetch = async (url, opts) => { + calls.push({ url, opts }); + return new Response(null); + }; + globalThis.fetch = fetch; + console.log = (...args) => logged.push(args); + const url = 'https://api.aem.live/o/sites/r/source/p.html?x=1'; + const env = { daadmin: { fetch }, IS_HELIX: mode }; + const opts = { method: 'HEAD', headers: new Headers({ 'X-test': 'value' }) }; + const expectedUrl = mode === 'local' + ? 'http://localhost:3000/o/sites/r/source/p.html?x=1' + : url; + try { + await getBackend(url, env).fetch(url, opts); + + assert.equal(calls.length, 1); + assert.equal(calls[0].url, expectedUrl); + assert.strictEqual(calls[0].opts, opts); + assert.deepStrictEqual(logged, mode === 'true' || mode === 'local' + ? [['*** Calling', expectedUrl]] + : []); + } finally { + globalThis.fetch = savedFetch; + console.log = savedLog; + } + }); + } + it('getBackend.putReqData builds multipart form-data for da-admin docs', () => { const backend = getBackend('https://admin.da.live/x.html', {}); const { body, size, headers } = backend.putReqData('hello world', 'text/html'); @@ -3123,13 +3248,25 @@ describe('Collab Test Suite', () => { assert.equal(isHelixDoc('http://localhost:8080/x.html'), false); }); + it('isHelixDoc IS_HELIX env overrides the URL heuristic', () => { + // IS_HELIX=true forces every doc onto Helix, even a da-admin URL. + assert.equal(isHelixDoc('https://admin.da.live/x.html', { IS_HELIX: 'true' }), true); + // IS_HELIX=local also routes through Helix (rewritten to localhost downstream). + assert.equal(isHelixDoc('https://admin.da.live/x.html', { IS_HELIX: 'local' }), true); + // IS_HELIX=false forces every doc onto da-admin, even an api.aem.live URL. + assert.equal(isHelixDoc('https://api.aem.live/o/r/p.html', { IS_HELIX: 'false' }), false); + // An env without IS_HELIX falls back to the URL heuristic. + assert.equal(isHelixDoc('https://api.aem.live/o/r/p.html', {}), true); + assert.equal(isHelixDoc('https://admin.da.live/x.html', {}), false); + }); + it('persistence.get routes to the global fetch for an api.aem.live doc', async () => { const savedFetch = globalThis.fetch; const calls = []; globalThis.fetch = async (url, opts) => { calls.push({ url, opts }); return { - ok: true, text: async () => 'helix content', status: 200, statusText: 'OK', + ok: true, text: async () => 'helix content', status: 200, statusText: 'OK', headers: new Headers({ etag: '"initial"' }), }; }; try { @@ -3139,7 +3276,7 @@ describe('Collab Test Suite', () => { const result = await persistence.get( 'https://api.aem.live/owner/repo/page.html', 'Bearer t', - daadmin, + { daadmin }, ); assert.equal(result, 'helix content'); assert.equal(1, calls.length); @@ -3155,13 +3292,16 @@ describe('Collab Test Suite', () => { const calls = []; globalThis.fetch = async (url, opts) => { calls.push({ url, opts }); - return { ok: true, status: 200, statusText: 'OK' }; + return { + ok: true, status: 200, statusText: 'OK', headers: new Headers({ etag: '"saved"' }), + }; }; try { const conns = new Map(); conns.set({ auth: 'Bearer abc' }, new Set()); const ydoc = { name: 'https://api.aem.live/owner/repo/page.html', + etag: 'W/"initial"', conns, daadmin: { fetch: async () => { assert.fail('daadmin.fetch must not be called for Helix docs'); }, @@ -3177,7 +3317,7 @@ describe('Collab Test Suite', () => { assert.equal(opts.method, 'POST'); assert.strictEqual(opts.body, body, 'Helix POST body must be the raw content string, not FormData'); assert.equal(opts.headers.get('Content-Type'), 'text/html'); - assert.equal(opts.headers.get('If-Match'), '*'); + assert.equal(opts.headers.get('If-Match'), '"initial"'); assert.equal(opts.headers.get('X-DA-Initiator'), 'collab'); assert.equal(opts.headers.get('Authorization'), 'Bearer abc'); } finally { @@ -3190,7 +3330,9 @@ describe('Collab Test Suite', () => { const calls = []; globalThis.fetch = async (url, opts) => { calls.push({ url, opts }); - return { ok: true, status: 200, statusText: 'OK' }; + return { + ok: true, status: 200, statusText: 'OK', headers: new Headers({ etag: '"saved"' }), + }; }; try { const conns = new Map(); @@ -3198,6 +3340,7 @@ describe('Collab Test Suite', () => { conns.set({ auth: 'Bearer xyz' }, new Set()); const ydoc = { name: 'https://api.aem.live/owner/repo/page.html', + etag: '"initial"', conns, daadmin: { fetch: async () => { assert.fail('daadmin.fetch must not be called for Helix docs'); }, @@ -3224,13 +3367,16 @@ describe('Collab Test Suite', () => { let captured; globalThis.fetch = async (url, opts) => { captured = opts; - return { ok: true, status: 200, statusText: 'OK' }; + return { + ok: true, status: 200, statusText: 'OK', headers: new Headers({ etag: '"saved"' }), + }; }; try { const conns = new Map(); conns.set({ auth: 'a' }, new Set()); const ydoc = { name: 'https://api.aem.live/o/r/d.json', + etag: '"initial"', conns, daadmin: {}, }; @@ -3248,7 +3394,9 @@ describe('Collab Test Suite', () => { const daadmin = { fetch: async (url, opts) => { calls.push({ url, opts }); - return { ok: true, status: 200, statusText: 'OK' }; + return { + ok: true, status: 200, statusText: 'OK', headers: new Headers(), + }; }, }; const conns = new Map(); @@ -3271,26 +3419,1023 @@ describe('Collab Test Suite', () => { ); }); + for (const mode of ['true', 'local', 'false']) { + it(`persistence.put honours IS_HELIX=${mode} for method, body and auth`, async () => { + const savedFetch = globalThis.fetch; + const calls = []; + const fetch = async (url, opts) => { + calls.push({ url, opts }); + return new Response(null, { headers: { etag: '"saved"' } }); + }; + globalThis.fetch = fetch; + const isHelix = mode !== 'false'; + const name = isHelix + ? 'https://admin.da.live/source/o/r/p.html' + : 'https://api.aem.live/o/sites/r/source/p.html'; + const conns = new Map([ + [{ auth: 'auth-a' }, new Set()], + [{ auth: 'auth-b' }, new Set()], + ]); + const ydoc = { + name, conns, IS_HELIX: mode, daadmin: { fetch }, etag: 'W/"initial"', + }; + const body = '

content long enough to avoid the empty-stub warning padding padding

'; + try { + const result = await persistence.put(ydoc, body); + + assert(result.ok); + assert.equal(calls.length, 1); + const { url, opts } = calls[0]; + const origin = mode === 'local' ? 'http://localhost:3000' : 'https://api.aem.live'; + assert.equal(url, isHelix ? `${origin}/o/sites/r/source/p.html` : name); + assert.equal(opts.method, isHelix ? 'POST' : 'PUT'); + assert.equal(opts.headers.get('If-Match'), isHelix ? '"initial"' : '*'); + assert.equal(opts.headers.get('Authorization'), isHelix ? 'auth-a' : 'auth-a,auth-b'); + if (isHelix) { + assert.strictEqual(opts.body, body); + assert.equal(opts.headers.get('Content-Type'), 'text/html'); + } else { + assert(opts.body instanceof FormData); + assert.equal(await opts.body.get('data').text(), body); + assert.equal(opts.headers.get('Content-Type'), null); + } + assert.equal(ydoc.etag, '"saved"'); + } finally { + globalThis.fetch = savedFetch; + } + }); + } + + for (const mode of [undefined, 'true', 'local', 'false']) { + it(`persistence.checkEtag uses backend-appropriate auth with IS_HELIX=${mode}`, async () => { + const savedFetch = globalThis.fetch; + const calls = []; + const fetch = async (url, opts) => { + calls.push({ url, opts }); + return new Response(null, { headers: { etag: '"current"' } }); + }; + globalThis.fetch = fetch; + const name = mode === 'true' || mode === 'local' + ? 'https://admin.da.live/source/o/r/p.html' + : 'https://api.aem.live/o/sites/r/source/p.html'; + const conns = new Map([ + [{ auth: 'auth-a' }, new Set()], + [{ auth: 'auth-b' }, new Set()], + ]); + const ydoc = { + name, conns, IS_HELIX: mode, daadmin: { fetch }, etag: '"current"', + }; + const docs = setYDoc(name, ydoc); + try { + await persistence.checkEtag(ydoc); + + assert.equal(calls.length, 1); + const { url, opts } = calls[0]; + const origin = mode === 'local' ? 'http://localhost:3000' : 'https://api.aem.live'; + assert.equal(url, mode === 'false' ? name : `${origin}/o/sites/r/source/p.html`); + assert.equal(opts.method, 'HEAD'); + assert.equal(opts.headers.get('Authorization'), mode === 'false' ? 'auth-a,auth-b' : 'auth-a'); + } finally { + globalThis.fetch = savedFetch; + docs.delete(name); + } + }); + } + + describe('Helix ETag polling', () => { + let savedFetch; + const opened = []; + + function makeDoc(slug, etag = '"old"', type = 'html') { + const name = `https://api.aem.live/o/sites/r/source/${slug}.${type}`; + const ydoc = new WSSharedDoc(name); + ydoc.etag = etag; + const conn = { + auth: 'auth-a', + readyState: 1, + closed: false, + send() {}, + close() { this.closed = true; }, + }; + ydoc.conns.set(conn, new Set()); + const docs = setYDoc(name, ydoc); + opened.push({ ydoc, docs }); + return { ydoc, conn, docs }; + } + + function deferred() { + let resolve; + let reject; + const promise = new Promise((r, j) => { + resolve = r; + reject = j; + }); + return { promise, resolve, reject }; + } + + beforeEach(() => { + savedFetch = globalThis.fetch; + }); + + afterEach(() => { + globalThis.fetch = savedFetch; + opened.splice(0).forEach(({ ydoc, docs }) => { + ydoc.conns.clear(); + ydoc.destroy(); + if (docs.get(ydoc.name) === ydoc) { + docs.delete(ydoc.name); + } + }); + }); + + for (const [stored, current] of [['W/"old"', '"old"'], ['"old"', 'W/"old"']]) { + it(`equivalent weak/strong ETags (${stored}, ${current}) do not invalidate`, async () => { + const { ydoc, conn } = makeDoc('equivalent-etag', stored); + globalThis.fetch = async () => new Response(null, { headers: { etag: current } }); + + await persistence.checkEtag(ydoc); + + assert.equal(conn.closed, false); + }); + } + + it('a missing stored ETag adopts a polling baseline without authorizing stale writes', async () => { + const { ydoc, conn } = makeDoc('unknown-baseline', null); + let calls = 0; + globalThis.fetch = async () => { + calls += 1; + return new Response(null, { headers: { etag: 'W/"observed"' } }); + }; + + await persistence.checkEtag(ydoc); + + assert.equal(conn.closed, false); + assert.equal(ydoc.etag, '"observed"'); + assert.equal(ydoc.etagUnverified, true); + await assert.rejects( + persistence.put(ydoc, '

Pending local edit

'), + /without a verified ETag/, + ); + assert.equal(calls, 1, 'An unverified baseline must not send a POST'); + }); + + for (const type of ['html', 'json']) { + it(`a ${type} GET without ETag establishes a matching fresh content/version baseline`, async () => { + const { ydoc, conn } = makeDoc(`missing-get-etag-${type}`, null, type); + const stale = type === 'json' ? '{"value":"stale"}' : '

Stale content

'; + const fresh = type === 'json' ? '{"value":"current"}' : '

Current content

'; + const responses = [ + new Response(stale), + new Response(null, { headers: { etag: '"current"' } }), + new Response(fresh), + new Response(null, { headers: { etag: 'W/"current"' } }), + ]; + const requests = []; + globalThis.fetch = async (_url, opts) => { + requests.push(opts.method ?? 'GET'); + assert.equal(opts.headers.get('Authorization'), conn.auth); + return responses.shift(); + }; + + const content = await persistence.get(ydoc.name, conn.auth, ydoc); + + assert.deepStrictEqual(content, type === 'json' ? { value: 'current' } : fresh); + assert.equal(ydoc.etag, '"current"'); + assert.equal(ydoc.etagUnverified, false); + assert.deepStrictEqual(requests, ['GET', 'HEAD', 'GET', 'HEAD']); + globalThis.fetch = async () => new Response(null, { headers: { etag: '"current"' } }); + await persistence.checkEtag(ydoc); + assert.equal(conn.closed, false); + }); + } + + it('a refreshed GET with an ETag supplies its own authoritative baseline', async () => { + const { ydoc, conn } = makeDoc('refreshed-get-etag', null); + const responses = [ + new Response('stale'), + new Response(null, { headers: { etag: '"before"' } }), + new Response('fresh', { headers: { etag: 'W/"fresh"' } }), + ]; + globalThis.fetch = async () => responses.shift(); + + assert.equal(await persistence.get(ydoc.name, conn.auth, ydoc), 'fresh'); + assert.equal(ydoc.etag, '"fresh"'); + assert.equal(ydoc.etagUnverified, false); + assert.equal(responses.length, 0); + }); + + it('retries a missing-GET-tag baseline when the backend changes during the fresh read', async () => { + const { ydoc, conn } = makeDoc('baseline-race', null); + const responses = [ + new Response('stale'), + new Response(null, { headers: { etag: '"first"' } }), + new Response('first'), + new Response(null, { headers: { etag: '"second"' } }), + new Response(null, { headers: { etag: '"second"' } }), + new Response('second'), + new Response(null, { headers: { etag: '"second"' } }), + ]; + globalThis.fetch = async () => responses.shift(); + + assert.equal(await persistence.get(ydoc.name, conn.auth, ydoc), 'second'); + assert.equal(ydoc.etag, '"second"'); + assert.equal(ydoc.etagUnverified, false); + assert.equal(responses.length, 0); + }); + + it('an unstable baseline is bounded, logged, and cannot authorize a save', async () => { + const { ydoc, conn } = makeDoc('unstable-baseline', null); + const responses = [ + new Response('stale'), + new Response(null, { headers: { etag: '"first"' } }), + new Response('first'), + new Response(null, { headers: { etag: '"second"' } }), + new Response(null, { headers: { etag: '"second"' } }), + new Response('second'), + new Response(null, { headers: { etag: '"third"' } }), + ]; + const warnings = []; + const savedWarn = console.warn; + console.warn = (...args) => warnings.push(args); + globalThis.fetch = async () => responses.shift(); + try { + assert.equal(await persistence.get(ydoc.name, conn.auth, ydoc), 'second'); + assert.equal(ydoc.etag, null); + assert.equal(ydoc.etagUnverified, true); + assert.equal(responses.length, 0); + assert.equal(warnings.length, 1); + globalThis.fetch = async () => assert.fail('No unconditional POST is allowed'); + await assert.rejects( + persistence.put(ydoc, '

Pending edit

'), + /without a verified ETag/, + ); + assert.equal(conn.closed, false); + } finally { + console.warn = savedWarn; + } + }); + + it('a missing HEAD tag during baseline establishment keeps the read available and blocks saves', async () => { + const { ydoc, conn } = makeDoc('no-tags', null); + const responses = [new Response('content'), new Response(null)]; + const savedWarn = console.warn; + const warnings = []; + console.warn = (...args) => warnings.push(args); + globalThis.fetch = async () => responses.shift(); + try { + assert.equal(await persistence.get(ydoc.name, conn.auth, ydoc), 'content'); + assert.equal(ydoc.etag, null); + assert.equal(ydoc.etagUnverified, true); + assert.equal(warnings.length, 1); + assert.equal(conn.closed, false); + } finally { + console.warn = savedWarn; + } + }); + + for (const stage of ['before', 'reload', 'after']) { + it(`baseline ${stage} failures preserve normal HTTP error reporting`, async () => { + const { ydoc, conn } = makeDoc(`baseline-failure-${stage}`, null); + const responses = [new Response('stale')]; + if (stage !== 'before') { + responses.push(new Response(null, { headers: { etag: '"initial"' } })); + } + if (stage === 'after') { + responses.push(new Response('reloaded')); + } + responses.push(new Response(null, { status: 403 })); + globalThis.fetch = async () => responses.shift(); + + await assert.rejects( + persistence.get(ydoc.name, conn.auth, ydoc), + (error) => error.status === 403, + ); + assert.equal(responses.length, 0); + }); + } + + it('a baseline fetch error cancels the unread document body and propagates the failure', async () => { + const { ydoc, conn } = makeDoc('baseline-network-failure', null); + let cancelled = false; + let calls = 0; + globalThis.fetch = async () => { + calls += 1; + if (calls === 1) { + return { + ok: true, + headers: new Headers(), + body: { cancel() { cancelled = true; } }, + }; + } + throw new Error('Baseline fetch failed'); + }; + + await assert.rejects(persistence.get(ydoc.name, conn.auth, ydoc), /Baseline fetch failed/); + assert(cancelled); + assert.equal(calls, 2); + }); + + it('successful Helix saves advance the normalized ETag and do not invalidate on the next poll', async () => { + const { ydoc, conn } = makeDoc('successful-versioned-save', 'W/"old"'); + const requests = []; + globalThis.fetch = async (_url, opts) => { + requests.push(opts.method); + if (opts.method === 'POST') { + assert.equal(opts.headers.get('If-Match'), '"old"'); + return new Response(null, { status: 201, headers: { etag: 'W/"saved"' } }); + } + return new Response(null, { headers: { etag: '"saved"' } }); + }; + + await persistence.put(ydoc, '

Content long enough to avoid the empty document warning when saved

'); + await persistence.checkEtag(ydoc); + + assert.equal(ydoc.etag, '"saved"'); + assert.equal(ydoc.etagUnverified, false); + assert.equal(conn.closed, false); + assert.deepStrictEqual(requests, ['POST', 'HEAD']); + }); + + it('an unexpected save response without ETag does not retain the old version or disconnect on the next poll', async () => { + const { ydoc, conn } = makeDoc('missing-save-etag'); + let posts = 0; + const warnings = []; + const savedWarn = console.warn; + console.warn = (...args) => warnings.push(args); + globalThis.fetch = async (_url, opts) => { + if (opts.method === 'POST') { + posts += 1; + return new Response(null, { status: 201 }); + } + return new Response(null, { headers: { etag: '"saved"' } }); + }; + try { + const content = '

Content long enough to avoid the empty document warning when saved

'; + await persistence.put(ydoc, content); + assert.equal(ydoc.etag, null); + assert.equal(ydoc.etagUnverified, true); + await persistence.checkEtag(ydoc); + + assert.equal(ydoc.etag, '"saved"'); + assert.equal(conn.closed, false); + assert.equal(warnings.length, 1); + await assert.rejects(persistence.put(ydoc, content), /without a verified ETag/); + assert.equal(posts, 1); + } finally { + console.warn = savedWarn; + } + }); + + for (const etag of [undefined, null, '*']) { + it(`a Helix save cannot fall back to unconditional If-Match for ETag ${etag}`, async () => { + const { ydoc } = makeDoc('no-write-validator'); + ydoc.etag = etag; + globalThis.fetch = async () => assert.fail('Do not send an unprotected Helix save'); + + await assert.rejects( + persistence.put(ydoc, '

Pending edit

'), + /without a verified ETag/, + ); + }); + } + + it('read-only Helix sessions still skip saves without requiring a write validator', async () => { + const { ydoc, conn } = makeDoc('readonly-no-validator', null); + conn.readOnly = true; + globalThis.fetch = async () => assert.fail('A read-only session must not POST'); + + assert.deepStrictEqual( + await persistence.put(ydoc, '

Read-only content

'), + { ok: true }, + ); + }); + + it('a version conflict before the next poll preserves external content and invalidates without another save', async () => { + const { ydoc, conn, docs } = makeDoc('version-conflict'); + aem2doc('

Original source content

', ydoc); + const original = doc2aem(ydoc); + const storage = makeStorage({ + doc: ydoc.name, + docstore: Y.encodeStateAsUpdate(ydoc), + lastsync: original, + }); + let remoteContent = original; + let remoteEtag = '"old"'; + let posts = 0; + globalThis.fetch = async (_url, opts) => { + if (opts.method === 'POST') { + posts += 1; + assert.equal(opts.headers.get('If-Match'), '"old"'); + assert.notEqual(opts.headers.get('If-Match'), remoteEtag); + return new Response(null, { status: 412 }); + } + return new Response(remoteContent, { headers: { etag: remoteEtag } }); + }; + await persistence.bindState(ydoc.name, ydoc, conn, storage); + ydoc.hasClientChanged = true; + const fragment = ydoc.getXmlFragment('prosemirror'); + fragment.delete(0, fragment.length); + aem2doc('

Unsaved client changes

', ydoc); + remoteContent = '

Externally replaced source

'; + remoteEtag = '"external"'; + const warnings = []; + const savedWarn = console.warn; + console.warn = (...args) => warnings.push(args); + try { + await ydoc.flushSave(); + await ydoc.flushSave(); + + assert.equal(posts, 1); + assert.equal(remoteContent, '

Externally replaced source

'); + assert(conn.closed); + assert.equal(docs.get(ydoc.name), undefined); + assert.equal(await storage.get('lastsync'), undefined); + assert.equal(ydoc.discardPendingChanges, true); + assert(ydoc.getMap('error').get('message').includes('changed externally')); + assert.equal(warnings.length, 1); + assert(warnings[0][0].includes('version conflict')); + } finally { + console.warn = savedWarn; + } + }); + + it('changed ETags cancel saves and finish closing all connections without flushing', async () => { + const { ydoc, conn, docs } = makeDoc('changed'); + const secondConn = { + closed: false, + close() { this.closed = true; }, + }; + ydoc.conns.set(secondConn, new Set()); + let cancelled = false; + ydoc.cancelSave = () => { + cancelled = true; + }; + ydoc.flushSave = async () => assert.fail('External invalidation must not flush stale edits'); + globalThis.fetch = async () => new Response(null, { headers: { etag: '"new"' } }); + + await persistence.checkEtag(ydoc); + + assert(cancelled); + assert(conn.closed); + assert(secondConn.closed); + assert.equal(docs.get(ydoc.name), undefined); + }); + + it('re-entrant awareness cleanup cannot flush edits during external invalidation', async () => { + const { ydoc, conn, docs } = makeDoc('reentrant-invalidation'); + const secondConn = { + readyState: 1, + closed: false, + send() {}, + close() { this.closed = true; }, + }; + ydoc.conns.set(secondConn, new Set()); + ydoc.awareness.setLocalState({ user: 'test' }); + ydoc.conns.get(conn).add(ydoc.clientID); + secondConn.readyState = 3; + let flushed = false; + ydoc.flushSave = async () => { + flushed = true; + }; + globalThis.fetch = async () => new Response(null, { headers: { etag: '"new"' } }); + + await persistence.checkEtag(ydoc); + + assert.equal(flushed, false); + assert(conn.closed); + assert(secondConn.closed); + assert.equal(docs.get(ydoc.name), undefined); + }); + + for (const [status, sameContent] of [[200, false], [200, true], [404, false], [410, false]]) { + it(`HEAD ${status} external invalidation${sameContent ? ' with identical content' : ''} cannot POST pending client edits`, async () => { + const { ydoc, conn, docs } = makeDoc(`pending-${status}`); + const snapshot = new Y.Doc(); + aem2doc('

Original document before client edits

', snapshot); + const original = doc2aem(snapshot); + const storage = makeStorage({ + doc: ydoc.name, + docstore: Y.encodeStateAsUpdate(snapshot), + lastsync: original, + }); + snapshot.destroy(); + let externalContent = null; + if (status === 200) { + externalContent = sameContent + ? original + : '

Externally replaced document

'; + } + let backendContent = original; + let backendEtag = '"old"'; + const writes = []; + globalThis.fetch = async (_url, opts) => { + if (opts.method === 'HEAD') { + return new Response(null, { + status, + headers: status === 200 ? { etag: '"external"' } : {}, + }); + } + if (opts.method === 'POST') { + writes.push(opts.body); + backendContent = opts.body; + return new Response(null, { headers: { etag: '"saved"' } }); + } + return new Response(backendContent, { headers: { etag: backendEtag } }); + }; + + await persistence.bindState(ydoc.name, ydoc, conn, storage); + ydoc.hasClientChanged = true; + const fragment = ydoc.getXmlFragment('prosemirror'); + fragment.delete(0, fragment.length); + aem2doc('

Unsaved client edit from the old document

', ydoc); + backendContent = externalContent; + backendEtag = '"external"'; + + await persistence.checkEtag(ydoc); + await ydoc.flushSave(); + + assert.deepStrictEqual(writes, []); + assert.strictEqual(backendContent, externalContent); + assert(conn.closed); + assert.equal(docs.get(ydoc.name), undefined); + assert.equal(await storage.get('lastsync'), undefined); + + if (sameContent) { + const { ydoc: reconnected, conn: newConn } = makeDoc(`pending-${status}`, backendEtag); + const restoring = []; + await persistence.bindState(reconnected.name, reconnected, newConn, storage, { + waitUntil: (promise) => restoring.push(promise), + }); + await Promise.all(restoring); + + assert.equal(doc2aem(reconnected), original, 'Stored pending edits must not survive invalidation'); + assert.equal(await storage.get('lastsync'), original); + assert.deepStrictEqual(writes, []); + } + }); + } + + for (const status of [404, 410]) { + it(`HEAD ${status} invalidates even without a stored ETag`, async () => { + const { ydoc, conn, docs } = makeDoc(`deleted-${status}`, null); + globalThis.fetch = async () => new Response(null, { status }); + + await persistence.checkEtag(ydoc); + + assert(conn.closed); + assert.equal(docs.get(ydoc.name), undefined); + }); + } + + for (const [status, level] of [[401, 'warn'], [403, 'log'], [429, 'error'], [500, 'error'], [503, 'error']]) { + it(`HEAD ${status} is logged without invalidating and permits the next poll`, async () => { + const { ydoc, conn, docs } = makeDoc(`failed-${status}`); + const logged = []; + const savedLog = console[level]; + console[level] = (...args) => logged.push(args); + let calls = 0; + globalThis.fetch = async () => { + calls += 1; + return calls === 1 + ? new Response(null, { status }) + : new Response(null, { headers: { etag: ydoc.etag } }); + }; + try { + await persistence.checkEtag(ydoc); + await persistence.checkEtag(ydoc); + + assert.equal(conn.closed, false); + assert.strictEqual(docs.get(ydoc.name), ydoc); + assert.equal(calls, 2); + assert.equal(logged.length, 1); + assert(logged[0][0].includes(String(status))); + } finally { + console[level] = savedLog; + } + }); + } + + it('a successful HEAD without an ETag warns without invalidating', async () => { + const { ydoc, conn, docs } = makeDoc('missing-etag'); + const logged = []; + const savedWarn = console.warn; + console.warn = (...args) => logged.push(args); + globalThis.fetch = async () => new Response(null); + try { + await persistence.checkEtag(ydoc); + + assert.equal(conn.closed, false); + assert.strictEqual(docs.get(ydoc.name), ydoc); + assert.equal(logged.length, 1); + assert(logged[0][0].includes('missing ETag')); + } finally { + console.warn = savedWarn; + } + }); + + it('a failed fetch is logged and does not block later polls', async () => { + const { ydoc, conn } = makeDoc('fetch-error'); + const logged = []; + const savedError = console.error; + console.error = (...args) => logged.push(args); + let calls = 0; + globalThis.fetch = async () => { + calls += 1; + if (calls === 1) { + throw new Error('Network failure'); + } + return new Response(null, { headers: { etag: ydoc.etag } }); + }; + try { + await persistence.checkEtag(ydoc); + await persistence.checkEtag(ydoc); + + assert.equal(conn.closed, false); + assert.equal(calls, 2); + assert.equal(logged.length, 1); + assert(logged[0][0].includes('Etag check failed')); + } finally { + console.error = savedError; + } + }); + + it('only one HEAD request can be in flight for a document', async () => { + const { ydoc, conn } = makeDoc('single-flight'); + const head = deferred(); + let calls = 0; + globalThis.fetch = async () => { + calls += 1; + return head.promise; + }; + + const checking = persistence.checkEtag(ydoc); + await persistence.checkEtag(ydoc); + assert.equal(calls, 1); + head.resolve(new Response(null, { headers: { etag: ydoc.etag } })); + await checking; + await persistence.checkEtag(ydoc); + + assert.equal(calls, 2); + assert.equal(conn.closed, false); + }); + + it('does not poll a closed session', async () => { + const { ydoc, conn } = makeDoc('closed'); + globalThis.fetch = async () => assert.fail('A closed session must not be polled'); + + await closeConn(ydoc, conn); + await persistence.checkEtag(ydoc); + }); + + it('does not poll a destroyed document still present in the cache', async () => { + const { ydoc } = makeDoc('destroyed'); + globalThis.fetch = async () => assert.fail('A destroyed document must not be polled'); + + ydoc.destroy(); + await persistence.checkEtag(ydoc); + }); + + for (const action of ['close', 'destroy']) { + it(`${action} aborts a pending HEAD without logging an expected cancellation`, async () => { + const { ydoc, conn } = makeDoc(`abort-${action}`); + const head = deferred(); + const logged = []; + const savedError = console.error; + console.error = (...args) => logged.push(args); + let signal; + globalThis.fetch = async (_url, opts) => { + signal = opts.signal; + signal?.addEventListener('abort', () => head.reject(signal.reason), { once: true }); + return head.promise; + }; + const checking = persistence.checkEtag(ydoc); + try { + assert(signal instanceof AbortSignal, 'The backend HEAD must receive an AbortSignal'); + if (action === 'close') { + await closeConn(ydoc, conn); + } else { + ydoc.destroy(); + } + assert(signal.aborted, 'Destruction must cancel the outstanding backend request'); + await checking; + + assert.equal(ydoc.etagCheckAbortController, undefined); + assert.deepStrictEqual(logged, []); + } finally { + head.resolve(new Response(null, { headers: { etag: ydoc.etag } })); + await checking; + console.error = savedError; + } + }); + } + + it('keeps a pending HEAD until the last collaborator disconnects', async () => { + const { ydoc, conn } = makeDoc('abort-last-collaborator'); + const secondConn = { + readyState: 1, + send() {}, + close() {}, + }; + ydoc.conns.set(secondConn, new Set()); + const head = deferred(); + let signal; + globalThis.fetch = async (_url, opts) => { + signal = opts.signal; + signal?.addEventListener('abort', () => head.reject(signal.reason), { once: true }); + return head.promise; + }; + const checking = persistence.checkEtag(ydoc); + try { + assert(signal instanceof AbortSignal); + await closeConn(ydoc, conn); + assert.equal(signal.aborted, false); + assert.strictEqual(ydoc.etagCheckAbortController.signal, signal); + + await closeConn(ydoc, secondConn); + assert(signal.aborted); + await checking; + assert.equal(ydoc.etagCheckAbortController, undefined); + } finally { + head.resolve(new Response(null, { headers: { etag: ydoc.etag } })); + await checking; + } + }); + + it('releases a completed HEAD controller without aborting it on later destruction', async () => { + const { ydoc, conn } = makeDoc('completed-head'); + let signal; + globalThis.fetch = async (_url, opts) => { + signal = opts.signal; + return new Response(null, { headers: { etag: ydoc.etag } }); + }; + + await persistence.checkEtag(ydoc); + assert(signal instanceof AbortSignal); + assert.equal(ydoc.etagCheckAbortController, undefined); + await closeConn(ydoc, conn); + assert.equal(signal.aborted, false); + }); + + it('still logs an unrelated HEAD error after the session closes', async () => { + const { ydoc, conn } = makeDoc('error-after-close'); + const head = deferred(); + const logged = []; + const savedError = console.error; + console.error = (...args) => logged.push(args); + const error = new Error('Transport failure'); + globalThis.fetch = async () => head.promise; + const checking = persistence.checkEtag(ydoc); + try { + await closeConn(ydoc, conn); + head.reject(error); + await checking; + + assert.equal(logged.length, 1); + assert(logged[0].includes(error)); + assert.equal(ydoc.etagCheckAbortController, undefined); + } finally { + head.resolve(new Response(null, { headers: { etag: ydoc.etag } })); + await checking; + console.error = savedError; + } + }); + + describe('initialization lifecycle', () => { + let savedSetInterval; + let savedClearInterval; + const scheduled = []; + const cleared = []; + + beforeEach(() => { + scheduled.length = 0; + cleared.length = 0; + savedSetInterval = globalThis.setInterval; + savedClearInterval = globalThis.clearInterval; + globalThis.setInterval = (callback, delay, ...args) => { + const handle = savedSetInterval(callback, delay, ...args); + scheduled.push({ handle, delay }); + return handle; + }; + globalThis.clearInterval = (handle) => { + cleared.push(handle); + savedClearInterval(handle); + }; + }); + + afterEach(() => { + scheduled.forEach(({ handle }) => savedClearInterval(handle)); + globalThis.setInterval = savedSetInterval; + globalThis.clearInterval = savedClearInterval; + }); + + for (const transition of ['disconnected', 'replaced', 'no-connections', 'destroyed']) { + it(`does not start polling when initialization finishes for a ${transition} session`, async () => { + const { ydoc, conn, docs } = makeDoc(`initializing-${transition}`); + const get = deferred(); + globalThis.fetch = async () => get.promise; + const loading = getYDoc(ydoc.name, conn, {}, makeStorage()); + try { + if (transition === 'disconnected') { + await closeConn(ydoc, conn); + } else if (transition === 'replaced') { + makeDoc(`initializing-${transition}`); + } else if (transition === 'no-connections') { + ydoc.conns.clear(); + } else { + ydoc.destroy(); + } + get.resolve(new Response('

Source content

', { + headers: { etag: '"source"' }, + })); + await loading; + + assert.equal(scheduled.filter(({ delay }) => delay === 5000).length, 0); + assert.equal(ydoc.etagCheckInterval, undefined); + if (transition === 'disconnected') { + assert.equal(docs.get(ydoc.name), undefined); + assert(scheduled.every(({ handle }) => cleared.includes(handle))); + } + if (transition === 'replaced') { + assert.notStrictEqual(docs.get(ydoc.name), ydoc); + } + } finally { + get.resolve(new Response(null)); + await loading; + } + }); + } + + it('does not start polling after disconnecting during the initial storage read', async () => { + const { ydoc, conn } = makeDoc('initializing-storage'); + const reading = deferred(); + const stored = deferred(); + const storage = makeStorage(); + const get = storage.get.bind(storage); + storage.get = async (key) => { + if (key === 'doc') { + reading.resolve(); + await stored.promise; + } + return get(key); + }; + globalThis.fetch = async () => new Response('

Source content

', { + headers: { etag: '"source"' }, + }); + const loading = getYDoc(ydoc.name, conn, {}, storage); + try { + await reading.promise; + await closeConn(ydoc, conn); + stored.resolve(); + await loading; + + assert.equal(scheduled.filter(({ delay }) => delay === 5000).length, 0); + assert.equal(ydoc.etagCheckInterval, undefined); + } finally { + stored.resolve(); + await loading; + } + }); + + it('starts a five-second poll for a live session and clears it on last close', async () => { + const { ydoc, conn } = makeDoc('live-polling'); + globalThis.fetch = async () => new Response('

Source content

', { + headers: { etag: ydoc.etag }, + }); + + await getYDoc(ydoc.name, conn, {}, makeStorage()); + const polling = scheduled.filter(({ delay }) => delay === 5000); + assert.equal(polling.length, 1); + assert.strictEqual(ydoc.etagCheckInterval, polling[0].handle); + + await closeConn(ydoc, conn); + assert(cleared.includes(polling[0].handle)); + }); + }); + + it('ignores an old HEAD after the session was replaced', async () => { + const { ydoc: oldDoc, conn: oldConn } = makeDoc('replacement'); + const head = deferred(); + globalThis.fetch = async () => head.promise; + + const checking = persistence.checkEtag(oldDoc); + await closeConn(oldDoc, oldConn); + const { ydoc: newDoc, conn: newConn, docs } = makeDoc('replacement', '"replacement"'); + head.resolve(new Response(null, { headers: { etag: '"new"' } })); + await checking; + + assert.equal(newConn.closed, false); + assert.strictEqual(docs.get(newDoc.name), newDoc); + }); + + it('ignores a HEAD whose ETag baseline changed while it was in flight', async () => { + const { ydoc, conn } = makeDoc('changed-baseline'); + const head = deferred(); + globalThis.fetch = async () => head.promise; + + const checking = persistence.checkEtag(ydoc); + ydoc.etag = '"new-baseline"'; + head.resolve(new Response(null, { headers: { etag: '"old"' } })); + await checking; + + assert.equal(conn.closed, false); + }); + + for (const [label, responseEtag] of [['changed', '"saved"'], ['same', '"old"'], ['missing', null]]) { + it(`ignores a HEAD overlapping a local save with response ETag ${responseEtag}`, async () => { + const { ydoc, conn } = makeDoc(`own-save-${label}`); + const head = deferred(); + globalThis.fetch = async (_url, opts) => { + if (opts.method === 'HEAD') { + return head.promise; + } + return new Response(null, { + headers: responseEtag ? { etag: responseEtag } : {}, + }); + }; + + const checking = persistence.checkEtag(ydoc); + await persistence.put(ydoc, '

New local content with enough padding to avoid empty warnings

'); + head.resolve(new Response(null, { headers: { etag: '"different"' } })); + await checking; + + assert.equal(conn.closed, false); + assert.equal(ydoc.pendingPuts, 0); + }); + } + + it('does not poll during an in-flight save and ignores a HEAD preceding that save', async () => { + const { ydoc, conn } = makeDoc('in-flight-save'); + const head = deferred(); + const post = deferred(); + let headCalls = 0; + globalThis.fetch = async (_url, opts) => { + if (opts.method === 'HEAD') { + headCalls += 1; + return head.promise; + } + return post.promise; + }; + + const checking = persistence.checkEtag(ydoc); + const saving = persistence.put(ydoc, '

New local content with enough padding to avoid empty warnings

'); + await persistence.checkEtag(ydoc); + assert.equal(headCalls, 1); + head.resolve(new Response(null, { headers: { etag: '"different"' } })); + await checking; + assert.equal(conn.closed, false); + post.resolve(new Response(null, { headers: { etag: '"saved"' } })); + await saving; + await persistence.checkEtag(ydoc); + + assert.equal(headCalls, 2); + assert.equal(ydoc.pendingPuts, 0); + assert(conn.closed, 'The next poll must still detect an external change'); + }); + + it('a failed local save releases the polling guard', async () => { + const { ydoc, conn } = makeDoc('failed-save'); + let headCalls = 0; + globalThis.fetch = async (_url, opts) => { + if (opts.method === 'HEAD') { + headCalls += 1; + return new Response(null, { headers: { etag: ydoc.etag } }); + } + throw new Error('Write failed'); + }; + + await assert.rejects( + persistence.put(ydoc, '

New local content with enough padding to avoid empty warnings

'), + /Write failed/, + ); + await persistence.checkEtag(ydoc); + + assert.equal(ydoc.pendingPuts, 0); + assert.equal(headCalls, 1); + assert.equal(conn.closed, false); + }); + }); + it('persistence.bindState reads a Helix doc through the global fetch', async () => { const savedFetch = globalThis.fetch; const savedUpdate = persistence.update; + const docName = 'https://api.aem.live/o/r/bindstate.html'; + const ydoc = new WSSharedDoc(docName); + const docs = setYDoc(docName, ydoc); const calls = []; globalThis.fetch = async (url, opts) => { calls.push({ url, opts }); return { - ok: true, text: async () => 'helix content', status: 200, statusText: 'OK', + ok: true, text: async () => 'helix content', status: 200, statusText: 'OK', headers: new Headers({ etag: '"initial"' }), }; }; persistence.update = async () => {}; try { - const docName = 'https://api.aem.live/o/r/bindstate.html'; - const ydoc = new Y.Doc(); ydoc.daadmin = { fetch: async () => { assert.fail('daadmin.fetch must not be called for Helix docs'); }, }; const mockConn = { auth: 'Bearer x' }; - setYDoc(docName, ydoc); - const storage = { list: async () => new Map() }; + ydoc.conns.set(mockConn, new Set()); + const storage = makeStorage(); await persistence.bindState(docName, ydoc, mockConn, storage); @@ -3300,6 +4445,9 @@ describe('Collab Test Suite', () => { } finally { globalThis.fetch = savedFetch; persistence.update = savedUpdate; + ydoc.conns.clear(); + ydoc.destroy(); + docs.delete(docName); } }); });