From 96d9cd2a6cfcde677fa3d1451afe6072ff149d0a Mon Sep 17 00:00:00 2001 From: Thomas Hart Date: Tue, 25 Aug 2026 18:02:06 +0000 Subject: [PATCH 1/2] feat: Add CockroachDB distributed SQL backend with hash-sharded keys Reuse the Postgres schema and feed SQL, then encode range hotspots, hash sharding, and SERIALIZABLE retries on 40001. --- README.md | 13 +++ src/cockroach/hotspot.ts | 141 +++++++++++++++++++++++++++ src/cockroach/retry.ts | 29 ++++++ src/cockroach/schema.ts | 10 ++ src/cockroach/store.ts | 107 +++++++++++++++++++++ src/index.ts | 10 ++ test/cockroach-store.test.ts | 179 +++++++++++++++++++++++++++++++++++ 7 files changed, 489 insertions(+) create mode 100644 src/cockroach/hotspot.ts create mode 100644 src/cockroach/retry.ts create mode 100644 src/cockroach/schema.ts create mode 100644 src/cockroach/store.ts create mode 100644 test/cockroach-store.test.ts diff --git a/README.md b/README.md index afcdc43..15f62e4 100644 --- a/README.md +++ b/README.md @@ -41,6 +41,12 @@ Polyglot persistence is the idea that one product rarely has one ideal database. - Dual-write compensation: a lost handle LWT deletes the `users_by_id` row; Cassandra has no cross-partition transaction - Fan-in as N partition reads (one `posts_by_author` partition per followee) plus a merge. An inbox table would fan out on write and snapshot follows, which would break the live-feed contract. +- Distributed SQL (CockroachDB): Postgres-compatible SQL over range-partitioned keyspaces and leaseholders +- Write hotspots: monotonically increasing PKs (`SERIAL`, packed timestamps) pin inserts to the rightmost range +- Hash-sharded indexes (`USING HASH WITH (bucket_count = n)`, power-of-two buckets) prepend a hash bucket so point writes scatter +- Ordered composite keys for prefix scans: `follows (follower_id, followee_id)` and `posts (author_id, created_at, id)` stay unhashed so `following()` and `postsByAuthor()` stay one-range +- Secondary-index hotspots: a global `created_at` index is sequential even when the PK is a UUID, so that index is the one that gets hashed +- SERIALIZABLE snapshot isolation: client retry on SQLSTATE `40001` (`restart transaction`) ## What's implemented - Project scaffold with TypeScript strict mode, Vitest, and CI @@ -63,6 +69,13 @@ Polyglot persistence is the idea that one product rarely has one ideal database. - Dual-write compensation: a lost handle LWT deletes the `users_by_id` row; Cassandra has no cross-partition transaction - Fan-in as N partition reads (one `posts_by_author` partition per followee) plus a merge. An inbox table would fan out on write and snapshot follows, which would break the live-feed contract. - Wide-column backend (Cassandra): query-first modeling, partition/clustering keys +- Distributed SQL (CockroachDB): Postgres-compatible SQL over range-partitioned keyspaces and leaseholders +- Write hotspots: monotonically increasing PKs (`SERIAL`, packed timestamps) pin inserts to the rightmost range +- Hash-sharded indexes (`USING HASH WITH (bucket_count = n)`, power-of-two buckets) prepend a hash bucket so point writes scatter +- Ordered composite keys for prefix scans: `follows (follower_id, followee_id)` and `posts (author_id, created_at, id)` stay unhashed so `following()` and `postsByAuthor()` stay one-range +- Secondary-index hotspots: a global `created_at` index is sequential even when the PK is a UUID, so that index is the one that gets hashed +- SERIALIZABLE snapshot isolation: client retry on SQLSTATE `40001` (`restart transaction`) +- Distributed SQL backend (CockroachDB): same SQL, hash-sharded point keys, ordered prefix scans, SERIALIZABLE retry ## Usage ```ts diff --git a/src/cockroach/hotspot.ts b/src/cockroach/hotspot.ts new file mode 100644 index 0000000..2760d8a --- /dev/null +++ b/src/cockroach/hotspot.ts @@ -0,0 +1,141 @@ +export type AccessShape = 'point' | 'prefix-scan' | 'monotonic-append' + +export const FEED_LAYOUT = { + usersPk: { access: 'point', sequential: true }, + postsPk: { access: 'point', sequential: true }, + followsByFollower: { access: 'prefix-scan', sequential: false }, + postsByAuthor: { access: 'prefix-scan', sequential: false }, + postsByCreatedAt: { access: 'monotonic-append', sequential: true }, +} as const + +export type FeedIndex = keyof typeof FEED_LAYOUT + +export function shouldHashShard(access: AccessShape, sequential: boolean): boolean { + if (access === 'prefix-scan') return false + if (access === 'monotonic-append') return true + return sequential +} + +export function layoutHash(index: FeedIndex): boolean { + const row = FEED_LAYOUT[index] + return shouldHashShard(row.access, row.sequential) +} + +export function fnv1a(input: string): number { + let h = 0x811c9dc5 + for (let i = 0; i < input.length; i++) { + h ^= input.charCodeAt(i) + h = Math.imul(h, 0x01000193) + } + return h >>> 0 +} + +export function hashBucket(key: string, bucketCount: number): number { + if (!Number.isInteger(bucketCount) || bucketCount < 2 || (bucketCount & (bucketCount - 1)) !== 0) { + throw new RangeError('bucketCount must be a power of two >= 2') + } + return fnv1a(key) & (bucketCount - 1) +} + +export function serialKey(n: number): string { + if (!Number.isInteger(n) || n < 0) throw new RangeError('serial must be a non-negative integer') + return n.toString(16).padStart(16, '0') +} + +export function uuidKey(seed: number): string { + const a = fnv1a(`u:${seed}`).toString(16).padStart(8, '0') + const b = fnv1a(`v:${seed}`).toString(16).padStart(8, '0') + return a + b +} + +export function hashShardedKey(key: string, bucketCount: number): string { + return `${hashBucket(key, bucketCount).toString(16).padStart(4, '0')}/${key}` +} + +export function orderedIndexKey(prefix: string, createdAt: number, id: string): string { + return `${prefix}/${serialKey(createdAt)}/${id}` +} + +export function hashShardedIndexKey(createdAt: number, id: string, bucketCount: number): string { + return hashShardedKey(orderedIndexKey('', createdAt, id), bucketCount) +} + +interface Range { + start: string + keys: string[] +} + +export class Keyspace { + private readonly ranges: Range[] = [{ start: '', keys: [] }] + + constructor(readonly splitAfter: number) { + if (!Number.isInteger(splitAfter) || splitAfter < 2) { + throw new RangeError('splitAfter must be an integer >= 2') + } + } + + rangeCount(): number { + return this.ranges.length + } + + rangeIndex(key: string): number { + let lo = 0 + let hi = this.ranges.length + while (lo < hi) { + const mid = (lo + hi) >> 1 + const start = this.ranges[mid]?.start + if (start !== undefined && start <= key) lo = mid + 1 + else hi = mid + } + return Math.max(0, lo - 1) + } + + insert(key: string): number { + let i = this.rangeIndex(key) + const range = this.ranges[i] + if (!range) throw new Error('missing range') + insertSorted(range.keys, key) + if (range.keys.length > this.splitAfter) { + this.split(i) + i = this.rangeIndex(key) + } + return i + } + + hottestShare(keys: readonly string[]): number { + if (keys.length === 0) return 0 + const counts = new Map() + let max = 0 + for (const key of keys) { + const i = this.rangeIndex(key) + const n = (counts.get(i) ?? 0) + 1 + counts.set(i, n) + if (n > max) max = n + } + return max / keys.length + } + + private split(i: number): void { + const range = this.ranges[i] + if (!range) return + const mid = Math.floor(range.keys.length / 2) + const rightKeys = range.keys.slice(mid) + const rightStart = rightKeys[0] + if (rightStart === undefined) return + range.keys = range.keys.slice(0, mid) + this.ranges.splice(i + 1, 0, { start: rightStart, keys: rightKeys }) + } +} + +function insertSorted(list: string[], key: string): void { + let lo = 0 + let hi = list.length + while (lo < hi) { + const mid = (lo + hi) >> 1 + const at = list[mid] + if (at !== undefined && at < key) lo = mid + 1 + else hi = mid + } + if (list[lo] === key) return + list.splice(lo, 0, key) +} diff --git a/src/cockroach/retry.ts b/src/cockroach/retry.ts new file mode 100644 index 0000000..95bc38d --- /dev/null +++ b/src/cockroach/retry.ts @@ -0,0 +1,29 @@ +export const SERIALIZATION_FAILURE = '40001' + +export interface RetryPolicy { + attempts: number +} + +export const DEFAULT_RETRY: RetryPolicy = { attempts: 8 } + +export function retryableSqlError(err: unknown): boolean { + if (typeof err !== 'object' || err === null || !('code' in err)) return false + return (err as { code: unknown }).code === SERIALIZATION_FAILURE +} + +export async function withSerializableRetry( + op: () => Promise, + policy: RetryPolicy = DEFAULT_RETRY, +): Promise { + const attempts = Math.max(1, Math.floor(policy.attempts)) + let last: unknown + for (let i = 0; i < attempts; i++) { + try { + return await op() + } catch (err) { + last = err + if (!retryableSqlError(err) || i === attempts - 1) throw err + } + } + throw last +} diff --git a/src/cockroach/schema.ts b/src/cockroach/schema.ts new file mode 100644 index 0000000..dbd02e3 --- /dev/null +++ b/src/cockroach/schema.ts @@ -0,0 +1,10 @@ +export { SCHEMA_STATEMENTS, SQL } from '../postgres/schema' + +export const HASH_BUCKETS = 16 + +// Follows PK and posts_author_timeline stay ordered so a prefix scan hits one range. +export const CRDB_HASH_STATEMENTS = [ + `ALTER TABLE users ALTER PRIMARY KEY USING COLUMNS (id) USING HASH WITH (bucket_count = ${HASH_BUCKETS})`, + `ALTER TABLE posts ALTER PRIMARY KEY USING COLUMNS (id) USING HASH WITH (bucket_count = ${HASH_BUCKETS})`, + `CREATE INDEX IF NOT EXISTS posts_created_at_hash_idx ON posts (created_at DESC) USING HASH WITH (bucket_count = 8)`, +] as const diff --git a/src/cockroach/store.ts b/src/cockroach/store.ts new file mode 100644 index 0000000..b093153 --- /dev/null +++ b/src/cockroach/store.ts @@ -0,0 +1,107 @@ +import type { Page, Post, PostId, User, UserId } from '../domain' +import { PostgresStore, type SqlQuery } from '../postgres/store' +import type { ActivityStore, CreateUserInput, PublishInput } from '../store' +import { CRDB_HASH_STATEMENTS } from './schema' +import { DEFAULT_RETRY, withSerializableRetry, type RetryPolicy } from './retry' + +export { CRDB_HASH_STATEMENTS, HASH_BUCKETS } from './schema' +export { + DEFAULT_RETRY, + SERIALIZATION_FAILURE, + retryableSqlError, + withSerializableRetry, +} from './retry' +export type { RetryPolicy } from './retry' +export { + FEED_LAYOUT, + Keyspace, + fnv1a, + hashBucket, + hashShardedIndexKey, + hashShardedKey, + layoutHash, + orderedIndexKey, + serialKey, + shouldHashShard, + uuidKey, +} from './hotspot' +export type { AccessShape, FeedIndex } from './hotspot' + +export class CockroachStore implements ActivityStore { + private constructor( + private readonly inner: PostgresStore, + private readonly policy: RetryPolicy, + ) {} + + static attach(sql: SqlQuery, policy: RetryPolicy = DEFAULT_RETRY): CockroachStore { + return new CockroachStore(PostgresStore.attach(sql), policy) + } + + static async migrate(sql: SqlQuery, applyHashLayout = false): Promise { + await PostgresStore.migrate(sql) + if (!applyHashLayout) return + for (const statement of CRDB_HASH_STATEMENTS) { + await sql.query(statement) + } + } + + static async create( + sql: SqlQuery, + policy: RetryPolicy = DEFAULT_RETRY, + ): Promise { + await CockroachStore.migrate(sql) + return CockroachStore.attach(sql, policy) + } + + createUser(input: CreateUserInput, now = Date.now()): Promise { + return this.retry(() => this.inner.createUser(input, now)) + } + + getUser(id: UserId): Promise { + return this.retry(() => this.inner.getUser(id)) + } + + getUserByHandle(handle: string): Promise { + return this.retry(() => this.inner.getUserByHandle(handle)) + } + + follow(followerId: UserId, followeeId: UserId): Promise { + return this.retry(() => this.inner.follow(followerId, followeeId)) + } + + unfollow(followerId: UserId, followeeId: UserId): Promise { + return this.retry(() => this.inner.unfollow(followerId, followeeId)) + } + + isFollowing(followerId: UserId, followeeId: UserId): Promise { + return this.retry(() => this.inner.isFollowing(followerId, followeeId)) + } + + following(userId: UserId): Promise { + return this.retry(() => this.inner.following(userId)) + } + + followers(userId: UserId): Promise { + return this.retry(() => this.inner.followers(userId)) + } + + publish(input: PublishInput, now = Date.now()): Promise { + return this.retry(() => this.inner.publish(input, now)) + } + + getPost(id: PostId): Promise { + return this.retry(() => this.inner.getPost(id)) + } + + postsByAuthor(authorId: UserId, page?: Page): Promise { + return this.retry(() => this.inner.postsByAuthor(authorId, page)) + } + + feed(userId: UserId, page?: Page): Promise { + return this.retry(() => this.inner.feed(userId, page)) + } + + private retry(op: () => Promise): Promise { + return withSerializableRetry(op, this.policy) + } +} diff --git a/src/index.ts b/src/index.ts index b10aa7a..868aa1a 100644 --- a/src/index.ts +++ b/src/index.ts @@ -55,3 +55,13 @@ export type { SelectOpts, TableSchema, } from './cassandra/store' + +export { + CockroachStore, + CRDB_HASH_STATEMENTS, + FEED_LAYOUT, + HASH_BUCKETS, + Keyspace, + SERIALIZATION_FAILURE, + withSerializableRetry, +} from './cockroach/store' diff --git a/test/cockroach-store.test.ts b/test/cockroach-store.test.ts new file mode 100644 index 0000000..1737cb8 --- /dev/null +++ b/test/cockroach-store.test.ts @@ -0,0 +1,179 @@ +import { PGlite } from '@electric-sql/pglite' +import { afterAll, beforeAll, describe, expect, it } from 'vitest' +import { + CockroachStore, + CRDB_HASH_STATEMENTS, + Keyspace, + SERIALIZATION_FAILURE, + hashShardedIndexKey, + hashShardedKey, + layoutHash, + orderedIndexKey, + serialKey, + shouldHashShard, + uuidKey, + withSerializableRetry, +} from '../src/cockroach/store' +import { defineStoreContract } from './contract' + +const db = new PGlite() + +beforeAll(async () => { + await db.waitReady + await CockroachStore.migrate(db) +}, 30_000) + +afterAll(async () => { + await db.close() +}) + +defineStoreContract('cockroach', async () => { + await db.query('TRUNCATE TABLE users CASCADE') + return CockroachStore.attach(db) +}) + +describe('cockroach key layout', () => { + it('hash-shards sequential point lookups and monotonic indexes, not prefix scans', () => { + expect(shouldHashShard('point', true)).toBe(true) + expect(shouldHashShard('point', false)).toBe(false) + expect(shouldHashShard('prefix-scan', true)).toBe(false) + expect(shouldHashShard('monotonic-append', false)).toBe(true) + expect(layoutHash('usersPk')).toBe(true) + expect(layoutHash('postsPk')).toBe(true) + expect(layoutHash('followsByFollower')).toBe(false) + expect(layoutHash('postsByAuthor')).toBe(false) + expect(layoutHash('postsByCreatedAt')).toBe(true) + }) + + it('declares HASH alters for users and posts PKs and the global created_at index', () => { + const sql = CRDB_HASH_STATEMENTS.join('\n') + expect(sql).toContain('ALTER TABLE users ALTER PRIMARY KEY') + expect(sql).toContain('ALTER TABLE posts ALTER PRIMARY KEY') + expect(sql).toContain('posts_created_at_hash_idx') + expect(sql).not.toContain('ALTER TABLE follows') + for (const statement of CRDB_HASH_STATEMENTS) { + expect(statement).toMatch(/USING HASH WITH \(bucket_count = \d+\)/) + } + }) +}) + +describe('cockroach ranges and hotspots', () => { + it('treats an empty key list as no hotspot', () => { + expect(new Keyspace(4).hottestShare([])).toBe(0) + }) + + it('rejects a non-power-of-two bucket count', () => { + expect(() => hashShardedKey('k', 3)).toThrow(RangeError) + expect(() => new Keyspace(1)).toThrow(RangeError) + }) + + it('pins sequential primary keys onto the rightmost range', () => { + const space = new Keyspace(8) + const keys = Array.from({ length: 40 }, (_, i) => serialKey(i)) + for (const key of keys) space.insert(key) + expect(space.rangeCount()).toBeGreaterThan(1) + expect(space.hottestShare(keys.slice(-8))).toBe(1) + }) + + it('spreads hash-sharded and uuid keys across ranges', () => { + const hashed = new Keyspace(8) + const hashedKeys = Array.from({ length: 64 }, (_, i) => hashShardedKey(uuidKey(i), 16)) + for (const key of hashedKeys) hashed.insert(key) + expect(hashed.hottestShare(hashedKeys.slice(-16))).toBeLessThan(0.5) + + const uuids = new Keyspace(8) + const uuidKeys = Array.from({ length: 64 }, (_, i) => uuidKey(i)) + for (const key of uuidKeys) uuids.insert(key) + expect(uuids.hottestShare(uuidKeys.slice(-16))).toBeLessThan(0.5) + }) + + it('keeps one author timeline on a contiguous span of ranges', () => { + const space = new Keyspace(8) + for (let t = 0; t < 32; t++) space.insert(orderedIndexKey('ada', t, `p${t}`)) + for (let t = 0; t < 32; t++) space.insert(orderedIndexKey('bob', t, `q${t}`)) + const ada = uniqueSorted( + Array.from({ length: 32 }, (_, t) => space.rangeIndex(orderedIndexKey('ada', t, `p${t}`))), + ) + const bob = uniqueSorted( + Array.from({ length: 32 }, (_, t) => space.rangeIndex(orderedIndexKey('bob', t, `q${t}`))), + ) + expect(contiguous(ada)).toBe(true) + expect(contiguous(bob)).toBe(true) + expect(Math.max(...ada)).toBeLessThanOrEqual(Math.min(...bob)) + }) + + it('scatters a hashed created_at index so monotonic time does not hotspot', () => { + const space = new Keyspace(8) + const keys = Array.from({ length: 48 }, (_, i) => hashShardedIndexKey(i, `p${i}`, 8)) + for (const key of keys) space.insert(key) + expect(space.rangeCount()).toBeGreaterThan(2) + const ranges = new Set(keys.map((key) => space.rangeIndex(key))) + expect(ranges.size).toBeGreaterThan(2) + }) +}) + +describe('serializable retry', () => { + it('retries 40001 and then returns', async () => { + let n = 0 + const value = await withSerializableRetry(async () => { + n += 1 + if (n < 3) throw Object.assign(new Error('restart transaction'), { code: SERIALIZATION_FAILURE }) + return 7 + }, { attempts: 5 }) + expect(value).toBe(7) + expect(n).toBe(3) + }) + + it('does not retry a constraint error', async () => { + let n = 0 + await expect( + withSerializableRetry(async () => { + n += 1 + throw Object.assign(new Error('unique'), { code: '23505' }) + }, { attempts: 5 }), + ).rejects.toMatchObject({ code: '23505' }) + expect(n).toBe(1) + }) + + it('retries a 40001 from the sql client on publish', async () => { + await db.query('TRUNCATE TABLE users CASCADE') + let fails = 1 + const sql = { + query: async >(query: string, params?: unknown[]) => { + if (query.includes('INSERT INTO posts') && fails > 0) { + fails -= 1 + throw Object.assign(new Error('restart transaction'), { code: SERIALIZATION_FAILURE }) + } + return db.query(query, params) + }, + } + const store = CockroachStore.attach(sql, { attempts: 4 }) + await store.createUser({ id: 'ada', handle: 'ada' }, 1) + const post = await store.publish({ id: 'p1', authorId: 'ada', body: 'retry' }, 2) + expect(post.id).toBe('p1') + expect(fails).toBe(0) + }) + + it('stops after the attempt budget', async () => { + let n = 0 + await expect( + withSerializableRetry(async () => { + n += 1 + throw Object.assign(new Error('restart transaction'), { code: SERIALIZATION_FAILURE }) + }, { attempts: 3 }), + ).rejects.toMatchObject({ code: SERIALIZATION_FAILURE }) + expect(n).toBe(3) + }) +}) + +function uniqueSorted(values: number[]): number[] { + return [...new Set(values)].sort((a, b) => a - b) +} + +function contiguous(indexes: number[]): boolean { + if (indexes.length <= 1) return true + const last = indexes[indexes.length - 1] + const first = indexes[0] + if (first === undefined || last === undefined) return false + return last - first === indexes.length - 1 +} From 91d9bd0ed0a51b44e2aaf89372894c198e3796b1 Mon Sep 17 00:00:00 2001 From: Thomas Hart Date: Tue, 25 Aug 2026 18:16:51 +0000 Subject: [PATCH 2/2] fix: keep users PK unique for FKs and drop leftover posts unique ALTER PRIMARY KEY USING HASH on users left a unique secondary index that FKs still need, so users stay on a plain PK. posts still hash-shard, then drop posts_id_key so inserts do not append to an unhashed unique. --- README.md | 10 +++ src/cockroach/hotspot.ts | 2 +- src/cockroach/retry.ts | 11 +++- src/cockroach/schema.ts | 4 +- src/cockroach/store.ts | 1 + test/cockroach-store.test.ts | 122 +++++++++++++++++++++++++++++++---- 6 files changed, 135 insertions(+), 15 deletions(-) diff --git a/README.md b/README.md index 15f62e4..03c582d 100644 --- a/README.md +++ b/README.md @@ -47,6 +47,11 @@ Polyglot persistence is the idea that one product rarely has one ideal database. - Ordered composite keys for prefix scans: `follows (follower_id, followee_id)` and `posts (author_id, created_at, id)` stay unhashed so `following()` and `postsByAuthor()` stay one-range - Secondary-index hotspots: a global `created_at` index is sequential even when the PK is a UUID, so that index is the one that gets hashed - SERIALIZABLE snapshot isolation: client retry on SQLSTATE `40001` (`restart transaction`) +- Write hotspots: monotonically increasing keys (`serialKey` in this lab, packed timestamps) pin inserts to the rightmost range. The live tables use TEXT ids assigned by the caller, not SERIAL columns +- Users keep a plain primary key because `follows` and `posts` reference `users(id)`. Hash-sharding that PK would leave a unique secondary index on `id`, which is still a sequential hotspot if ids pack +- posts PK is hash-sharded, then the leftover unique on `posts(id)` from `ALTER PRIMARY KEY` is dropped so inserts do not keep appending to an unhashed unique +- Secondary-index hotspots: a global `created_at` index is sequential even when the PK is hashed, so that index is the one that gets hashed + ## What's implemented - Project scaffold with TypeScript strict mode, Vitest, and CI @@ -76,6 +81,11 @@ Polyglot persistence is the idea that one product rarely has one ideal database. - Secondary-index hotspots: a global `created_at` index is sequential even when the PK is a UUID, so that index is the one that gets hashed - SERIALIZABLE snapshot isolation: client retry on SQLSTATE `40001` (`restart transaction`) - Distributed SQL backend (CockroachDB): same SQL, hash-sharded point keys, ordered prefix scans, SERIALIZABLE retry +- Write hotspots: monotonically increasing keys (`serialKey` in this lab, packed timestamps) pin inserts to the rightmost range. The live tables use TEXT ids assigned by the caller, not SERIAL columns +- Users keep a plain primary key because `follows` and `posts` reference `users(id)`. Hash-sharding that PK would leave a unique secondary index on `id`, which is still a sequential hotspot if ids pack +- posts PK is hash-sharded, then the leftover unique on `posts(id)` from `ALTER PRIMARY KEY` is dropped so inserts do not keep appending to an unhashed unique +- Secondary-index hotspots: a global `created_at` index is sequential even when the PK is hashed, so that index is the one that gets hashed +- Distributed SQL backend (CockroachDB): same SQL, hash-sharded posts PK, ordered prefix scans, SERIALIZABLE retry ## Usage ```ts diff --git a/src/cockroach/hotspot.ts b/src/cockroach/hotspot.ts index 2760d8a..26b8d15 100644 --- a/src/cockroach/hotspot.ts +++ b/src/cockroach/hotspot.ts @@ -1,7 +1,7 @@ export type AccessShape = 'point' | 'prefix-scan' | 'monotonic-append' export const FEED_LAYOUT = { - usersPk: { access: 'point', sequential: true }, + usersPk: { access: 'point', sequential: false }, postsPk: { access: 'point', sequential: true }, followsByFollower: { access: 'prefix-scan', sequential: false }, postsByAuthor: { access: 'prefix-scan', sequential: false }, diff --git a/src/cockroach/retry.ts b/src/cockroach/retry.ts index 95bc38d..773a04d 100644 --- a/src/cockroach/retry.ts +++ b/src/cockroach/retry.ts @@ -2,9 +2,11 @@ export const SERIALIZATION_FAILURE = '40001' export interface RetryPolicy { attempts: number + baseDelayMs?: number + sleep?: (ms: number) => Promise } -export const DEFAULT_RETRY: RetryPolicy = { attempts: 8 } +export const DEFAULT_RETRY: RetryPolicy = { attempts: 8, baseDelayMs: 25 } export function retryableSqlError(err: unknown): boolean { if (typeof err !== 'object' || err === null || !('code' in err)) return false @@ -16,6 +18,8 @@ export async function withSerializableRetry( policy: RetryPolicy = DEFAULT_RETRY, ): Promise { const attempts = Math.max(1, Math.floor(policy.attempts)) + const baseDelayMs = policy.baseDelayMs ?? 0 + const sleep = policy.sleep ?? delay let last: unknown for (let i = 0; i < attempts; i++) { try { @@ -23,7 +27,12 @@ export async function withSerializableRetry( } catch (err) { last = err if (!retryableSqlError(err) || i === attempts - 1) throw err + if (baseDelayMs > 0) await sleep(baseDelayMs * 2 ** i) } } throw last } + +function delay(ms: number): Promise { + return new Promise((resolve) => setTimeout(resolve, ms)) +} diff --git a/src/cockroach/schema.ts b/src/cockroach/schema.ts index dbd02e3..441b400 100644 --- a/src/cockroach/schema.ts +++ b/src/cockroach/schema.ts @@ -4,7 +4,7 @@ export const HASH_BUCKETS = 16 // Follows PK and posts_author_timeline stay ordered so a prefix scan hits one range. export const CRDB_HASH_STATEMENTS = [ - `ALTER TABLE users ALTER PRIMARY KEY USING COLUMNS (id) USING HASH WITH (bucket_count = ${HASH_BUCKETS})`, `ALTER TABLE posts ALTER PRIMARY KEY USING COLUMNS (id) USING HASH WITH (bucket_count = ${HASH_BUCKETS})`, - `CREATE INDEX IF NOT EXISTS posts_created_at_hash_idx ON posts (created_at DESC) USING HASH WITH (bucket_count = 8)`, + `DROP INDEX IF EXISTS posts_id_key CASCADE`, + `CREATE INDEX IF NOT EXISTS posts_created_at_hash_idx ON posts (created_at DESC) USING HASH WITH (bucket_count = ${HASH_BUCKETS})`, ] as const diff --git a/src/cockroach/store.ts b/src/cockroach/store.ts index b093153..09bc0c6 100644 --- a/src/cockroach/store.ts +++ b/src/cockroach/store.ts @@ -49,6 +49,7 @@ export class CockroachStore implements ActivityStore { sql: SqlQuery, policy: RetryPolicy = DEFAULT_RETRY, ): Promise { + // PGlite cannot run USING HASH; a cluster calls migrate(sql, true) then attach. await CockroachStore.migrate(sql) return CockroachStore.attach(sql, policy) } diff --git a/test/cockroach-store.test.ts b/test/cockroach-store.test.ts index 1737cb8..91d482c 100644 --- a/test/cockroach-store.test.ts +++ b/test/cockroach-store.test.ts @@ -3,6 +3,7 @@ import { afterAll, beforeAll, describe, expect, it } from 'vitest' import { CockroachStore, CRDB_HASH_STATEMENTS, + HASH_BUCKETS, Keyspace, SERIALIZATION_FAILURE, hashShardedIndexKey, @@ -14,6 +15,7 @@ import { uuidKey, withSerializableRetry, } from '../src/cockroach/store' +import { SCHEMA_STATEMENTS } from '../src/postgres/schema' import { defineStoreContract } from './contract' const db = new PGlite() @@ -38,22 +40,28 @@ describe('cockroach key layout', () => { expect(shouldHashShard('point', false)).toBe(false) expect(shouldHashShard('prefix-scan', true)).toBe(false) expect(shouldHashShard('monotonic-append', false)).toBe(true) - expect(layoutHash('usersPk')).toBe(true) + expect(layoutHash('usersPk')).toBe(false) expect(layoutHash('postsPk')).toBe(true) expect(layoutHash('followsByFollower')).toBe(false) expect(layoutHash('postsByAuthor')).toBe(false) expect(layoutHash('postsByCreatedAt')).toBe(true) }) - it('declares HASH alters for users and posts PKs and the global created_at index', () => { + it('hash-shards posts PK and created_at, drops the leftover unique, leaves users plain', () => { const sql = CRDB_HASH_STATEMENTS.join('\n') - expect(sql).toContain('ALTER TABLE users ALTER PRIMARY KEY') + expect(sql).not.toContain('ALTER TABLE users') expect(sql).toContain('ALTER TABLE posts ALTER PRIMARY KEY') + expect(sql).toMatch(/DROP INDEX IF EXISTS posts_id_key/) expect(sql).toContain('posts_created_at_hash_idx') expect(sql).not.toContain('ALTER TABLE follows') - for (const statement of CRDB_HASH_STATEMENTS) { - expect(statement).toMatch(/USING HASH WITH \(bucket_count = \d+\)/) + const hashed = CRDB_HASH_STATEMENTS.filter((statement) => statement.includes('USING HASH')) + expect(hashed.length).toBeGreaterThan(0) + for (const statement of hashed) { + expect(statement).toContain(`USING HASH WITH (bucket_count = ${HASH_BUCKETS})`) } + const layout = postsIdLayout([...SCHEMA_STATEMENTS, ...CRDB_HASH_STATEMENTS]) + expect(layout.hashedPostsPk).toBe(true) + expect(layout.plainUniqueOnPostsId).toBe(false) }) }) @@ -103,12 +111,56 @@ describe('cockroach ranges and hotspots', () => { }) it('scatters a hashed created_at index so monotonic time does not hotspot', () => { - const space = new Keyspace(8) - const keys = Array.from({ length: 48 }, (_, i) => hashShardedIndexKey(i, `p${i}`, 8)) - for (const key of keys) space.insert(key) - expect(space.rangeCount()).toBeGreaterThan(2) - const ranges = new Set(keys.map((key) => space.rangeIndex(key))) - expect(ranges.size).toBeGreaterThan(2) + const sequential = new Keyspace(8) + const seqKeys = Array.from({ length: 40 }, (_, t) => orderedIndexKey('', t, `p${t}`)) + for (const key of seqKeys) sequential.insert(key) + expect(sequential.hottestShare(seqKeys.slice(-8))).toBe(1) + + const hashed = new Keyspace(8) + const hashedKeys = Array.from({ length: 64 }, (_, t) => + hashShardedIndexKey(t, `p${t}`, HASH_BUCKETS), + ) + for (const key of hashedKeys) hashed.insert(key) + expect(hashed.hottestShare(hashedKeys.slice(-16))).toBeLessThan(0.5) + }) +}) + +describe('cockroach migrate', () => { + it('runs schema then hash statements only when applyHashLayout is true', async () => { + const recorded: string[] = [] + const sql = { + query: async (statement: string) => { + recorded.push(statement) + return { rows: [] } + }, + } + + await CockroachStore.migrate(sql, true) + expect(recorded.slice(0, SCHEMA_STATEMENTS.length)).toEqual([...SCHEMA_STATEMENTS]) + expect(recorded.slice(SCHEMA_STATEMENTS.length)).toEqual([...CRDB_HASH_STATEMENTS]) + + recorded.length = 0 + await CockroachStore.migrate(sql) + expect(recorded).toEqual([...SCHEMA_STATEMENTS]) + expect(recorded.join('\n')).not.toContain('USING HASH') + + recorded.length = 0 + await CockroachStore.migrate(sql, false) + expect(recorded).toEqual([...SCHEMA_STATEMENTS]) + expect(recorded.join('\n')).not.toContain('USING HASH') + }) + + it('create migrates without hash layout so PGlite can attach', async () => { + const recorded: string[] = [] + const sql = { + query: async (statement: string) => { + recorded.push(statement) + return { rows: [] } + }, + } + await CockroachStore.create(sql) + expect(recorded).toEqual([...SCHEMA_STATEMENTS]) + expect(recorded.join('\n')).not.toContain('USING HASH') }) }) @@ -154,6 +206,27 @@ describe('serializable retry', () => { expect(fails).toBe(0) }) + it('backs off exponentially between 40001 retries', async () => { + const delays: number[] = [] + let n = 0 + const value = await withSerializableRetry( + async () => { + n += 1 + if (n < 4) throw Object.assign(new Error('restart transaction'), { code: SERIALIZATION_FAILURE }) + return 1 + }, + { + attempts: 5, + baseDelayMs: 10, + sleep: async (ms) => { + delays.push(ms) + }, + }, + ) + expect(value).toBe(1) + expect(delays).toEqual([10, 20, 40]) + }) + it('stops after the attempt budget', async () => { let n = 0 await expect( @@ -166,6 +239,33 @@ describe('serializable retry', () => { }) }) +function postsIdLayout(statements: readonly string[]): { + hashedPostsPk: boolean + plainUniqueOnPostsId: boolean +} { + let hashedPostsPk = false + let plainUniqueOnPostsId = false + for (const raw of statements) { + const s = raw.replace(/\s+/g, ' ') + if (/CREATE TABLE IF NOT EXISTS posts /.test(s) && /PRIMARY KEY/.test(s)) { + if (/USING HASH/.test(s)) { + hashedPostsPk = true + plainUniqueOnPostsId = false + } else { + plainUniqueOnPostsId = true + } + } + if (/ALTER TABLE posts ALTER PRIMARY KEY/.test(s) && /USING HASH/.test(s)) { + hashedPostsPk = true + plainUniqueOnPostsId = true + } + if (/DROP INDEX.*posts_id_key/.test(s) || /DROP CONSTRAINT.*posts_id_key/.test(s)) { + plainUniqueOnPostsId = false + } + } + return { hashedPostsPk, plainUniqueOnPostsId } +} + function uniqueSorted(values: number[]): number[] { return [...new Set(values)].sort((a, b) => a - b) }