Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 23 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,17 @@ 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`)
- 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
Expand All @@ -63,6 +74,18 @@ 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
- 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
Expand Down
141 changes: 141 additions & 0 deletions src/cockroach/hotspot.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,141 @@
export type AccessShape = 'point' | 'prefix-scan' | 'monotonic-append'

export const FEED_LAYOUT = {
usersPk: { access: 'point', sequential: false },
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<number, number>()
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)
}
38 changes: 38 additions & 0 deletions src/cockroach/retry.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
export const SERIALIZATION_FAILURE = '40001'

export interface RetryPolicy {
attempts: number
baseDelayMs?: number
sleep?: (ms: number) => Promise<void>
}

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
return (err as { code: unknown }).code === SERIALIZATION_FAILURE
}

export async function withSerializableRetry<T>(
op: () => Promise<T>,
policy: RetryPolicy = DEFAULT_RETRY,
): Promise<T> {
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 {
return await op()
} 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<void> {
return new Promise((resolve) => setTimeout(resolve, ms))
}
10 changes: 10 additions & 0 deletions src/cockroach/schema.ts
Original file line number Diff line number Diff line change
@@ -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 posts ALTER PRIMARY KEY USING COLUMNS (id) USING HASH WITH (bucket_count = ${HASH_BUCKETS})`,
`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
108 changes: 108 additions & 0 deletions src/cockroach/store.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,108 @@
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<void> {
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<CockroachStore> {
// PGlite cannot run USING HASH; a cluster calls migrate(sql, true) then attach.
await CockroachStore.migrate(sql)
return CockroachStore.attach(sql, policy)
}

createUser(input: CreateUserInput, now = Date.now()): Promise<User> {
return this.retry(() => this.inner.createUser(input, now))
}

getUser(id: UserId): Promise<User | null> {
return this.retry(() => this.inner.getUser(id))
}

getUserByHandle(handle: string): Promise<User | null> {
return this.retry(() => this.inner.getUserByHandle(handle))
}

follow(followerId: UserId, followeeId: UserId): Promise<void> {
return this.retry(() => this.inner.follow(followerId, followeeId))
}

unfollow(followerId: UserId, followeeId: UserId): Promise<void> {
return this.retry(() => this.inner.unfollow(followerId, followeeId))
}

isFollowing(followerId: UserId, followeeId: UserId): Promise<boolean> {
return this.retry(() => this.inner.isFollowing(followerId, followeeId))
}

following(userId: UserId): Promise<UserId[]> {
return this.retry(() => this.inner.following(userId))
}

followers(userId: UserId): Promise<UserId[]> {
return this.retry(() => this.inner.followers(userId))
}

publish(input: PublishInput, now = Date.now()): Promise<Post> {
return this.retry(() => this.inner.publish(input, now))
}

getPost(id: PostId): Promise<Post | null> {
return this.retry(() => this.inner.getPost(id))
}

postsByAuthor(authorId: UserId, page?: Page): Promise<Post[]> {
return this.retry(() => this.inner.postsByAuthor(authorId, page))
}

feed(userId: UserId, page?: Page): Promise<Post[]> {
return this.retry(() => this.inner.feed(userId, page))
}

private retry<T>(op: () => Promise<T>): Promise<T> {
return withSerializableRetry(op, this.policy)
}
}
10 changes: 10 additions & 0 deletions src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Loading
Loading