diff --git a/README.md b/README.md index b46fd6e..afcdc43 100644 --- a/README.md +++ b/README.md @@ -32,6 +32,15 @@ Polyglot persistence is the idea that one product rarely has one ideal database. - Dual-write follow: `$addToSet` / `$pull` on embedded `following[]` for the feed `$in` path; unique `{ followerId, followeeId }` edges answer `followers()`. `isFollowing()` reads the edge; `following()` and `feed()` read the embed. A missed embed is repaired on the next follow; until then the two sources can disagree. - Unique index on `handle` (`E11000` / code 11000) - Compound keyset via `$or` (`createdAt $lt`, or same `createdAt` and `_id $lt`); feed is `posts.find({ authorId: { $in: following } })`. Empty `$in` matches nothing. +- Query-first data modeling (Cassandra / CQL): one table per query, denormalized writes, no joins +- Partition key as the unit of distribution and colocation (`author_id`, `follower_id`, `handle`) +- Clustering key and `CLUSTERING ORDER BY (created_at DESC, post_id DESC)` so a partition is already a newest-first timeline +- Primary key shape `PRIMARY KEY ((partition), clustering...)` versus a single-column `PRIMARY KEY (handle)` lookup table +- Clustering tuple bound `(created_at, post_id) < cursor` for keyset paging inside one partition +- Lightweight transactions (`INSERT ... IF NOT EXISTS`) for per-partition uniqueness (user id, handle, post id) +- 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. + ## What's implemented - Project scaffold with TypeScript strict mode, Vitest, and CI @@ -45,6 +54,15 @@ Polyglot persistence is the idea that one product rarely has one ideal database. - Unique index on `handle` (`E11000` / code 11000) - Compound keyset via `$or` (`createdAt $lt`, or same `createdAt` and `_id $lt`); feed is `posts.find({ authorId: { $in: following } })`. Empty `$in` matches nothing. - Document backend (MongoDB): document shape, embedding vs referencing +- Query-first data modeling (Cassandra / CQL): one table per query, denormalized writes, no joins +- Partition key as the unit of distribution and colocation (`author_id`, `follower_id`, `handle`) +- Clustering key and `CLUSTERING ORDER BY (created_at DESC, post_id DESC)` so a partition is already a newest-first timeline +- Primary key shape `PRIMARY KEY ((partition), clustering...)` versus a single-column `PRIMARY KEY (handle)` lookup table +- Clustering tuple bound `(created_at, post_id) < cursor` for keyset paging inside one partition +- Lightweight transactions (`INSERT ... IF NOT EXISTS`) for per-partition uniqueness (user id, handle, post id) +- 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 ## Usage ```ts diff --git a/src/cassandra/memory.ts b/src/cassandra/memory.ts new file mode 100644 index 0000000..0f9ad15 --- /dev/null +++ b/src/cassandra/memory.ts @@ -0,0 +1,202 @@ +export type CqlType = 'text' | 'bigint' +export type Cell = string | number +export type Row = Record + +export interface TableSchema { + readonly name: string + readonly columns: Readonly> + readonly partition: readonly string[] + readonly clustering: readonly string[] + readonly clusteringDesc: readonly boolean[] +} + +function tbl( + name: string, + columns: Readonly>, + partition: readonly string[], + clustering: readonly string[] = [], + clusteringDesc: readonly boolean[] = [], +): TableSchema { + return { name, columns, partition, clustering, clusteringDesc } +} + +export const TABLES = { + usersById: tbl('users_by_id', { user_id: 'text', handle: 'text', created_at: 'bigint' }, ['user_id']), + usersByHandle: tbl('users_by_handle', { handle: 'text', user_id: 'text', created_at: 'bigint' }, ['handle']), + followingByUser: tbl('following_by_user', { follower_id: 'text', followee_id: 'text' }, ['follower_id'], ['followee_id'], [false]), + followersByUser: tbl('followers_by_user', { followee_id: 'text', follower_id: 'text' }, ['followee_id'], ['follower_id'], [false]), + postsByAuthor: tbl( + 'posts_by_author', + { author_id: 'text', created_at: 'bigint', post_id: 'text', body: 'text' }, + ['author_id'], + ['created_at', 'post_id'], + [true, true], + ), + postsById: tbl('posts_by_id', { post_id: 'text', author_id: 'text', body: 'text', created_at: 'bigint' }, ['post_id']), +} as const satisfies Record + +export const TABLE_LIST: readonly TableSchema[] = Object.values(TABLES) + +export function toCreateCql(t: TableSchema): string { + const pkSet = new Set([...t.partition, ...t.clustering]) + const extra = Object.keys(t.columns).filter((c) => !pkSet.has(c)) + const cols = [...t.partition, ...t.clustering, ...extra].map((name) => { + const ty = t.columns[name] + if (!ty) throw new Error(`unknown column ${name}`) + return `${name} ${ty}` + }) + const pk = + t.clustering.length === 0 + ? (t.partition[0] ?? '') + : `(${t.partition.join(', ')}), ${t.clustering.join(', ')}` + const order = t.clustering + .map((col, i) => `${col} ${t.clusteringDesc[i] ? 'DESC' : 'ASC'}`) + .join(', ') + const clustering = order ? ` WITH CLUSTERING ORDER BY (${order})` : '' + return `CREATE TABLE ${t.name} (${cols.join(', ')}, PRIMARY KEY (${pk}))${clustering}` +} + +export class InvalidQueryError extends Error { + constructor(message: string) { + super(message) + this.name = 'InvalidQueryError' + } +} + +export interface SelectOpts { + eq: Record + clusteringLt?: Record + limit?: number +} + +export class MemoryCassandra { + private readonly tables = new Map>() + private readonly schemas = new Map() + + constructor() { + for (const t of TABLE_LIST) { + this.schemas.set(t.name, t) + this.tables.set(t.name, new Map()) + } + } + + insert(table: string, row: Row, opts?: { ifNotExists?: boolean }): { applied: boolean } { + const schema = this.must(table) + const key = partitionKey(schema, row) + const copy = { ...row } + const parts = this.tables.get(table) + if (!parts) throw new InvalidQueryError(`unknown table ${table}`) + let rows = parts.get(key) + if (!rows) { + rows = [] + parts.set(key, rows) + } + const idx = rows.findIndex((existing) => sameCk(schema, existing, copy)) + if (idx >= 0) { + if (opts?.ifNotExists) return { applied: false } + rows[idx] = copy + return { applied: true } + } + let lo = 0 + let hi = rows.length + while (lo < hi) { + const mid = (lo + hi) >> 1 + const at = rows[mid] + if (at !== undefined && compareClustering(schema, at, copy) <= 0) lo = mid + 1 + else hi = mid + } + rows.splice(lo, 0, copy) + return { applied: true } + } + + select(table: string, opts: SelectOpts): Row[] { + const schema = this.must(table) + for (const col of Object.keys(opts.eq)) { + if (!schema.partition.includes(col) && !schema.clustering.includes(col)) { + throw new InvalidQueryError(`cannot restrict ${col} on ${table}`) + } + } + for (const col of schema.partition) { + if (!(col in opts.eq)) throw new InvalidQueryError(`missing partition key ${col}`) + } + const bound = opts.clusteringLt + if (bound) { + for (const col of Object.keys(bound)) { + if (!schema.clustering.includes(col)) throw new InvalidQueryError(`cannot range ${col} on ${table}`) + } + } + const part = this.tables.get(table)?.get(partitionKey(schema, opts.eq)) ?? [] + const eqCk = schema.clustering.filter((c) => c in opts.eq) + let rows = part.filter((row) => eqCk.every((c) => row[c] === opts.eq[c])) + if (bound) rows = rows.filter((row) => naturallyLess(schema, row, bound)) + if (opts.limit !== undefined) rows = rows.slice(0, Math.max(0, opts.limit)) + return rows.map((row) => ({ ...row })) + } + + delete(table: string, where: Record): void { + const schema = this.must(table) + const parts = this.tables.get(table) + if (!parts) return + const key = partitionKey(schema, where) + const rows = parts.get(key) + if (!rows) return + const ck = schema.clustering.filter((c) => c in where) + if (ck.length === 0) { + parts.delete(key) + return + } + const next = rows.filter((row) => !ck.every((c) => row[c] === where[c])) + if (next.length === 0) parts.delete(key) + else parts.set(key, next) + } + + partitionCount(table: string): number { + return this.tables.get(table)?.size ?? 0 + } + + private must(table: string): TableSchema { + const schema = this.schemas.get(table) + if (!schema) throw new InvalidQueryError(`unknown table ${table}`) + return schema + } +} + +function partitionKey(schema: TableSchema, row: Record): string { + return schema.partition + .map((col) => { + if (!(col in row)) throw new InvalidQueryError(`missing partition key ${col}`) + return String(row[col]) + }) + .join('\x1f') +} + +function sameCk(schema: TableSchema, a: Row, b: Row): boolean { + return schema.clustering.every((c) => a[c] === b[c]) +} + +function compareCell(a: Cell | undefined, b: Cell | undefined): number { + if (a === b) return 0 + if (a === undefined) return -1 + if (b === undefined) return 1 + if (typeof a === 'number' && typeof b === 'number') return a - b + return String(a) < String(b) ? -1 : 1 +} + +function compareClustering(schema: TableSchema, a: Row, b: Row): number { + for (let i = 0; i < schema.clustering.length; i++) { + const col = schema.clustering[i] + if (!col) continue + const c = compareCell(a[col], b[col]) + if (c !== 0) return schema.clusteringDesc[i] ? -c : c + } + return 0 +} + +function naturallyLess(schema: TableSchema, row: Row, bound: Record): boolean { + for (const col of schema.clustering) { + if (!(col in bound)) continue + const c = compareCell(row[col], bound[col]) + if (c !== 0) return c < 0 + } + return false +} diff --git a/src/cassandra/store.ts b/src/cassandra/store.ts new file mode 100644 index 0000000..c320cca --- /dev/null +++ b/src/cassandra/store.ts @@ -0,0 +1,177 @@ +import { + comparePosts, + normalizeBody, + normalizeHandle, + normalizeId, + pageLimit, + StoreError, + tryNormalizeHandle, + tryNormalizeId, + type Page, + type Post, + type PostId, + type User, + type UserId, +} from '../domain' +import type { ActivityStore, CreateUserInput, PublishInput } from '../store' +import { MemoryCassandra, TABLES, type Row } from './memory' + +export { InvalidQueryError, MemoryCassandra, TABLES, TABLE_LIST, toCreateCql } from './memory' +export type { Cell, Row, SelectOpts, TableSchema } from './memory' + +export class CassandraStore implements ActivityStore { + private constructor(private readonly ks: MemoryCassandra) {} + + static attach(ks: MemoryCassandra): CassandraStore { + return new CassandraStore(ks) + } + + async createUser(input: CreateUserInput, now = Date.now()): Promise { + const id = normalizeId(input.id) + const handle = normalizeHandle(input.handle) + if (!this.ks.insert(TABLES.usersById.name, { user_id: id, handle, created_at: now }, { ifNotExists: true }).applied) { + throw new StoreError('user_exists') + } + if ( + !this.ks.insert(TABLES.usersByHandle.name, { handle, user_id: id, created_at: now }, { ifNotExists: true }).applied + ) { + this.ks.delete(TABLES.usersById.name, { user_id: id }) + throw new StoreError('handle_taken') + } + return { id, handle, createdAt: now } + } + + async getUser(id: UserId): Promise { + const key = tryNormalizeId(id) + if (!key) return null + const row = this.ks.select(TABLES.usersById.name, { eq: { user_id: key } })[0] + return row ? toUser(row) : null + } + + async getUserByHandle(handle: string): Promise { + const key = tryNormalizeHandle(handle) + if (!key) return null + const row = this.ks.select(TABLES.usersByHandle.name, { eq: { handle: key } })[0] + return row ? this.getUser(String(row.user_id)) : null + } + + async follow(followerId: UserId, followeeId: UserId): Promise { + const from = await this.requireUser(followerId) + const to = await this.requireUser(followeeId) + if (from === to) throw new StoreError('self_follow') + this.ks.insert(TABLES.followingByUser.name, { follower_id: from, followee_id: to }) + this.ks.insert(TABLES.followersByUser.name, { followee_id: to, follower_id: from }) + } + + async unfollow(followerId: UserId, followeeId: UserId): Promise { + const from = await this.requireUser(followerId) + const to = await this.requireUser(followeeId) + this.ks.delete(TABLES.followingByUser.name, { follower_id: from, followee_id: to }) + this.ks.delete(TABLES.followersByUser.name, { followee_id: to, follower_id: from }) + } + + async isFollowing(followerId: UserId, followeeId: UserId): Promise { + const from = await this.requireUser(followerId) + const to = await this.requireUser(followeeId) + return this.ks.select(TABLES.followingByUser.name, { + eq: { follower_id: from, followee_id: to }, + limit: 1, + }).length > 0 + } + + async following(userId: UserId): Promise { + const id = await this.requireUser(userId) + return this.ks.select(TABLES.followingByUser.name, { eq: { follower_id: id } }).map((row) => String(row.followee_id)) + } + + async followers(userId: UserId): Promise { + const id = await this.requireUser(userId) + return this.ks.select(TABLES.followersByUser.name, { eq: { followee_id: id } }).map((row) => String(row.follower_id)) + } + + async publish(input: PublishInput, now = Date.now()): Promise { + const id = normalizeId(input.id) + const authorId = await this.requireUser(input.authorId) + const body = normalizeBody(input.body) + const row = { post_id: id, author_id: authorId, body, created_at: now } + if (!this.ks.insert(TABLES.postsById.name, row, { ifNotExists: true }).applied) { + throw new StoreError('post_exists') + } + this.ks.insert(TABLES.postsByAuthor.name, row) + return { id, authorId, body, createdAt: now } + } + + async getPost(id: PostId): Promise { + const key = tryNormalizeId(id) + if (!key) return null + const row = this.ks.select(TABLES.postsById.name, { eq: { post_id: key } })[0] + return row ? toPost(row) : null + } + + async postsByAuthor(authorId: UserId, page?: Page): Promise { + return this.timeline(await this.requireUser(authorId), page) + } + + async feed(userId: UserId, page?: Page): Promise { + const id = await this.requireUser(userId) + const followees = this.ks.select(TABLES.followingByUser.name, { eq: { follower_id: id } }) + return mergeNewest( + followees.map((row) => this.timeline(String(row.followee_id), page)), + pageLimit(page), + ) + } + + private timeline(authorId: UserId, page?: Page): Post[] { + const before = page?.before + return this.ks + .select(TABLES.postsByAuthor.name, { + eq: { author_id: authorId }, + clusteringLt: before ? { created_at: before.createdAt, post_id: before.id } : undefined, + limit: pageLimit(page), + }) + .map(toPost) + } + + private async requireUser(id: UserId): Promise { + const key = tryNormalizeId(id) + if (!key) throw new StoreError('user_not_found') + if (!this.ks.select(TABLES.usersById.name, { eq: { user_id: key } })[0]) { + throw new StoreError('user_not_found') + } + return key + } +} + +function toUser(row: Row): User { + return { id: String(row.user_id), handle: String(row.handle), createdAt: Number(row.created_at) } +} + +function toPost(row: Row): Post { + return { + id: String(row.post_id), + authorId: String(row.author_id), + body: String(row.body), + createdAt: Number(row.created_at), + } +} + +function mergeNewest(lists: Post[][], limit: number): Post[] { + const heads = lists.map(() => 0) + const out: Post[] = [] + while (out.length < limit) { + let bestI = -1 + let best: Post | undefined + for (let i = 0; i < lists.length; i++) { + const cand = lists[i]?.[heads[i] ?? 0] + if (!cand) continue + if (!best || comparePosts(cand, best) < 0) { + best = cand + bestI = i + } + } + if (!best || bestI < 0) break + out.push(best) + heads[bestI] = (heads[bestI] ?? 0) + 1 + } + return out +} diff --git a/src/index.ts b/src/index.ts index d7cd098..b10aa7a 100644 --- a/src/index.ts +++ b/src/index.ts @@ -40,3 +40,18 @@ export type { PostDoc, UserDoc, } from './mongo/store' + +export { + CassandraStore, + InvalidQueryError, + MemoryCassandra, + TABLES, + toCreateCql, +} from './cassandra/store' + +export type { + Cell, + Row, + SelectOpts, + TableSchema, +} from './cassandra/store' diff --git a/test/cassandra-store.test.ts b/test/cassandra-store.test.ts new file mode 100644 index 0000000..ea7ff03 --- /dev/null +++ b/test/cassandra-store.test.ts @@ -0,0 +1,87 @@ +import { describe, expect, it } from 'vitest' +import { + CassandraStore, + InvalidQueryError, + MemoryCassandra, + TABLES, + toCreateCql, +} from '../src/cassandra/store' +import { defineStoreContract } from './contract' + +defineStoreContract('cassandra', () => CassandraStore.attach(new MemoryCassandra())) + +describe('cassandra query-first tables', () => { + it('emits compound partition and DESC clustering CQL for the author timeline', () => { + const cql = toCreateCql(TABLES.postsByAuthor) + expect(cql).toContain('PRIMARY KEY ((author_id), created_at, post_id)') + expect(cql).toContain('CLUSTERING ORDER BY (created_at DESC, post_id DESC)') + expect(toCreateCql(TABLES.usersByHandle)).toContain('PRIMARY KEY (handle)') + expect(toCreateCql(TABLES.followingByUser)).toContain( + 'PRIMARY KEY ((follower_id), followee_id)', + ) + }) + + it('puts each author in their own partition, clustered newest-first', async () => { + const ks = new MemoryCassandra() + const store = CassandraStore.attach(ks) + await store.createUser({ id: 'ada', handle: 'ada' }, 1) + await store.createUser({ id: 'bob', handle: 'bob' }, 1) + await store.publish({ id: 'p1', authorId: 'bob', body: 'oldest' }, 1) + await store.publish({ id: 'p3', authorId: 'bob', body: 'tie-hi' }, 10) + await store.publish({ id: 'p2', authorId: 'bob', body: 'tie-lo' }, 10) + await store.publish({ id: 'a1', authorId: 'ada', body: 'ada' }, 9) + + expect(ks.partitionCount('posts_by_author')).toBe(2) + expect(ks.select('posts_by_author', { eq: { author_id: 'bob' } }).map((row) => row.post_id)).toEqual([ + 'p3', + 'p2', + 'p1', + ]) + expect(ks.select('posts_by_author', { eq: { author_id: 'ada' } }).map((row) => row.post_id)).toEqual(['a1']) + }) + + it('rejects a clustering-only read and a non-key restriction', () => { + const ks = new MemoryCassandra() + expect(() => + ks.select('posts_by_author', { eq: { post_id: 'p1' } }), + ).toThrow(InvalidQueryError) + expect(() => ks.select('posts_by_author', { eq: { body: 'x' } })).toThrow( + /cannot restrict body/, + ) + }) + + it('uses IF NOT EXISTS on the handle partition and rolls back users_by_id', async () => { + const ks = new MemoryCassandra() + const store = CassandraStore.attach(ks) + await store.createUser({ id: 'u1', handle: 'ada' }, 1) + await expect(store.createUser({ id: 'u2', handle: 'ADA' })).rejects.toMatchObject({ + code: 'handle_taken', + }) + expect(ks.select('users_by_id', { eq: { user_id: 'u2' } })).toEqual([]) + expect(ks.partitionCount('users_by_handle')).toBe(1) + expect( + ks.insert('users_by_handle', { handle: 'ada', user_id: 'x', created_at: 0 }, { + ifNotExists: true, + }).applied, + ).toBe(false) + }) + + it('pages a partition with a clustering tuple bound, not a table scan', async () => { + const ks = new MemoryCassandra() + const store = CassandraStore.attach(ks) + await store.createUser({ id: 'bob', handle: 'bob' }, 1) + await store.publish({ id: 'p1', authorId: 'bob', body: 'older-id' }, 10) + await store.publish({ id: 'p2', authorId: 'bob', body: 'newer-id' }, 10) + await store.publish({ id: 'p3', authorId: 'bob', body: 'earlier' }, 4) + + const rest = ks + .select('posts_by_author', { + eq: { author_id: 'bob' }, + clusteringLt: { created_at: 10, post_id: 'p2' }, + limit: 2, + }) + .map((row) => row.post_id) + expect(rest).toEqual(['p1', 'p3']) + }) + +})