diff --git a/bun.lock b/bun.lock index 606005a..01bc9db 100644 --- a/bun.lock +++ b/bun.lock @@ -27,6 +27,7 @@ "elysia": "^1.4.28", "fflate": "^0.8.3", "katex": "^0.16.47", + "lib0": "^0.2.117", "markdown-it": "^14.1.1", "markdown-it-footnote": "^4.0.0", "markdown-it-mark": "^4.0.0", @@ -34,6 +35,9 @@ "sanitize-html": "^2.17.4", "sharp": "^0.34.5", "shiki": "^4.1.0", + "y-codemirror.next": "^0.3.5", + "y-protocols": "^1.0.7", + "yjs": "^13.6.30", }, "devDependencies": { "@biomejs/biome": "^2.4.15", @@ -594,6 +598,8 @@ "is-plain-object": ["is-plain-object@5.0.0", "", {}, "sha512-VRSzKkbMm5jMDoKLbltAkFQ5Qr7VDiTFGXxYFXXowVj387GeGNOCsOH6Msy00SGZ3Fp84b1Naa1psqgcCIEP5Q=="], + "isomorphic.js": ["isomorphic.js@0.2.5", "", {}, "sha512-PIeMbHqMt4DnUP3MA/Flc0HElYjMXArsw1qwJZcm9sqR8mq3l8NYizFMty0pWwE/tzIGH3EKK5+jes5mAr85yw=="], + "jiti": ["jiti@2.7.0", "", { "bin": { "jiti": "lib/jiti-cli.mjs" } }, "sha512-AC/7JofJvZGrrneWNaEnJeOLUx+JlGt7tNa0wZiRPT4MY1wmfKjt2+6O2p2uz2+skll8OZZmJMNqeke7kKbNgQ=="], "katex": ["katex@0.16.47", "", { "dependencies": { "commander": "^8.3.0" }, "bin": { "katex": "cli.js" } }, "sha512-Eeo8Ys1doU1z+x8AZsPpQu+p/QcZBI5PeOo7QGQdy2x2m0MU/hYagBbGOmXwr5KVbEfVuWv9LpnQWeehogurjg=="], @@ -602,6 +608,8 @@ "launder": ["launder@1.7.1", "", { "dependencies": { "dayjs": "^1.11.7" } }, "sha512-mU6WRz5EusL9ZZuiZ5SO4Y6C0P9PAUR9iwdb6bzj4KDihm28DiHFw+/yk9DBH4f+Pv1wuzQ4e2jV3oQ7mkIqvw=="], + "lib0": ["lib0@0.2.117", "", { "dependencies": { "isomorphic.js": "^0.2.4" }, "bin": { "0serve": "bin/0serve.js", "0gentesthtml": "bin/gentesthtml.js", "0ecdsa-generate-keypair": "bin/0ecdsa-generate-keypair.js" } }, "sha512-DeXj9X5xDCjgKLU/7RR+/HQEVzuuEUiwldwOGsHK/sfAfELGWEyTcf0x+uOvCvK3O2zPmZePXWL85vtia6GyZw=="], + "lightningcss": ["lightningcss@1.32.0", "", { "dependencies": { "detect-libc": "^2.0.3" }, "optionalDependencies": { "lightningcss-android-arm64": "1.32.0", "lightningcss-darwin-arm64": "1.32.0", "lightningcss-darwin-x64": "1.32.0", "lightningcss-freebsd-x64": "1.32.0", "lightningcss-linux-arm-gnueabihf": "1.32.0", "lightningcss-linux-arm64-gnu": "1.32.0", "lightningcss-linux-arm64-musl": "1.32.0", "lightningcss-linux-x64-gnu": "1.32.0", "lightningcss-linux-x64-musl": "1.32.0", "lightningcss-win32-arm64-msvc": "1.32.0", "lightningcss-win32-x64-msvc": "1.32.0" } }, "sha512-NXYBzinNrblfraPGyrbPoD19C1h9lfI/1mzgWYvXUTe414Gz/X1FD2XBZSZM7rRTrMA8JL3OtAaGifrIKhQ5yQ=="], "lightningcss-android-arm64": ["lightningcss-android-arm64@1.32.0", "", { "os": "android", "cpu": "arm64" }, "sha512-YK7/ClTt4kAK0vo6w3X+Pnm0D2cf2vPHbhOXdoNti1Ga0al1P4TBZhwjATvjNwLEBCnKvjJc2jQgHXH0NEwlAg=="], @@ -764,8 +772,14 @@ "walk-up-path": ["walk-up-path@4.0.0", "", {}, "sha512-3hu+tD8YzSLGuFYtPRb48vdhKMi0KQV5sn+uWr8+7dMEq/2G/dtLrdDinkLjqq5TIbIBjYJ4Ax/n3YiaW7QM8A=="], + "y-codemirror.next": ["y-codemirror.next@0.3.5", "", { "dependencies": { "lib0": "^0.2.42" }, "peerDependencies": { "@codemirror/state": "^6.0.0", "@codemirror/view": "^6.0.0", "yjs": "^13.5.6" } }, "sha512-VluNu3e5HfEXybnypnsGwKAj+fKLd4iAnR7JuX1Sfyydmn1jCBS5wwEL/uS04Ch2ib0DnMAOF6ZRR/8kK3wyGw=="], + + "y-protocols": ["y-protocols@1.0.7", "", { "dependencies": { "lib0": "^0.2.85" }, "peerDependencies": { "yjs": "^13.0.0" } }, "sha512-YSVsLoXxO67J6eE/nV4AtFtT3QEotZf5sK5BHxFBXso7VDUT3Tx07IfA6hsu5Q5OmBdMkQVmFZ9QOA7fikWvnw=="], + "yaml": ["yaml@2.9.0", "", { "bin": { "yaml": "bin.mjs" } }, "sha512-2AvhNX3mb8zd6Zy7INTtSpl1F15HW6Wnqj0srWlkKLcpYl/gMIMJiyuGq2KeI2YFxUPjdlB+3Lc10seMLtL4cA=="], + "yjs": ["yjs@13.6.30", "", { "dependencies": { "lib0": "^0.2.99" } }, "sha512-vv/9h42eCMC81ZHDFswuu/MKzkl/vyq1BhaNGfHyOonwlG4CJbQF4oiBBJPvfdeCt/PlVDWh7Nov9D34YY09uQ=="], + "yocto-queue": ["yocto-queue@1.2.2", "", {}, "sha512-4LCcse/U2MHZ63HAJVE+v71o7yOdIe4cZ70Wpf8D/IyjDKYQLV5GD46B+hSTjJsvV5PztjvHoU580EftxjDZFQ=="], "zod": ["zod@4.4.3", "", {}, "sha512-ytENFjIJFl2UwYglde2jchW2Hwm4GJFLDiSXWdTrJQBIN9Fcyp7n4DhxJEiWNAJMV1/BqWfW/kkg71UDcHJyTQ=="], diff --git a/package.json b/package.json index 3678d5f..e9375c5 100644 --- a/package.json +++ b/package.json @@ -55,13 +55,17 @@ "elysia": "^1.4.28", "fflate": "^0.8.3", "katex": "^0.16.47", + "lib0": "^0.2.117", "markdown-it": "^14.1.1", "markdown-it-footnote": "^4.0.0", "markdown-it-mark": "^4.0.0", "markdown-it-task-lists": "^2.1.1", "sanitize-html": "^2.17.4", "sharp": "^0.34.5", - "shiki": "^4.1.0" + "shiki": "^4.1.0", + "y-codemirror.next": "^0.3.5", + "y-protocols": "^1.0.7", + "yjs": "^13.6.30" }, "overrides": { "@codemirror/state": "6.6.0", diff --git a/public/editor/collab.ts b/public/editor/collab.ts new file mode 100644 index 0000000..be2cd0c --- /dev/null +++ b/public/editor/collab.ts @@ -0,0 +1,182 @@ +import type { Extension } from "@codemirror/state"; +import { yCollab } from "y-codemirror.next"; +import { Awareness, applyAwarenessUpdate } from "y-protocols/awareness"; +import * as Y from "yjs"; +import { + encodeAwareness, + encodeControl, + encodeSyncStep1, + encodeUpdate, + MSG_AWARENESS, + MSG_CONTROL, + MSG_SYNC, + readAwarenessFrame, + readControlFrame, + readFrameType, + readSyncFrame, +} from "../../src/shared/collab-protocol.ts"; + +// Distinct, reasonably high-contrast cursor colors. Each DID maps deterministically +// to one so the same person keeps the same color across sessions. +const PALETTE = [ + "#e0567f", + "#2d8c6f", + "#d9822b", + "#5566d6", + "#0a9396", + "#9b5de5", + "#c1121f", + "#3a7d44", +]; + +function colorForDid(did: string): string { + let hash = 0; + for (let i = 0; i < did.length; i++) { + hash = (hash * 31 + did.charCodeAt(i)) >>> 0; + } + return PALETTE[hash % PALETTE.length] as string; +} + +interface CollabUser { + did: string; + handle: string; + displayName: string; +} + +export interface CollabHandle { + extension: Extension; + /** Current merged document text (source of truth for the save). */ + text(): string; + /** Handles of other connected editors, for the co-author message. */ + coAuthors(): string[]; + /** Tell the server the note was just saved — clears the draft, reloads peers. */ + markSaved(): void; + /** Discard the shared draft for everyone. */ + discard(): void; + destroy(): void; +} + +interface ConnectOptions { + wsPath: string; + user: CollabUser; + /** Called when the server signals the session ended (saved or discarded). */ + onReload: () => void; + onStatus?: (status: "connecting" | "connected" | "closed") => void; +} + +// Marks doc/awareness changes that originated from the server, so we don't echo +// them straight back. +const REMOTE_ORIGIN = "remote"; + +export function connectCollab(opts: ConnectOptions): CollabHandle { + const ydoc = new Y.Doc(); + const ytext = ydoc.getText("content"); + const awareness = new Awareness(ydoc); + const docOrigin = {}; // identity marker for updates applied from the server + + awareness.setLocalStateField("user", { + name: opts.user.displayName || opts.user.handle, + handle: opts.user.handle, + did: opts.user.did, + color: colorForDid(opts.user.did), + colorLight: `${colorForDid(opts.user.did)}33`, + }); + + let ws: WebSocket | null = null; + let closedByUs = false; + let reconnectTimer: ReturnType | undefined; + + const wsUrl = `${location.protocol === "https:" ? "wss" : "ws"}://${location.host}${opts.wsPath}`; + + function send(data: Uint8Array): void { + // lib0 encoders return a Uint8Array backed by a plain ArrayBuffer; the DOM + // lib's stricter Uint8Array typing can't see that, so assert it. + if (ws && ws.readyState === WebSocket.OPEN) { + ws.send(data as unknown as Uint8Array); + } + } + + const onDocUpdate = (update: Uint8Array, origin: unknown): void => { + if (origin === docOrigin) return; // came from the server; don't loop it back + send(encodeUpdate(update)); + }; + ydoc.on("update", onDocUpdate); + + const onAwarenessUpdate = ( + changes: { added: number[]; updated: number[]; removed: number[] }, + origin: unknown, + ): void => { + if (origin === REMOTE_ORIGIN) return; + const changed = changes.added.concat(changes.updated, changes.removed); + send(encodeAwareness(awareness, changed)); + }; + awareness.on("update", onAwarenessUpdate); + + function handleFrame(data: Uint8Array): void { + const { type, decoder } = readFrameType(data); + if (type === MSG_SYNC) { + const reply = readSyncFrame(decoder, ydoc, docOrigin); + if (reply) send(reply); + } else if (type === MSG_AWARENESS) { + applyAwarenessUpdate( + awareness, + readAwarenessFrame(decoder), + REMOTE_ORIGIN, + ); + } else if (type === MSG_CONTROL) { + if (readControlFrame(decoder).t === "reload") opts.onReload(); + } + } + + function connect(): void { + opts.onStatus?.("connecting"); + ws = new WebSocket(wsUrl); + ws.binaryType = "arraybuffer"; + ws.onopen = () => { + opts.onStatus?.("connected"); + // Both sides open with sync step 1; also announce our presence. + send(encodeSyncStep1(ydoc)); + send(encodeAwareness(awareness, [ydoc.clientID])); + }; + ws.onmessage = (ev) => { + if (ev.data instanceof ArrayBuffer) handleFrame(new Uint8Array(ev.data)); + }; + ws.onclose = () => { + opts.onStatus?.("closed"); + if (!closedByUs) reconnectTimer = setTimeout(connect, 1500); + }; + ws.onerror = () => { + try { + ws?.close(); + } catch {} + }; + } + connect(); + + return { + extension: yCollab(ytext, awareness), + text: () => ytext.toString(), + coAuthors: () => { + const handles = new Set(); + for (const [clientId, state] of awareness.getStates()) { + if (clientId === ydoc.clientID) continue; + const user = (state as { user?: { handle?: string } }).user; + if (user?.handle) handles.add(user.handle); + } + return [...handles]; + }, + markSaved: () => send(encodeControl({ t: "saved" })), + discard: () => send(encodeControl({ t: "discard" })), + destroy: () => { + closedByUs = true; + if (reconnectTimer) clearTimeout(reconnectTimer); + ydoc.off("update", onDocUpdate); + awareness.off("update", onAwarenessUpdate); + try { + ws?.close(); + } catch {} + awareness.destroy(); + ydoc.destroy(); + }, + }; +} diff --git a/public/editor/editor.ts b/public/editor/editor.ts index 336e047..3f0d4b7 100644 --- a/public/editor/editor.ts +++ b/public/editor/editor.ts @@ -15,12 +15,77 @@ import { lineNumbers, } from "@codemirror/view"; +import { yUndoManagerKeymap } from "y-codemirror.next"; +import { type CollabHandle, connectCollab } from "./collab.ts"; import { renderPreview } from "./preview.ts"; import { createToolbar } from "./toolbar.ts"; import { showToast, syncBlobMetadata, uploadImage } from "./upload.ts"; import { wikilinkCompletion } from "./wikilink-autocomplete.ts"; let activeView: EditorView | null = null; +let activeCollab: CollabHandle | null = null; + +/** + * Save flow for live-collaboration mode. Reuses the normal edit route — the + * shared Y.Doc text is posted through the same `POST .../edit` endpoint as the + * solo form (PDS first, then DB). On success the server clears the draft and + * broadcasts a reload, which navigates everyone to the saved note. + */ +async function collabSave( + form: HTMLFormElement, + collab: CollabHandle, + noteHref: string, +): Promise { + syncBlobMetadata(); + const titleInput = form.querySelector("#title"); + const messageInput = form.querySelector("#message"); + const blobInput = document.getElementById( + "blob_metadata", + ) as HTMLInputElement | null; + + let message = messageInput?.value ?? ""; + const peers = collab.coAuthors(); + if (peers.length > 0) { + const suffix = `Co-edited by ${peers.map((h) => `@${h}`).join(", ")}`; + message = message ? `${message} — ${suffix}` : suffix; + } + + const body = new FormData(); + body.set("title", titleInput?.value ?? ""); + body.set("content", collab.text()); + body.set("message", message); + body.set("blob_metadata", blobInput?.value ?? ""); + + const csrf = + document + .querySelector('meta[name="csrf-token"]') + ?.getAttribute("content") ?? ""; + + try { + const resp = await fetch(form.action, { + method: "POST", + body, + headers: csrf ? { "X-CSRF-Token": csrf } : {}, + redirect: "manual", + }); + // The edit route 302-redirects on success (opaqueredirect under manual mode). + if ( + resp.type === "opaqueredirect" || + (resp.status >= 200 && resp.status < 300) + ) { + collab.markSaved(); + // Server broadcasts "reload" → onReload navigates. Fallback if the WS dropped. + if (noteHref) setTimeout(() => location.assign(noteHref), 2500); + } else { + const text = await resp.text().catch(() => ""); + showToast(`Save failed: ${text || resp.status}`); + } + } catch (err) { + showToast( + `Save failed: ${err instanceof Error ? err.message : String(err)}`, + ); + } +} const themedEditor = EditorView.theme({ "&": { @@ -132,22 +197,49 @@ function initEditor(root: Document | Element = document): void { ] : []; + // Live-collaboration mode (opt-in via the Collaborate button → ?collab=1). + // The editor connects to the shared Y.Doc; content arrives over the wire, so + // CodeMirror starts empty and uses the Yjs undo stack instead of CM history. + const collabEnabled = textarea.dataset["collab"] === "1"; + const collabWs = textarea.dataset["collabWs"]; + const userDid = textarea.dataset["did"]; + const userHandle = textarea.dataset["handle"]; + const displayName = textarea.dataset["displayName"] ?? ""; + const noteHref = textarea.dataset["noteUrl"] ?? ""; + + let collabHandle: CollabHandle | null = null; + if (collabEnabled && collabWs && userDid && userHandle) { + collabHandle = connectCollab({ + wsPath: collabWs, + user: { did: userDid, handle: userHandle, displayName }, + onReload: () => { + if (noteHref) location.assign(noteHref); + }, + }); + activeCollab = collabHandle; + } + const view = new EditorView({ state: EditorState.create({ - doc: textarea.value, + doc: collabHandle ? "" : textarea.value, extensions: [ lineNumbers(), drawSelection(), highlightActiveLine(), - history(), + ...(collabHandle ? [] : [history()]), - keymap.of([indentWithTab, ...defaultKeymap, ...historyKeymap]), + keymap.of([ + indentWithTab, + ...defaultKeymap, + ...(collabHandle ? yUndoManagerKeymap : historyKeymap), + ]), markdown(), themedEditor, ...autocompleteExt, + ...(collabHandle ? [collabHandle.extension] : []), EditorView.lineWrapping, @@ -238,11 +330,30 @@ function initEditor(root: Document | Element = document): void { renderPreview(preview, textarea.value); - form.addEventListener("submit", () => { + form.addEventListener("submit", (event) => { + if (activeCollab) { + event.preventDefault(); + void collabSave(form, activeCollab, noteHref); + return; + } textarea.value = view.state.doc.toString(); syncBlobMetadata(); }); + + if (activeCollab) { + const discardBtn = root.querySelector("[data-collab-discard]"); + discardBtn?.addEventListener("click", () => { + const collab = activeCollab; + if (!collab) return; + const confirmText = + discardBtn.dataset["confirm"] ?? "Discard the shared draft?"; + if (confirm(confirmText)) { + collab.discard(); + if (noteHref) setTimeout(() => location.assign(noteHref), 2500); + } + }); + } } /** @@ -319,6 +430,8 @@ document.addEventListener("htmx:beforeSwap", (event: Event) => { if (target instanceof Element && target.contains(activeView.dom)) { activeView.destroy(); activeView = null; + activeCollab?.destroy(); + activeCollab = null; } }); diff --git a/src/atproto/session.ts b/src/atproto/session.ts index ffae4fa..b14d4b0 100644 --- a/src/atproto/session.ts +++ b/src/atproto/session.ts @@ -80,12 +80,19 @@ export async function getClient(): Promise { return oauthClient; } -export async function getSessionFromRequest( - request: Request, +export async function getSessionFromCookie( + cookieHeader: string | null | undefined, ): Promise { + const cookie = cookieHeader ?? undefined; const client = await getClient(); if (client) { - return getSession(client, request.headers.get("cookie") ?? undefined); + return getSession(client, cookie); } - return getDevSession(request.headers.get("cookie") ?? undefined); + return getDevSession(cookie); +} + +export async function getSessionFromRequest( + request: Request, +): Promise { + return getSessionFromCookie(request.headers.get("cookie")); } diff --git a/src/lib/i18n/en.ts b/src/lib/i18n/en.ts index dd2a0da..0360df1 100644 --- a/src/lib/i18n/en.ts +++ b/src/lib/i18n/en.ts @@ -94,6 +94,11 @@ export const en: Messages = { titlePlaceholder: "My Note Title", importFile: "Import from file", importFileBrowse: "Browse", + collaborate: "Collaborate", + collabActive: "Live collaboration", + discardDraft: "Discard draft", + confirmDiscardDraft: + "Discard the shared draft? Everyone's unsaved changes will be lost.", }, createWiki: { heading: "Create a Wiki", diff --git a/src/lib/i18n/fr.ts b/src/lib/i18n/fr.ts index 631766a..dfe8655 100644 --- a/src/lib/i18n/fr.ts +++ b/src/lib/i18n/fr.ts @@ -96,6 +96,11 @@ export const fr: PartialMessages = { titlePlaceholder: "Titre de ma note", importFile: "Importer depuis un fichier", importFileBrowse: "Parcourir", + collaborate: "Collaborer", + collabActive: "Collaboration en direct", + discardDraft: "Abandonner le brouillon", + confirmDiscardDraft: + "Abandonner le brouillon partagé ? Les modifications non enregistrées de chacun seront perdues.", }, createWiki: { heading: "Créer un wiki", diff --git a/src/lib/i18n/index.ts b/src/lib/i18n/index.ts index 08d2402..8a36068 100644 --- a/src/lib/i18n/index.ts +++ b/src/lib/i18n/index.ts @@ -94,6 +94,10 @@ export interface Messages { titlePlaceholder: string; importFile: string; importFileBrowse: string; + collaborate: string; + collabActive: string; + discardDraft: string; + confirmDiscardDraft: string; }; createWiki: { heading: string; diff --git a/src/lib/urls.ts b/src/lib/urls.ts index 362d629..7116f56 100644 --- a/src/lib/urls.ts +++ b/src/lib/urls.ts @@ -64,6 +64,27 @@ export function historyNoteUrl( return `/@${handle}/${wikiSlug}/${noteSlug}/-/history`; } +/** Opens the editor in live-collaboration mode. */ +export function editCollabUrl( + handle: string, + wikiSlug: string, + noteSlug: string, +): string { + return `${editNoteUrl(handle, wikiSlug, noteSlug)}?collab=1`; +} + +/** + * Path for the collab WebSocket. The client prepends ws(s):// + host. Note: no + * `@` before the handle here — it matches the `/collab/:handle/...` route param. + */ +export function collabWsPath( + handle: string, + wikiSlug: string, + noteSlug: string, +): string { + return `/collab/${handle}/${wikiSlug}/${noteSlug}`; +} + export function redirect(url: string): Response { return new Response(null, { status: 302, diff --git a/src/server/app.ts b/src/server/app.ts index 937bc66..4a95824 100644 --- a/src/server/app.ts +++ b/src/server/app.ts @@ -9,6 +9,7 @@ import { canonicalHandlePlugin } from "./canonical-handle-plugin.ts"; import { getDb } from "./db/index.ts"; import { blobRoutes } from "./routes/blob.ts"; import { bookmarkRoutes } from "./routes/bookmark.ts"; +import { collabRoutes } from "./routes/collab.ts"; import { exploreRoutes } from "./routes/explore.ts"; import { homeRoute } from "./routes/home.ts"; import { localeRoutes } from "./routes/locale.ts"; @@ -72,6 +73,7 @@ export function buildApp() { .use(staticPlugin({ prefix: "/public", assets: "public" })) .use(canonicalHandlePlugin) .use(atprotoRoutes()) + .use(collabRoutes) .use(blobRoutes) .use(bookmarkRoutes) .use(localeRoutes) diff --git a/src/server/collab/manager.ts b/src/server/collab/manager.ts new file mode 100644 index 0000000..576b52f --- /dev/null +++ b/src/server/collab/manager.ts @@ -0,0 +1,272 @@ +import { + Awareness, + applyAwarenessUpdate, + removeAwarenessStates, +} from "y-protocols/awareness"; +import * as Y from "yjs"; +import * as protocol from "../../shared/collab-protocol.ts"; +import { + deleteDraft, + getCurrentNoteContent, + loadDraft, + saveDraftState, +} from "../db/queries/index.ts"; + +/** + * In-memory registry of live collaborative documents, one Y.Doc per note that + * currently has at least one editor connected. The server is a relay hub: it + * applies every client's sync/awareness messages to its own Y.Doc and rebroadcasts + * to peers, and batches the Y.Doc state to SQLite (the `drafts` table) so an + * in-progress draft survives disconnects and appview restarts. + * + * This holds drafts only — never the canonical record. Saving goes through the + * normal edit route (PDS first), after which `clearAfterSave` drops the draft. + */ + +const PERSIST_DEBOUNCE_MS = 5000; +const WS_OPEN = 1; + +/** Minimal view of a Bun ServerWebSocket — identity-keyed, used to push frames. */ +interface CollabSocket { + send(data: Uint8Array): unknown; + readonly readyState: number; +} + +interface Conn { + did: string; + /** Awareness client ids announced over this socket, removed on disconnect. */ + controlledIds: Set; +} + +interface LiveDoc { + noteAtUri: string; + wikiAtUri: string; + ydoc: Y.Doc; + awareness: Awareness; + conns: Map; + persistTimer: ReturnType | null; + lastModifierDid: string | null; +} + +const docs = new Map(); + +// --- Seed / persist (exported so the load/seed/persist logic is testable +// without a socket) --- + +/** + * Populate a fresh Y.Doc for a note: load an existing draft if one is stored, + * otherwise seed the "content" text from the note's last saved revision. + * Both paths use a non-client origin so the registry's update handler does not + * treat seeding as an edit to persist. + */ +export function seedDoc(ydoc: Y.Doc, noteAtUri: string): void { + const draft = loadDraft(noteAtUri); + if (draft) { + Y.applyUpdate(ydoc, draft.yjs_state, "load"); + return; + } + const content = getCurrentNoteContent(noteAtUri); + if (content) { + ydoc.transact(() => { + ydoc.getText("content").insert(0, content); + }, "seed"); + } +} + +/** Encode and store the current Y.Doc state to the drafts table. */ +export function persistDoc( + ydoc: Y.Doc, + noteAtUri: string, + wikiAtUri: string, + modifierDid: string | null, +): void { + saveDraftState( + noteAtUri, + wikiAtUri, + Y.encodeStateAsUpdate(ydoc), + modifierDid, + ); +} + +// --- Registry lifecycle --- + +function getOrLoad(noteAtUri: string, wikiAtUri: string): LiveDoc { + const existing = docs.get(noteAtUri); + if (existing) return existing; + + const ydoc = new Y.Doc(); + const awareness = new Awareness(ydoc); + awareness.setLocalState(null); // the server is a relay, not a participant + + const live: LiveDoc = { + noteAtUri, + wikiAtUri, + ydoc, + awareness, + conns: new Map(), + persistTimer: null, + lastModifierDid: null, + }; + + // Seed before attaching handlers so the seed/load doesn't broadcast or persist. + seedDoc(ydoc, noteAtUri); + + ydoc.on("update", (update: Uint8Array, origin: unknown) => { + broadcast(live, protocol.encodeUpdate(update)); + if (origin && live.conns.has(origin as CollabSocket)) { + live.lastModifierDid = + live.conns.get(origin as CollabSocket)?.did ?? live.lastModifierDid; + } + schedulePersist(live); + }); + + awareness.on( + "update", + ( + changes: { added: number[]; updated: number[]; removed: number[] }, + origin: unknown, + ) => { + const conn = origin ? live.conns.get(origin as CollabSocket) : undefined; + if (conn) { + for (const id of changes.added) conn.controlledIds.add(id); + for (const id of changes.removed) conn.controlledIds.delete(id); + } + const changed = changes.added.concat(changes.updated, changes.removed); + broadcast(live, protocol.encodeAwareness(awareness, changed)); + }, + ); + + docs.set(noteAtUri, live); + return live; +} + +function broadcast(live: LiveDoc, frame: Uint8Array): void { + for (const sock of live.conns.keys()) { + if (sock.readyState !== WS_OPEN) continue; + try { + sock.send(frame); + } catch { + // A dead socket is reaped on its own close event; ignore send failures. + } + } +} + +function schedulePersist(live: LiveDoc): void { + if (live.persistTimer) return; // batch: first edit opens a window, flush at its end + live.persistTimer = setTimeout(() => { + live.persistTimer = null; + persistDoc(live.ydoc, live.noteAtUri, live.wikiAtUri, live.lastModifierDid); + }, PERSIST_DEBOUNCE_MS); +} + +function dispose(live: LiveDoc): void { + if (live.persistTimer) { + clearTimeout(live.persistTimer); + live.persistTimer = null; + } + live.awareness.destroy(); + live.ydoc.destroy(); + docs.delete(live.noteAtUri); +} + +// --- Connection handlers (called from the WebSocket route) --- + +/** Register a socket, then send it sync step 1 and the current awareness states. */ +export function openConn( + noteAtUri: string, + wikiAtUri: string, + sock: CollabSocket, + did: string, +): void { + const live = getOrLoad(noteAtUri, wikiAtUri); + live.conns.set(sock, { did, controlledIds: new Set() }); + + if (sock.readyState === WS_OPEN) { + sock.send(protocol.encodeSyncStep1(live.ydoc)); + const states = live.awareness.getStates(); + if (states.size > 0) { + sock.send(protocol.encodeAwareness(live.awareness, [...states.keys()])); + } + } +} + +/** Handle one inbound frame from a client socket. */ +export function recvMessage( + noteAtUri: string, + sock: CollabSocket, + data: Uint8Array, +): void { + const live = docs.get(noteAtUri); + if (!live) return; + + const { type, decoder } = protocol.readFrameType(data); + switch (type) { + case protocol.MSG_SYNC: { + const reply = protocol.readSyncFrame(decoder, live.ydoc, sock); + if (reply && sock.readyState === WS_OPEN) sock.send(reply); + break; + } + case protocol.MSG_AWARENESS: + applyAwarenessUpdate( + live.awareness, + protocol.readAwarenessFrame(decoder), + sock, + ); + break; + case protocol.MSG_CONTROL: + handleControl(live, protocol.readControlFrame(decoder)); + break; + } +} + +function handleControl(live: LiveDoc, msg: protocol.ControlMessage): void { + // A save or discard ends the shared draft. Clear it, then tell ALL editors + // (including the initiator) to reload to the saved note. Broadcasting back to + // the initiator lets it navigate only once the server has confirmed the draft + // was cleared — no navigate-before-flush race on the WS. + if (msg.t === "saved" || msg.t === "discard") { + const all = [...live.conns.keys()]; + clearAfterSave(live.noteAtUri); + const frame = protocol.encodeControl({ t: "reload" }); + for (const peer of all) { + if (peer.readyState !== WS_OPEN) continue; + try { + peer.send(frame); + } catch {} + } + } +} + +/** Remove a socket; flush + dispose the doc once the last editor leaves. */ +export function closeConn(noteAtUri: string, sock: CollabSocket): void { + const live = docs.get(noteAtUri); + if (!live) return; + + const conn = live.conns.get(sock); + live.conns.delete(sock); + if (conn && conn.controlledIds.size > 0) { + // Broadcasts the removal to remaining peers so their cursors clear. + removeAwarenessStates(live.awareness, [...conn.controlledIds], null); + } + + if (live.conns.size === 0) { + if (live.persistTimer) { + clearTimeout(live.persistTimer); + live.persistTimer = null; + persistDoc( + live.ydoc, + live.noteAtUri, + live.wikiAtUri, + live.lastModifierDid, + ); + } + dispose(live); + } +} + +/** Drop the draft (DB + memory) after a save or discard. */ +export function clearAfterSave(noteAtUri: string): void { + const live = docs.get(noteAtUri); + deleteDraft(noteAtUri); + if (live) dispose(live); +} diff --git a/src/server/db/queries/draft.ts b/src/server/db/queries/draft.ts new file mode 100644 index 0000000..82bbfb7 --- /dev/null +++ b/src/server/db/queries/draft.ts @@ -0,0 +1,42 @@ +import { getDb } from "../index.ts"; +import type { DraftRow } from "../types.ts"; + +/** + * Collaborative-draft persistence. A draft holds the serialized Y.Doc state for + * a note while editors have a live collab session open. It is appview-only, + * ephemeral scratch state — the canonical record is the saved revision on the + * PDS. Drafts are seeded from the note's current content on first connect, + * batched to SQLite while editing, and deleted on save/discard. + */ + +export function loadDraft(noteAtUri: string): DraftRow | null { + const db = getDb(); + return ( + (db + .query("SELECT * FROM drafts WHERE note_at_uri = ?") + .get(noteAtUri) as DraftRow) ?? null + ); +} + +export function saveDraftState( + noteAtUri: string, + wikiAtUri: string, + state: Uint8Array, + modifierDid: string | null, +): void { + const db = getDb(); + db.run( + `INSERT INTO drafts (note_at_uri, wiki_at_uri, yjs_state, last_modifier_did, last_modified) + VALUES (?, ?, ?, ?, datetime('now')) + ON CONFLICT(note_at_uri) DO UPDATE SET + yjs_state = excluded.yjs_state, + last_modifier_did = excluded.last_modifier_did, + last_modified = excluded.last_modified`, + [noteAtUri, wikiAtUri, state, modifierDid], + ); +} + +export function deleteDraft(noteAtUri: string): void { + const db = getDb(); + db.run("DELETE FROM drafts WHERE note_at_uri = ?", [noteAtUri]); +} diff --git a/src/server/db/queries/index.ts b/src/server/db/queries/index.ts index 4ad0f9d..cef2479 100644 --- a/src/server/db/queries/index.ts +++ b/src/server/db/queries/index.ts @@ -9,6 +9,7 @@ export { upsertBookmark, } from "./bookmark.ts"; export { getCursor, setCursor } from "./cursor.ts"; +export { deleteDraft, loadDraft, saveDraftState } from "./draft.ts"; export { deleteMembership, deleteMembershipByUri, @@ -29,6 +30,7 @@ export { createNote, deleteNoteByAtUri, getCurrentNote, + getCurrentNoteContent, getNoteByAtUri, getNoteBySlug, getNoteWithCurrent, diff --git a/src/server/db/queries/note.ts b/src/server/db/queries/note.ts index ef85a7d..c01ea06 100644 --- a/src/server/db/queries/note.ts +++ b/src/server/db/queries/note.ts @@ -94,6 +94,18 @@ export function searchNotes( }); } +/** + * Current saved content keyed directly by note AT-URI. Used to seed a fresh + * collaborative draft from the last saved revision. Returns "" if absent. + */ +export function getCurrentNoteContent(noteAtUri: string): string { + const db = getDb(); + const row = db + .query("SELECT content FROM current_note WHERE note_at_uri = ?") + .get(noteAtUri) as { content: string } | null; + return row?.content ?? ""; +} + export function getCurrentNote( wikiAtUri: string, noteSlug: string, diff --git a/src/server/db/schema.ts b/src/server/db/schema.ts index 56460aa..d458a7a 100644 --- a/src/server/db/schema.ts +++ b/src/server/db/schema.ts @@ -129,6 +129,20 @@ export function initSchema(db: Database): void { ) `); + // Live collaborative-editing drafts. Appview-only, ephemeral: holds the + // shared Y.Doc state for a note while one or more editors have a collab + // session open. Cleared on save/discard. The canonical record is always the + // saved revision on the PDS — this is never an archive, just scratch state. + db.run(` + CREATE TABLE IF NOT EXISTS drafts ( + note_at_uri TEXT PRIMARY KEY REFERENCES notes(at_uri) ON DELETE CASCADE, + wiki_at_uri TEXT NOT NULL, + yjs_state BLOB NOT NULL, + last_modifier_did TEXT, + last_modified TEXT NOT NULL DEFAULT (datetime('now')) + ) + `); + db.run(` CREATE TABLE IF NOT EXISTS firehose_cursor ( id INTEGER PRIMARY KEY CHECK (id = 1), diff --git a/src/server/db/types.ts b/src/server/db/types.ts index 329f89c..950b6ed 100644 --- a/src/server/db/types.ts +++ b/src/server/db/types.ts @@ -76,3 +76,12 @@ export interface RequestRow { at_uri: string; created_at: string; } + +export interface DraftRow { + note_at_uri: string; + wiki_at_uri: string; + // Serialized Y.Doc (Y.encodeStateAsUpdate) — bun:sqlite returns BLOBs as Uint8Array. + yjs_state: Uint8Array; + last_modifier_did: string | null; + last_modified: string; +} diff --git a/src/server/routes/collab.ts b/src/server/routes/collab.ts new file mode 100644 index 0000000..2eae164 --- /dev/null +++ b/src/server/routes/collab.ts @@ -0,0 +1,129 @@ +import { Elysia } from "elysia"; +import { getAtprotoEnv } from "../../atproto/env.ts"; +import { getSessionFromCookie } from "../../atproto/session.ts"; +import { canEdit, getAccessLevel } from "../../lib/access.ts"; +import { resolveHandleToDid } from "../../lib/profile.ts"; +import { closeConn, openConn, recvMessage } from "../collab/manager.ts"; +import { getMemberRole, getNoteBySlug, getWiki } from "../db/queries/index.ts"; + +/** + * WebSocket transport for live collaborative editing. One connection per editor; + * the topic is the note's AT-URI. Auth is checked on connect (origin + session + + * canEdit); editing/cursors are relayed by the collab manager. Saving is NOT done + * here — the client saves through the normal edit route, then sends a `saved` + * control frame so the manager can clear the draft and reload peers. + */ + +interface ConnCtx { + noteAtUri: string; + wikiAtUri: string; + did: string; +} + +// Resolved per-connection state, keyed by the raw socket (identity-stable across +// open/message/close), so message handling doesn't re-resolve access per frame. +const connCtx = new WeakMap(); + +/** Reject cross-origin upgrades (anti-CSWSH). Browsers always send Origin on WS. */ +export function isSameOrigin(request: Request): boolean { + const origin = request.headers.get("origin"); + if (!origin) return true; // non-browser client (no ambient cookies to abuse) + let originHost: string; + try { + originHost = new URL(origin).host; + } catch { + return false; + } + const env = getAtprotoEnv(); + if (env) { + // Production runs behind Caddy on the apex domain — compare to PUBLIC_URL. + try { + return originHost === new URL(env.publicUrl).host; + } catch { + return false; + } + } + // Dev (no proxy): the request URL host is the real host. + try { + return originHost === new URL(request.url).host; + } catch { + return false; + } +} + +async function resolveCollabAccess( + handle: string, + wikiSlug: string, + noteSlug: string, + did: string, +): Promise { + const ownerDid = await resolveHandleToDid(handle); + if (!ownerDid) return null; + const wiki = getWiki(ownerDid, wikiSlug); + if (!wiki) return null; + const note = getNoteBySlug(wiki.at_uri, noteSlug); + if (!note) return null; + const role = getMemberRole(wiki.at_uri, did); + if (!canEdit(getAccessLevel(wiki, did, role))) return null; + return { noteAtUri: note.at_uri, wikiAtUri: wiki.at_uri, did }; +} + +function toBytes(message: unknown): Uint8Array | null { + if (message instanceof Uint8Array) return message; // Buffer is a Uint8Array + if (message instanceof ArrayBuffer) return new Uint8Array(message); + if (ArrayBuffer.isView(message)) { + return new Uint8Array( + message.buffer, + message.byteOffset, + message.byteLength, + ); + } + return null; +} + +export const collabRoutes = new Elysia().ws( + "/collab/:handle/:wikiSlug/:noteSlug", + { + async open(ws) { + const { handle, wikiSlug, noteSlug } = ws.data.params; + const request = ws.data.request; + + if (!isSameOrigin(request)) { + ws.close(1008, "bad origin"); + return; + } + const session = await getSessionFromCookie(request.headers.get("cookie")); + if (!session) { + ws.close(1008, "auth required"); + return; + } + const ctx = await resolveCollabAccess( + handle, + wikiSlug, + noteSlug, + session.did, + ); + if (!ctx) { + ws.close(1008, "forbidden"); + return; + } + + connCtx.set(ws.raw, ctx); + openConn(ctx.noteAtUri, ctx.wikiAtUri, ws.raw, ctx.did); + }, + + message(ws, message) { + const ctx = connCtx.get(ws.raw); + if (!ctx) return; + const bytes = toBytes(message); + if (bytes) recvMessage(ctx.noteAtUri, ws.raw, bytes); + }, + + close(ws) { + const ctx = connCtx.get(ws.raw); + if (!ctx) return; + closeConn(ctx.noteAtUri, ws.raw); + connCtx.delete(ws.raw); + }, + }, +); diff --git a/src/server/routes/note.ts b/src/server/routes/note.ts index c2e0b37..e2ac322 100644 --- a/src/server/routes/note.ts +++ b/src/server/routes/note.ts @@ -12,7 +12,7 @@ import { parseNoteFormFields, } from "../../lib/orchestrators/note.ts"; import { htmlResponse } from "../../lib/response.ts"; -import { noteUrl, redirect, wikiUrl } from "../../lib/urls.ts"; +import { collabWsPath, noteUrl, redirect, wikiUrl } from "../../lib/urls.ts"; import { editNotePage } from "../../views/edit-note.ts"; import { HISTORY_PAGE_SIZE, @@ -75,24 +75,47 @@ export const noteRoutes = new Elysia() throw err; } }) - .get("/@:handle/:wikiSlug/:noteSlug/edit", async ({ params, request }) => { - const { handle: urlHandle, wikiSlug, noteSlug } = hwnp(params); - const ctx = await resolveWikiContext(request, urlHandle, wikiSlug, "edit"); + .get( + "/@:handle/:wikiSlug/:noteSlug/edit", + async ({ params, request, query }) => { + const { handle: urlHandle, wikiSlug, noteSlug } = hwnp(params); + const ctx = await resolveWikiContext( + request, + urlHandle, + wikiSlug, + "edit", + ); - const data = getNoteWithCurrent(ctx.wiki.at_uri, noteSlug); - if (!data) - throw new NotFoundError("Note not found", { i18nKey: "noteNotFound" }); + const data = getNoteWithCurrent(ctx.wiki.at_uri, noteSlug); + if (!data) + throw new NotFoundError("Note not found", { i18nKey: "noteNotFound" }); - return htmlResponse( - editNotePage( - wikiIdentity(ctx), - noteSlug, - data.note.title, - data.current.content, - { ...wikiLayoutOptions(ctx), ...EDITOR_LAYOUT_EXTRAS }, - ), - ); - }) + // Opt-in live collaboration (?collab=1). Requires a session; the WS is + // access-checked again on connect. + const collab = + query["collab"] === "1" && ctx.did + ? { + enabled: true, + wsPath: collabWsPath(ctx.ownerHandle, ctx.wiki.slug, noteSlug), + did: ctx.did, + handle: ctx.session?.handle ?? ctx.did, + displayName: "", + noteUrl: noteUrl(ctx.ownerHandle, ctx.wiki.slug, noteSlug), + } + : undefined; + + return htmlResponse( + editNotePage( + wikiIdentity(ctx), + noteSlug, + data.note.title, + data.current.content, + { ...wikiLayoutOptions(ctx), ...EDITOR_LAYOUT_EXTRAS }, + collab, + ), + ); + }, + ) .post("/@:handle/:wikiSlug/:noteSlug/edit", async ({ params, request }) => { const { handle: urlHandle, wikiSlug, noteSlug } = hwnp(params); const ctx = await resolveWikiContext(request, urlHandle, wikiSlug, "edit"); diff --git a/src/shared/collab-protocol.ts b/src/shared/collab-protocol.ts new file mode 100644 index 0000000..a5acef5 --- /dev/null +++ b/src/shared/collab-protocol.ts @@ -0,0 +1,100 @@ +import * as decoding from "lib0/decoding"; +import * as encoding from "lib0/encoding"; +import type { Awareness } from "y-protocols/awareness"; +import { encodeAwarenessUpdate } from "y-protocols/awareness"; +import { readSyncMessage, writeSyncStep1, writeUpdate } from "y-protocols/sync"; +import type * as Y from "yjs"; + +/** + * Wire format for the collab WebSocket, shared by the server relay + * (src/server/collab) and the browser provider (public/editor/collab.ts). + * Every frame starts with a varUint message type. SYNC and AWARENESS payloads + * use the standard y-protocols encoding (the same framing y-websocket uses, + * minus its room/auth layers); CONTROL carries a small JSON payload for + * save/discard/reload coordination. + */ + +export const MSG_SYNC = 0; +export const MSG_AWARENESS = 1; +export const MSG_CONTROL = 2; + +export type ControlMessage = + | { t: "saved" } + | { t: "discard" } + | { t: "reload" }; + +/** Sync step 1 (state-vector request) — sent by both sides on connect. */ +export function encodeSyncStep1(doc: Y.Doc): Uint8Array { + const encoder = encoding.createEncoder(); + encoding.writeVarUint(encoder, MSG_SYNC); + writeSyncStep1(encoder, doc); + return encoding.toUint8Array(encoder); +} + +/** A document update to relay to peers. */ +export function encodeUpdate(update: Uint8Array): Uint8Array { + const encoder = encoding.createEncoder(); + encoding.writeVarUint(encoder, MSG_SYNC); + writeUpdate(encoder, update); + return encoding.toUint8Array(encoder); +} + +/** Awareness (cursor/presence) update for the given client ids. */ +export function encodeAwareness( + awareness: Awareness, + clients: number[], +): Uint8Array { + const encoder = encoding.createEncoder(); + encoding.writeVarUint(encoder, MSG_AWARENESS); + encoding.writeVarUint8Array( + encoder, + encodeAwarenessUpdate(awareness, clients), + ); + return encoding.toUint8Array(encoder); +} + +/** A control frame (JSON payload) for save/discard/reload coordination. */ +export function encodeControl(msg: ControlMessage): Uint8Array { + const encoder = encoding.createEncoder(); + encoding.writeVarUint(encoder, MSG_CONTROL); + encoding.writeVarString(encoder, JSON.stringify(msg)); + return encoding.toUint8Array(encoder); +} + +export interface IncomingFrame { + type: number; + decoder: decoding.Decoder; +} + +/** Read the leading message type; the returned decoder is positioned after it. */ +export function readFrameType(data: Uint8Array): IncomingFrame { + const decoder = decoding.createDecoder(data); + const type = decoding.readVarUint(decoder); + return { type, decoder }; +} + +/** + * Apply an inbound SYNC frame to `doc` (decoder positioned after the type byte), + * returning any reply bytes to send back to the sender, or null if none. + */ +export function readSyncFrame( + decoder: decoding.Decoder, + doc: Y.Doc, + origin: unknown, +): Uint8Array | null { + const encoder = encoding.createEncoder(); + encoding.writeVarUint(encoder, MSG_SYNC); + readSyncMessage(decoder, encoder, doc, origin); + // length 1 means only the type byte was written — nothing to reply. + return encoding.length(encoder) > 1 ? encoding.toUint8Array(encoder) : null; +} + +/** Read an AWARENESS frame's update payload (decoder positioned after the type byte). */ +export function readAwarenessFrame(decoder: decoding.Decoder): Uint8Array { + return decoding.readVarUint8Array(decoder); +} + +/** Read a CONTROL frame's JSON payload (decoder positioned after the type byte). */ +export function readControlFrame(decoder: decoding.Decoder): ControlMessage { + return JSON.parse(decoding.readVarString(decoder)) as ControlMessage; +} diff --git a/src/views/edit-note.ts b/src/views/edit-note.ts index ed42a77..7efc49c 100644 --- a/src/views/edit-note.ts +++ b/src/views/edit-note.ts @@ -5,16 +5,55 @@ import type { WikiIdentity } from "../lib/wiki-identity.ts"; import { type LayoutOptions, layout } from "./layout.ts"; import { inputClass, primaryButtonClass, THEME } from "./theme/index.ts"; +export interface CollabEditOptions { + enabled: boolean; + /** WebSocket path for the collab session (client prepends ws(s)://host). */ + wsPath: string; + /** Current editor's identity, for presence/awareness. */ + did: string; + handle: string; + displayName: string; + /** Where to navigate after save/discard. */ + noteUrl: string; +} + export function editNotePage( wiki: WikiIdentity, noteSlug: string, noteTitle: string, currentContent: string, options?: LayoutOptions, + collab?: CollabEditOptions, ): string { const locale = options?.locale ?? "en"; const msg = t(locale); + const collabAttrs = collab?.enabled + ? ` + data-collab="1" + data-collab-ws="${escapeHtml(collab.wsPath)}" + data-did="${escapeHtml(collab.did)}" + data-handle="${escapeHtml(collab.handle)}" + data-display-name="${escapeHtml(collab.displayName)}" + data-note-url="${escapeHtml(collab.noteUrl)}"` + : ""; + + const collabBanner = collab?.enabled + ? `
+ + ${msg.editor.collabActive} +
` + : ""; + + const discardButton = collab?.enabled + ? `` + : ""; + // CSS fix: ensures the edit form has a fixed height on all viewport widths // so that the CodeMirror editor can constrain its scroller and show a scrollbar. const scrollFixStyle = ` @@ -44,6 +83,7 @@ export function editNotePage( const formHtml = `
+ ${collabBanner}
${escapeHtml(currentContent)}
@@ -83,6 +123,7 @@ export function editNotePage( class="${primaryButtonClass}" >${msg.editor.save} ${msg.editor.cancel} + ${discardButton}