-
Notifications
You must be signed in to change notification settings - Fork 0
Performance pass: stored-JSON pokes, bootstrap reuse, filter indexes, cheaper client apply #6
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
5cbec55
bbc1d2f
1869ac1
88b9ada
367adef
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -40,6 +40,38 @@ function request<T>(req: IDBRequest<T>): Promise<T> { | |
| }) | ||
| } | ||
|
|
||
| /** | ||
| * Row writes per batch. Each put structured-clones its value synchronously | ||
| * on the calling thread, so a bootstrap's worth of puts issued at once is | ||
| * one long main-thread task; batches let the page breathe between them. | ||
| */ | ||
| const WRITE_BATCH = 1_000 | ||
|
|
||
| /** | ||
| * Issues the row ops in batches, each queued from the previous batch's last | ||
| * request callback — the transaction stays active inside request callbacks, | ||
| * so every op still commits in the one transaction. | ||
| */ | ||
| function writeRows(rowStore: IDBObjectStore, ops: PokePersist['ops']): Promise<void> { | ||
| return new Promise((resolve, reject) => { | ||
| let i = 0 | ||
| const next = (): void => { | ||
| let last: IDBRequest | null = null | ||
| for (const end = Math.min(i + WRITE_BATCH, ops.length); i < end; i++) { | ||
| const op = ops[i]! | ||
| last = op.op === 'put' ? rowStore.put(op.value, [op.tbl, op.id]) : rowStore.delete([op.tbl, op.id]) | ||
| } | ||
| if (!last || i >= ops.length) { | ||
| resolve() | ||
| return | ||
| } | ||
| last.onsuccess = next | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🩺 Stability & Availability | 🟡 Minor | ⚡ Quick win 🔎 Supported by static analysis🏁 Script executed: sed -n '1,110p' packages/client/src/idb-store.ts
sed -n '165,220p' packages/client/src/idb-store.ts
rg -n 'txnDone|applyPoke|#queue|queue|rowValueSchema' packages/client/src/idb-store.ts packages/client/src/client.ts packages/protocol/src/messages.tsRepository: InfinityBowman/cf-sync Length of output: 11106 🏁 Script executed: #!/bin/bash
set -eu
printf '%s\n' '--- queue and store source ---'
sed -n '120,145p' packages/client/src/idb-store.ts
sed -n '175,205p' packages/client/src/idb-store.ts
printf '%s\n' '--- protocol validation ---'
sed -n '170,210p' packages/protocol/src/messages.ts
printf '%s\n' '--- relevant tests ---'
rg -n -C 5 'writeRows|applyPoke|non-clone|clone|DataClone|1000|IndexedDBSyncStore|PokePersist' packages --glob '*test*' --glob '*spec*'
printf '%s\n' '--- PR diff summary and changed hunks ---'
git diff --stat e76acc15dc47b1db68c07c7533bd7bf4f51216da 1869ac18fc8c483540b2cd2966b247363c5999ae -- packages/client/src/idb-store.ts packages/protocol/src/messages.ts
git diff --unified=25 e76acc15dc47b1db68c07c7533bd7bf4f51216da 1869ac18fc8c483540b2cd2966b247363c5999ae -- packages/client/src/idb-store.tsRepository: InfinityBowman/cf-sync Length of output: 31072 🏁 Script executed: #!/bin/bash
set -eu
printf '%s\n' '--- PokePersist and client application path ---'
rg -n -C 8 'export (interface|type) PokePersist|PokePersist|patchOpSchema|safeParse|parse\(' packages/client/src packages/protocol/src packages/server/src
printf '%s\n' '--- exact client poke handling ---'
sed -n '1260,1420p' packages/client/src/client.ts
printf '%s\n' '--- protocol/server serialization references ---'
rg -n -C 8 'JSON\.stringify|JSON\.parse|patch|poke|value' packages/server/src packages/protocol/src | head -240
printf '%s\n' '--- package versions relevant to IndexedDB test implementation ---'
rg -n -C 3 'fake-indexeddb|typescript|vitest' package.json packages/client/package.json pnpm-lock.yaml package-lock.json yarn.lock 2>/dev/null | head -160Repository: InfinityBowman/cf-sync Length of output: 42187 🌐 Web query:
💡 Result: 🏁 Script executed: #!/bin/bash
set -eu
printf '%s\n' '--- exact poke assembly and persistence call ---'
sed -n '1310,1415p' packages/client/src/client.ts
printf '%s\n' '--- protocol message definitions around poke completion ---'
sed -n '210,285p' packages/protocol/src/messages.ts
printf '%s\n' '--- all construction sites for PokePersist/applyPoke ---'
rg -n -C 6 'applyPoke\\(|ops:.*patch|patch\\.map|PokePersist' packages/client/src packages/client/test packages/server/srcRepository: InfinityBowman/cf-sync Length of output: 6824 🏁 Script executed: #!/bin/bash
set -eu
printf '%s\n' '--- exact poke assembly and persistence call ---'
sed -n '1310,1415p' packages/client/src/client.ts
printf '%s\n' '--- protocol message definitions around poke completion ---'
sed -n '210,285p' packages/protocol/src/messages.ts
printf '%s\n' '--- all construction sites for PokePersist/applyPoke ---'
rg -n -C 6 'applyPoke\(|ops:.*patch|patch\.map|PokePersist' packages/client/src packages/client/test packages/server/srcRepository: InfinityBowman/cf-sync Length of output: 31354 Reject synchronous When operation 1001 contains a non-cloneable value, Incoming JSON pokes cannot contain non-cloneable values. Direct callers of the exported Suggested fix let i = 0
const next = (): void => {
- let last: IDBRequest | null = null
- for (const end = Math.min(i + WRITE_BATCH, ops.length); i < end; i++) {
- const op = ops[i]!
- last = op.op === 'put' ? rowStore.put(op.value, [op.tbl, op.id]) : rowStore.delete([op.tbl, op.id])
- }
- if (!last || i >= ops.length) {
- resolve()
- return
+ try {
+ let last: IDBRequest | null = null
+ for (const end = Math.min(i + WRITE_BATCH, ops.length); i < end; i++) {
+ const op = ops[i]!
+ last = op.op === 'put' ? rowStore.put(op.value, [op.tbl, op.id]) : rowStore.delete([op.tbl, op.id])
+ }
+ if (!last || i >= ops.length) {
+ resolve()
+ return
+ }
+ last.onsuccess = next
+ last.onerror = () => reject(last!.error ?? new Error('IndexedDB request failed'))
+ } catch (err) {
+ rowStore.transaction.abort()
+ reject(err)
}
- last.onsuccess = next
- last.onerror = () => reject(last!.error ?? new Error('IndexedDB request failed'))
}🤖 Prompt for AI Agents |
||
| last.onerror = () => reject(last!.error ?? new Error('IndexedDB request failed')) | ||
| } | ||
| next() | ||
| }) | ||
| } | ||
|
|
||
| function txnDone(txn: IDBTransaction): Promise<void> { | ||
| return new Promise((resolve, reject) => { | ||
| txn.oncomplete = () => resolve() | ||
|
|
@@ -155,10 +187,7 @@ export class IndexedDBSyncStore implements SyncStore { | |
| if (!subsumed) { | ||
| const rowStore = txn.objectStore(ROWS) | ||
| if (update.clear) rowStore.clear() | ||
| for (const op of update.ops) { | ||
| if (op.op === 'put') rowStore.put(op.value, [op.tbl, op.id]) | ||
| else rowStore.delete([op.tbl, op.id]) | ||
| } | ||
| await writeRows(rowStore, update.ops) | ||
| metaStore.put({ schemaVersion: update.schemaVersion, cursor: update.cursor } satisfies MetaRecord, 'state') | ||
| } | ||
| txn.objectStore(OUTBOX).put( | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,139 @@ | ||
| /** | ||
| * Client-side cost of a large workspace (~60k rows). Opt-in, never part of | ||
| * `pnpm test`: | ||
| * | ||
| * CF_SYNC_BENCH=1 pnpm vitest run test/bench | ||
| * | ||
| * Frames are built the way the server chunks them, then fed through a fake | ||
| * socket into a SyncClient with real TanStack DB collections. | ||
| */ | ||
| import { MAX_PART_PATCH_BYTES, chunkBySize, jsonByteSize, serverMsgSchema } from '@cf-sync/protocol/internal' | ||
| import { it } from 'vitest' | ||
| import { SyncClient } from '../../src/client' | ||
| import { createCollections } from '../../src/collection' | ||
| import { MemorySyncStore } from '../../src/store' | ||
| import { FakeSocket, flushMicrotasks } from '../fake-socket' | ||
| import { testApp } from '../test-schema' | ||
|
|
||
| const ROWS = 60_000 | ||
| const CHECKLISTS = 840 | ||
|
|
||
| function uuid(n: number): string { | ||
| return `00000000-0000-4000-8000-${n.toString(16).padStart(12, '0')}` | ||
| } | ||
|
|
||
| type Op = { op: 'clear' } | { op: 'put'; tbl: string; id: string; value: Record<string, unknown> } | ||
|
|
||
| function snapshotOps(): Op[] { | ||
| const ops: Op[] = [{ op: 'clear' }] | ||
| for (let c = 0; c < CHECKLISTS; c++) { | ||
| ops.push({ op: 'put', tbl: 'notes', id: uuid(c), value: { id: uuid(c), studyId: uuid(c >> 2), type: 'ROB2' } }) | ||
| } | ||
| for (let i = 0; i < ROWS; i++) { | ||
| const checklistId = uuid(i % CHECKLISTS) | ||
| const key = `domain${i % 7}.q${i % 13}` | ||
| const id = `${checklistId}:${key}:${i}` | ||
| ops.push({ | ||
| op: 'put', | ||
| tbl: 'todos', | ||
| id, | ||
| value: { id, studyId: uuid((i % CHECKLISTS) >> 2), checklistId, key, value: i % 3 ? 'Y' : 'PN' }, | ||
| }) | ||
| } | ||
| // The server sends a snapshot ORDER BY tbl, id. | ||
| const [clear, ...puts] = ops | ||
| puts.sort((a, b) => { | ||
| const x = a as Extract<Op, { op: 'put' }> | ||
| const y = b as Extract<Op, { op: 'put' }> | ||
| return x.tbl < y.tbl ? -1 : x.tbl > y.tbl ? 1 : x.id < y.id ? -1 : x.id > y.id ? 1 : 0 | ||
| }) | ||
| return [clear!, ...puts] | ||
| } | ||
|
|
||
| function snapshotFrames(clientId: string): string[] { | ||
| const ops = snapshotOps() | ||
| const pokeId = 'bench-poke' | ||
| const frames = [JSON.stringify({ type: 'pokeStart', pokeId, baseCursor: null })] | ||
| const chunks = chunkBySize(ops, { maxBytes: MAX_PART_PATCH_BYTES, sizeOf: jsonByteSize }) | ||
| let sent = 0 | ||
| chunks.forEach((patch, i) => { | ||
| sent += patch.length | ||
| const part: Record<string, unknown> = { type: 'pokePart', pokeId, patch, remaining: ops.length - sent } | ||
| if (i === 0) part.lastMutationIdChanges = { [clientId]: 0 } | ||
| frames.push(JSON.stringify(part)) | ||
| }) | ||
| frames.push( | ||
| JSON.stringify({ type: 'pokeEnd', pokeId, cursor: { backendId: 'b1', version: 1 }, pageInfo: { more: false } }), | ||
| ) | ||
| return frames | ||
| } | ||
|
|
||
| function time(fn: () => void): number { | ||
| const t0 = performance.now() | ||
| fn() | ||
| return Math.round(performance.now() - t0) | ||
| } | ||
|
|
||
| async function timeAsync(fn: () => Promise<void>): Promise<number> { | ||
| const t0 = performance.now() | ||
| await fn() | ||
| return Math.round(performance.now() - t0) | ||
| } | ||
|
|
||
| function makeClient(store?: MemorySyncStore) { | ||
| let socket: FakeSocket | null = null | ||
| const client = new SyncClient({ | ||
| url: 'ws://bench', | ||
| workspaceId: 'bench', | ||
| clientId: 'bench-client', | ||
| autoStart: false, | ||
| app: testApp, | ||
| ...(store ? { store } : {}), | ||
| createSocket: () => (socket = new FakeSocket()), | ||
| }) | ||
| const collections = createCollections(client, { startSync: true }) | ||
| return { client, collections, socket: () => socket! } | ||
| } | ||
|
|
||
| it('client bootstrap timings', async () => { | ||
| const frames = snapshotFrames('bench-client') | ||
| const bytes = frames.reduce((n, f) => n + f.length, 0) | ||
| const results: Record<string, string> = { 'snapshot size': `${Math.round(bytes / 1024)} KB in ${frames.length} frames` } | ||
|
|
||
| // Warm the JIT on the parse path before timing it. | ||
| for (const frame of frames.slice(0, 3)) serverMsgSchema.safeParse(JSON.parse(frame)) | ||
|
|
||
| results['JSON.parse all frames'] = `${time(() => frames.forEach((f) => JSON.parse(f)))} ms` | ||
| results['JSON.parse + serverMsgSchema.safeParse'] = `${time(() => frames.forEach((f) => serverMsgSchema.safeParse(JSON.parse(f))))} ms` | ||
|
|
||
| const runs: number[] = [] | ||
| for (let r = 0; r < 3; r++) { | ||
| const store = new MemorySyncStore() | ||
| const { client, collections, socket } = makeClient(store) | ||
| client.start() | ||
| await flushMicrotasks(20) | ||
| socket().open() | ||
| runs.push( | ||
| await timeAsync(async () => { | ||
| for (const frame of frames) socket().receiveRaw(frame) | ||
| await flushMicrotasks(20) | ||
| }), | ||
| ) | ||
| if (collections.todos.size !== ROWS) throw new Error(`expected ${ROWS} todos, got ${collections.todos.size}`) | ||
| if (r === 2) { | ||
| const hydrate = await timeAsync(async () => { | ||
| const next = makeClient(store) | ||
| next.client.start() | ||
| await next.client.whenHydrated | ||
| if (next.collections.todos.size !== ROWS) throw new Error(`hydrated ${next.collections.todos.size} todos`) | ||
| await next.client.destroy() | ||
| }) | ||
| results['hydrate 60k rows from a store into collections'] = `${hydrate} ms` | ||
| } | ||
| await client.destroy() | ||
| } | ||
| runs.sort((a, b) => a - b) | ||
| results['receive + apply bootstrap poke (median of 3)'] = `${runs[1]} ms` | ||
|
|
||
| console.log(`\n[cf-sync client bench] ${ROWS} rows\n${Object.entries(results).map(([k, v]) => ` ${k.padEnd(48)} ${v}`).join('\n')}\n`) | ||
| }) |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,17 @@ | ||
| import { describe, expect, it } from 'vitest' | ||
| import { patchOpSchema } from '../src/internal' | ||
|
|
||
| describe('patchOpSchema', () => { | ||
| it('passes a put value through by reference, with no copy', () => { | ||
| const value = { title: 'x', nested: { n: 1 } } | ||
| const parsed = patchOpSchema.parse({ op: 'put', tbl: 'todos', id: 't1', value }) | ||
| expect(parsed).toEqual({ op: 'put', tbl: 'todos', id: 't1', value }) | ||
| expect((parsed as { value: unknown }).value).toBe(value) | ||
| }) | ||
|
|
||
| it('rejects put values that are not objects', () => { | ||
| for (const value of [null, [], 'x', 1, true, undefined]) { | ||
| expect(patchOpSchema.safeParse({ op: 'put', tbl: 'todos', id: 't1', value }).success).toBe(false) | ||
| } | ||
| }) | ||
| }) |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win
Qualify the indexed-filter guarantee.
A field such as
a-bdoes not passFIELD_RE.SqlRowStore.listskips its SQL filter and index, thenWriteSetfilters the returned rows in JavaScript. State that automatic SQL indexes apply only to identifier-like field names; other fields can require a table scan.🤖 Prompt for AI Agents