File
Blob: src/client/lib/doc-sync-session.ts
| 1 | import { useCallback, useEffect, useEffectEvent, useRef, useState } from "react"; |
| 2 | import * as Y from "yjs"; |
| 3 | import { IndexeddbPersistence } from "y-indexeddb"; |
| 4 | import YProvider from "y-partyserver/provider"; |
| 5 | import { api, refreshSession } from "@/client/lib/api"; |
| 6 | import { docCache } from "@/client/lib/doc-cache-registry"; |
| 7 | import { reconcileDocSyncProvider } from "@/client/lib/doc-sync-provider"; |
| 8 | import { reportClientError } from "@/client/lib/report-client-error"; |
| 9 | import { useOnline } from "@/client/hooks/use-online"; |
| 10 | import { useAuthStore } from "@/client/stores/auth-store"; |
| 11 | import { YJS_PAGE_TITLE } from "@/shared/constants"; |
| 12 | |
| 13 | export type DocSyncBootstrapStatus = "pending" | "resolved" | "error"; |
| 14 | |
| 15 | export interface DocSyncPhaseInputs { |
| 16 | hasLocalBodyState: boolean; |
| 17 | wantsConnection: boolean; |
| 18 | workspaceId: string | undefined; |
| 19 | bootstrapStatus: DocSyncBootstrapStatus; |
| 20 | } |
| 21 | |
| 22 | export interface DocSyncPhaseSnapshot { |
| 23 | ready: boolean; |
| 24 | shouldConnect: boolean; |
| 25 | snapshotFetch: { workspaceId: string } | null; |
| 26 | error: boolean; |
| 27 | } |
| 28 | |
| 29 | // Cold-bootstrap rule: if the session mounts against an empty local Y.Doc, |
| 30 | // any mount-time local mutation can merge into the authoritative document |
| 31 | // before remote content arrives. So when wantsConnection && !hasLocalBodyState, |
| 32 | // fetch the persisted snapshot over HTTP before letting the provider connect. |
| 33 | // A missing snapshot means the server is empty, which is safe to mount against. |
| 34 | export function deriveDocSyncPhase(i: DocSyncPhaseInputs): DocSyncPhaseSnapshot { |
| 35 | if (!i.wantsConnection) { |
| 36 | return { ready: true, shouldConnect: false, snapshotFetch: null, error: false }; |
| 37 | } |
| 38 | if (i.hasLocalBodyState) { |
| 39 | return { ready: true, shouldConnect: true, snapshotFetch: null, error: false }; |
| 40 | } |
| 41 | if (!i.workspaceId) { |
| 42 | return { ready: false, shouldConnect: false, snapshotFetch: null, error: false }; |
| 43 | } |
| 44 | switch (i.bootstrapStatus) { |
| 45 | case "pending": |
| 46 | return { ready: false, shouldConnect: false, snapshotFetch: { workspaceId: i.workspaceId }, error: false }; |
| 47 | case "error": |
| 48 | return { ready: false, shouldConnect: false, snapshotFetch: null, error: true }; |
| 49 | case "resolved": |
| 50 | return { ready: true, shouldConnect: true, snapshotFetch: null, error: false }; |
| 51 | } |
| 52 | } |
| 53 | |
| 54 | export interface DocSyncRefreshDecisionInput { |
| 55 | isOnline: boolean; |
| 56 | isProviderActive: boolean; |
| 57 | hasShareToken: boolean; |
| 58 | currentAccessToken: string | null; |
| 59 | lastRefreshAttemptedFor: string | null; |
| 60 | } |
| 61 | |
| 62 | export type DocSyncRefreshDecision = |
| 63 | | { kind: "skip"; reason: "share" | "offline" | "inactive" | "no_token" | "already_attempted" } |
| 64 | | { kind: "refresh"; tokenAtAttempt: string }; |
| 65 | |
| 66 | // Reconnect-time refresh policy. The browser cannot read the HTTP 401 from a |
| 67 | // failed WS upgrade, so every authenticated `connection-close` is treated as |
| 68 | // potentially auth-related. The gate keys on the access-token value to avoid |
| 69 | // looping on the same stale token, while still allowing a future expiration |
| 70 | // of the rotated token to refresh again. |
| 71 | export function decideDocSyncRefresh(i: DocSyncRefreshDecisionInput): DocSyncRefreshDecision { |
| 72 | if (!i.isOnline) return { kind: "skip", reason: "offline" }; |
| 73 | if (!i.isProviderActive) return { kind: "skip", reason: "inactive" }; |
| 74 | if (i.hasShareToken) return { kind: "skip", reason: "share" }; |
| 75 | if (!i.currentAccessToken) return { kind: "skip", reason: "no_token" }; |
| 76 | if (i.lastRefreshAttemptedFor === i.currentAccessToken) { |
| 77 | return { kind: "skip", reason: "already_attempted" }; |
| 78 | } |
| 79 | return { kind: "refresh", tokenAtAttempt: i.currentAccessToken }; |
| 80 | } |
| 81 | |
| 82 | export interface DocSyncSessionOptions<TRuntime> { |
| 83 | pageId: string; |
| 84 | initialTitle: string; |
| 85 | onTitleChange?: (title: string) => void; |
| 86 | onProvider?: (provider: YProvider | null) => void; |
| 87 | shareToken?: string; |
| 88 | workspaceId?: string; |
| 89 | enabled?: boolean; |
| 90 | /** IDB cache key for this session; must be stable for a given pageId. */ |
| 91 | cacheKey: string; |
| 92 | /** Reporting label for snapshot fetch failures. */ |
| 93 | errorSource: string; |
| 94 | /** Projects the doc-specific Y roots from the Y.Doc. */ |
| 95 | roots: (ydoc: Y.Doc) => TRuntime; |
| 96 | /** Returns true when `runtime` has body content worth mounting against. */ |
| 97 | hasBody: (runtime: TRuntime) => boolean; |
| 98 | } |
| 99 | |
| 100 | interface DocSyncSessionBase { |
| 101 | title: string; |
| 102 | onTitleInput: (event: React.ChangeEvent<HTMLTextAreaElement>) => void; |
| 103 | } |
| 104 | |
| 105 | export interface DocSyncSessionLoadingState extends DocSyncSessionBase { |
| 106 | kind: "loading"; |
| 107 | } |
| 108 | |
| 109 | export type DocSyncSessionReadyState<TRuntime> = DocSyncSessionBase & { |
| 110 | kind: "ready"; |
| 111 | ydoc: Y.Doc; |
| 112 | provider: YProvider; |
| 113 | } & TRuntime; |
| 114 | |
| 115 | export interface DocSyncSessionErrorState extends DocSyncSessionBase { |
| 116 | kind: "error"; |
| 117 | onRetry: () => void; |
| 118 | } |
| 119 | |
| 120 | export type DocSyncSessionState<TRuntime> = |
| 121 | | DocSyncSessionLoadingState |
| 122 | | DocSyncSessionReadyState<TRuntime> |
| 123 | | DocSyncSessionErrorState; |
| 124 | |
| 125 | export function useDocSyncSession<TRuntime extends object>({ |
| 126 | pageId, |
| 127 | initialTitle, |
| 128 | onTitleChange, |
| 129 | onProvider, |
| 130 | shareToken, |
| 131 | workspaceId, |
| 132 | enabled = true, |
| 133 | cacheKey, |
| 134 | errorSource, |
| 135 | roots, |
| 136 | hasBody, |
| 137 | }: DocSyncSessionOptions<TRuntime>): DocSyncSessionState<TRuntime> { |
| 138 | const online = useOnline(); |
| 139 | const isAuthed = useAuthStore((s) => !!s.accessToken); |
| 140 | const wantsConnection = online && (!!shareToken || isAuthed); |
| 141 | const sessionKey = `${cacheKey}:${shareToken ?? ""}`; |
| 142 | |
| 143 | const [titleState, setTitleState] = useState(() => ({ key: sessionKey, title: initialTitle })); |
| 144 | const [runtimeState, setRuntimeState] = useState<{ |
| 145 | key: string; |
| 146 | runtime: { ydoc: Y.Doc; provider: YProvider } & TRuntime; |
| 147 | } | null>(null); |
| 148 | const [hasCachedBodyState, setHasCachedBodyState] = useState(() => ({ key: sessionKey, value: false })); |
| 149 | const [bootstrapState, setBootstrapState] = useState<{ key: string; status: DocSyncBootstrapStatus }>(() => ({ |
| 150 | key: sessionKey, |
| 151 | status: "pending", |
| 152 | })); |
| 153 | |
| 154 | const title = titleState.key === sessionKey ? titleState.title : initialTitle; |
| 155 | const runtime = enabled && runtimeState?.key === sessionKey ? runtimeState.runtime : null; |
| 156 | const hasCachedBody = enabled && hasCachedBodyState.key === sessionKey ? hasCachedBodyState.value : false; |
| 157 | const bootstrapStatus = enabled && bootstrapState.key === sessionKey ? bootstrapState.status : "pending"; |
| 158 | |
| 159 | const readInitialTitle = useEffectEvent(() => initialTitle); |
| 160 | const emitTitleChange = useEffectEvent((nextTitle: string) => { |
| 161 | onTitleChange?.(nextTitle); |
| 162 | }); |
| 163 | const emitProvider = useEffectEvent((provider: YProvider | null) => { |
| 164 | onProvider?.(provider); |
| 165 | }); |
| 166 | const readWantsConnection = useEffectEvent(() => wantsConnection); |
| 167 | const readOnline = useEffectEvent(() => online); |
| 168 | const projectRoots = useEffectEvent((ydoc: Y.Doc) => roots(ydoc)); |
| 169 | const runtimeHasBody = useEffectEvent((runtime: TRuntime) => hasBody(runtime)); |
| 170 | const lastRefreshAttemptedFor = useRef<string | null>(null); |
| 171 | |
| 172 | useEffect(() => { |
| 173 | if (!enabled) { |
| 174 | return; |
| 175 | } |
| 176 | |
| 177 | const currentSessionKey = sessionKey; |
| 178 | const ydoc = new Y.Doc(); |
| 179 | const idb = new IndexeddbPersistence(cacheKey, ydoc); |
| 180 | const projected = projectRoots(ydoc); |
| 181 | const titleText = ydoc.getText(YJS_PAGE_TITLE); |
| 182 | const wsProvider = new YProvider(window.location.host, pageId, ydoc, { |
| 183 | party: "doc-sync", |
| 184 | connect: false, |
| 185 | params: shareToken ? () => ({ share: shareToken }) : () => ({ token: useAuthStore.getState().accessToken || "" }), |
| 186 | }); |
| 187 | |
| 188 | let mounted = true; |
| 189 | let seededTitle = false; |
| 190 | let seedTitleTimeout: number | null = null; |
| 191 | |
| 192 | const titleObserver = () => { |
| 193 | if (!mounted) return; |
| 194 | const nextTitle = titleText.toString(); |
| 195 | setTitleState({ key: currentSessionKey, title: nextTitle }); |
| 196 | emitTitleChange(nextTitle); |
| 197 | }; |
| 198 | titleText.observe(titleObserver); |
| 199 | |
| 200 | const maybeSeedTitle = () => { |
| 201 | if (!mounted || seededTitle) return; |
| 202 | seededTitle = true; |
| 203 | if (seedTitleTimeout !== null) { |
| 204 | window.clearTimeout(seedTitleTimeout); |
| 205 | seedTitleTimeout = null; |
| 206 | } |
| 207 | const seed = readInitialTitle(); |
| 208 | if (titleText.length === 0 && seed) { |
| 209 | titleText.insert(0, seed); |
| 210 | } |
| 211 | }; |
| 212 | |
| 213 | const handleProviderSync = (isSynced: boolean) => { |
| 214 | if (!isSynced) return; |
| 215 | maybeSeedTitle(); |
| 216 | docCache.mark(pageId); |
| 217 | }; |
| 218 | wsProvider.on("sync", handleProviderSync); |
| 219 | |
| 220 | const handleConnectionClose = () => { |
| 221 | if (!mounted) return; |
| 222 | const decision = decideDocSyncRefresh({ |
| 223 | isOnline: readOnline(), |
| 224 | isProviderActive: readWantsConnection(), |
| 225 | hasShareToken: !!shareToken, |
| 226 | currentAccessToken: useAuthStore.getState().accessToken, |
| 227 | lastRefreshAttemptedFor: lastRefreshAttemptedFor.current, |
| 228 | }); |
| 229 | if (decision.kind !== "refresh") return; |
| 230 | // Mark before the async call. Do NOT overwrite on success: a future |
| 231 | // expiration of the rotated token must remain refreshable. |
| 232 | lastRefreshAttemptedFor.current = decision.tokenAtAttempt; |
| 233 | void refreshSession(); |
| 234 | }; |
| 235 | wsProvider.on("connection-close", handleConnectionClose); |
| 236 | |
| 237 | const handleIdbSync = () => { |
| 238 | if (!mounted) return; |
| 239 | const bodyReady = runtimeHasBody(projected); |
| 240 | |
| 241 | if (titleText.length > 0) { |
| 242 | const nextTitle = titleText.toString(); |
| 243 | setTitleState({ key: currentSessionKey, title: nextTitle }); |
| 244 | emitTitleChange(nextTitle); |
| 245 | docCache.mark(pageId); |
| 246 | } |
| 247 | |
| 248 | setHasCachedBodyState({ key: currentSessionKey, value: bodyReady }); |
| 249 | emitProvider(wsProvider); |
| 250 | setRuntimeState({ key: currentSessionKey, runtime: { ...projected, ydoc, provider: wsProvider } }); |
| 251 | |
| 252 | if (!readWantsConnection()) { |
| 253 | seedTitleTimeout = window.setTimeout(maybeSeedTitle, 2000); |
| 254 | } |
| 255 | }; |
| 256 | |
| 257 | if (idb.synced) { |
| 258 | handleIdbSync(); |
| 259 | } else { |
| 260 | idb.on("synced", handleIdbSync); |
| 261 | } |
| 262 | |
| 263 | return () => { |
| 264 | mounted = false; |
| 265 | idb.off("synced", handleIdbSync); |
| 266 | titleText.unobserve(titleObserver); |
| 267 | if (seedTitleTimeout !== null) window.clearTimeout(seedTitleTimeout); |
| 268 | wsProvider.off("sync", handleProviderSync); |
| 269 | wsProvider.off("connection-close", handleConnectionClose); |
| 270 | emitProvider(null); |
| 271 | wsProvider.destroy(); |
| 272 | idb.destroy(); |
| 273 | ydoc.destroy(); |
| 274 | }; |
| 275 | }, [enabled, pageId, shareToken, cacheKey, sessionKey]); |
| 276 | |
| 277 | const phase = deriveDocSyncPhase({ |
| 278 | hasLocalBodyState: hasCachedBody, |
| 279 | wantsConnection, |
| 280 | workspaceId, |
| 281 | bootstrapStatus, |
| 282 | }); |
| 283 | |
| 284 | const snapshotWorkspaceId = phase.snapshotFetch?.workspaceId ?? null; |
| 285 | useEffect(() => { |
| 286 | if (!runtime || !snapshotWorkspaceId) return; |
| 287 | |
| 288 | const controller = new AbortController(); |
| 289 | let cancelled = false; |
| 290 | |
| 291 | void api.pages |
| 292 | .snapshot(snapshotWorkspaceId, pageId, shareToken, controller.signal) |
| 293 | .then((result) => { |
| 294 | if (cancelled) return; |
| 295 | if (result.kind === "found") { |
| 296 | Y.applyUpdate(runtime.ydoc, new Uint8Array(result.snapshot)); |
| 297 | docCache.mark(pageId); |
| 298 | } |
| 299 | setBootstrapState({ key: sessionKey, status: "resolved" }); |
| 300 | }) |
| 301 | .catch((error) => { |
| 302 | if (controller.signal.aborted || cancelled) return; |
| 303 | reportClientError({ |
| 304 | source: errorSource, |
| 305 | error, |
| 306 | context: { |
| 307 | pageId, |
| 308 | workspaceId: snapshotWorkspaceId, |
| 309 | hasShareToken: !!shareToken, |
| 310 | }, |
| 311 | }); |
| 312 | setBootstrapState({ key: sessionKey, status: "error" }); |
| 313 | }); |
| 314 | |
| 315 | return () => { |
| 316 | cancelled = true; |
| 317 | controller.abort(); |
| 318 | }; |
| 319 | }, [runtime, snapshotWorkspaceId, pageId, shareToken, errorSource, sessionKey]); |
| 320 | |
| 321 | useEffect(() => { |
| 322 | if (!runtime) return; |
| 323 | reconcileDocSyncProvider(runtime.provider, phase.shouldConnect); |
| 324 | }, [runtime, phase.shouldConnect]); |
| 325 | |
| 326 | const onTitleInput = useCallback( |
| 327 | (event: React.ChangeEvent<HTMLTextAreaElement>) => { |
| 328 | const nextTitle = event.target.value; |
| 329 | setTitleState({ key: sessionKey, title: nextTitle }); |
| 330 | if (!runtime) return; |
| 331 | |
| 332 | const titleText = runtime.ydoc.getText(YJS_PAGE_TITLE); |
| 333 | runtime.ydoc.transact(() => { |
| 334 | titleText.delete(0, titleText.length); |
| 335 | titleText.insert(0, nextTitle); |
| 336 | }); |
| 337 | }, |
| 338 | [runtime, sessionKey], |
| 339 | ); |
| 340 | |
| 341 | const retrySnapshot = useCallback(() => { |
| 342 | setBootstrapState({ key: sessionKey, status: "pending" }); |
| 343 | }, [sessionKey]); |
| 344 | |
| 345 | if (phase.error) { |
| 346 | return { |
| 347 | kind: "error", |
| 348 | title, |
| 349 | onTitleInput, |
| 350 | onRetry: retrySnapshot, |
| 351 | }; |
| 352 | } |
| 353 | |
| 354 | if (!runtime || !phase.ready) { |
| 355 | return { |
| 356 | kind: "loading", |
| 357 | title, |
| 358 | onTitleInput, |
| 359 | }; |
| 360 | } |
| 361 | |
| 362 | return { |
| 363 | kind: "ready", |
| 364 | title, |
| 365 | onTitleInput, |
| 366 | ...runtime, |
| 367 | } as DocSyncSessionReadyState<TRuntime>; |
| 368 | } |