{"_id":"@deptyped/turso-cdc-drizzle","_rev":"2-9bbf8acac4cfca48e443b8ed62d2d03f","name":"@deptyped/turso-cdc-drizzle","dist-tags":{"latest":"0.2.0"},"versions":{"0.1.0":{"name":"@deptyped/turso-cdc-drizzle","version":"0.1.0","_id":"@deptyped/turso-cdc-drizzle@0.1.0","maintainers":[{"name":"deptyped","email":"deptyped@gmail.com"}],"dist":{"shasum":"674975aa6a19cbb04e11a23666bc59bf077d6d85","tarball":"https://registry.npmjs.org/@deptyped/turso-cdc-drizzle/-/turso-cdc-drizzle-0.1.0.tgz","fileCount":15,"integrity":"sha512-aFXQYvdtRjA3iNLi+oWFDtal0FeCPcF0eyvSfgNtGI705ljaZHswM+1wYBGSRUKFijZBIr2z62Em9BSmjyXGqA==","signatures":[{"sig":"MEUCICHPuZx26LWMWo/4HLROYcYSfI4mziHB+YaJ270XMOJXAiEA6CfP9rF3gzuvC0LJOtaAFLqE+314+fhGmER/DbpUE78=","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":35172},"main":"./dist/src/index.js","type":"module","types":"./dist/src/index.d.ts","exports":{".":{"types":"./dist/src/index.d.ts","import":"./dist/src/index.js"}},"gitHead":"3acb2bd7e278057f48c518311a6d009c5ff526af","scripts":{"test":"tsx --test test/*.test.ts","build":"tsc","format":"biome check --write","format:check":"biome check"},"_npmUser":{"name":"deptyped","email":"deptyped@gmail.com"},"_npmVersion":"11.8.0","description":"[Turso CDC](https://docs.turso.tech/tursodb/cdc) tracks every INSERT, UPDATE, and DELETE in your database as events. This library exposes those events through Drizzle ORM — read them in batches, consume them as a stream, or delete them after processing. A","directories":{},"_nodeVersion":"24.13.1","_hasShrinkwrap":false,"devDependencies":{"tsx":"^4.0.0","typescript":"^7.0.0","@types/node":"^26.0.0","drizzle-orm":"rc","@biomejs/biome":"^2.0.0","@tursodatabase/database":"^0.7.0"},"peerDependencies":{"drizzle-orm":"rc"},"_npmOperationalInternal":{"tmp":"tmp/turso-cdc-drizzle_0.1.0_1784587938895_0.4238822814279506","host":"s3://npm-registry-packages-npm-production"}},"0.2.0":{"name":"@deptyped/turso-cdc-drizzle","version":"0.2.0","type":"module","main":"./dist/src/index.js","types":"./dist/src/index.d.ts","exports":{".":{"types":"./dist/src/index.d.ts","import":"./dist/src/index.js"}},"scripts":{"build":"tsc","test":"tsx --test test/*.test.ts","format":"biome check --write","format:check":"biome check"},"peerDependencies":{"drizzle-orm":"rc"},"devDependencies":{"@biomejs/biome":"^2.0.0","@tursodatabase/database":"^0.7.0","@types/node":"^26.0.0","drizzle-orm":"rc","tsx":"^4.0.0","typescript":"^7.0.0"},"gitHead":"9168d14f1b65e6559df21ad2f8ffde61c6541fc6","_id":"@deptyped/turso-cdc-drizzle@0.2.0","description":"[Turso CDC](https://docs.turso.tech/tursodb/cdc) tracks every INSERT, UPDATE, and DELETE in your database as events. This library exposes those events through Drizzle ORM — read them in batches, consume them as a stream, or delete them after processing. A","_nodeVersion":"24.13.1","_npmVersion":"11.8.0","dist":{"integrity":"sha512-0veqQVBetM8A9TE5YbC2f5j8fi6zoDbwkoj9TuNpWV5F9q5i99ZSd4Qwu+WWjaY2o7uOKzNOHpk+cdFUAeGlUw==","shasum":"3d564a56a5679c62e3ee546121866aa6a2c9bc0b","tarball":"https://registry.npmjs.org/@deptyped/turso-cdc-drizzle/-/turso-cdc-drizzle-0.2.0.tgz","fileCount":15,"unpackedSize":38833,"signatures":[{"keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U","sig":"MEYCIQC6DZ4GFcfQn5HTPh4phSAkcQukibnz2n2dkM5f8pIXRAIhAMY2NioByVuorkKItAfVH+0Bj5KTiH1iMT72orli0RJs"}]},"_npmUser":{"name":"deptyped","email":"deptyped@gmail.com"},"directories":{},"maintainers":[{"name":"deptyped","email":"deptyped@gmail.com"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages-npm-production","tmp":"tmp/turso-cdc-drizzle_0.2.0_1784644697842_0.3836806608131662"},"_hasShrinkwrap":false}},"time":{"created":"2026-07-20T22:52:18.734Z","modified":"2026-07-21T14:38:18.200Z","0.1.0":"2026-07-20T22:52:19.048Z","0.2.0":"2026-07-21T14:38:18.017Z"},"description":"[Turso CDC](https://docs.turso.tech/tursodb/cdc) tracks every INSERT, UPDATE, and DELETE in your database as events. This library exposes those events through Drizzle ORM — read them in batches, consume them as a stream, or delete them after processing. A","maintainers":[{"name":"deptyped","email":"deptyped@gmail.com"}],"readme":"# Drizzle ORM integration for [Turso CDC](https://docs.turso.tech/tursodb/cdc)\n\n[Turso CDC](https://docs.turso.tech/tursodb/cdc) tracks every INSERT, UPDATE, and DELETE in your database as events. This library exposes those events through Drizzle ORM — read them in batches, consume them as a stream, or delete them after processing. All events are typed to your Drizzle table schemas.\n\n```\nnpm install @deptyped/turso-cdc-drizzle\n```\n\nRequires `drizzle-orm` as a peer dependency.\n\n## Quick start\n\nImports, table schema, and database instance.\n\n```ts\nimport { drizzle } from 'drizzle-orm/tursodatabase/database';\nimport { sqliteTable, int, text } from 'drizzle-orm/sqlite-core';\nimport { enableCdc, getEvents, streamEvents } from '@deptyped/turso-cdc-drizzle';\n\nconst users = sqliteTable('users', {\n  id: int('id').primaryKey(),\n  name: text('name'),\n});\n\nconst db = drizzle({ client });\n\n// Enable CDC (required). By default, captures only the primary key/rowid of changed rows.\nawait enableCdc(db);\n\n// With row data capture:\nawait enableCdc(db, 'full');\n```\n\nQuery with decoded row data, filtered by kind.\n\n```ts\nawait db.insert(users).values({ id: 1, name: 'Alice' });\n\nconst events = await getEvents(db, users, { mode: 'full', limit: 10 });\n// events[0]!.after — { id: 1, name: 'Alice' }\n```\n\nWith auto-delete after read.\n\n```ts\nconst events = await getEvents(db, users, { deleteAfterRead: true, limit: 10 });\n// events are removed from the CDC table — next poll won't see them again\n```\n\nStreaming — poll for new changes.\n\n```ts\nfor await (const event of streamEvents(db, users)) {\n  console.log(event.changeType, event.rowId);\n}\n\n// With decoded data (requires enableCdc(db, 'full'))\nfor await (const event of streamEvents(db, users, { mode: 'full' })) {\n  console.log(event.after?.name);\n}\n```\n\nStreaming filtered by change kind.\n\n```ts\nfor await (const event of streamEvents(db, users, { kinds: ['DELETE'] })) {\n  console.log(event.rowId, 'was deleted');\n}\n```\n\nStreaming with auto-delete after read.\n\n```ts\nfor await (const event of streamEvents(db, users, { deleteAfterRead: true })) {\n  process(event); // events are deleted after being yielded\n}\n```\n\n## Resilience\n\n### Crash recovery with checkpoint\n\nUse `CheckpointStrategy` to persist progress. The stream calls `restore` on start (if no `afterId` is given) and `save` after each batch is fully consumed. On abort, `save` fires one final time with the last yielded event's `changeId` before exit.\n\n```ts\nimport type { CheckpointStrategy } from '@deptyped/turso-cdc-drizzle';\nimport { sql } from 'drizzle-orm';\n\nconst checkpoint: CheckpointStrategy = {\n  save: async (changeId, db) => {\n    await db.run(sql.raw(`UPDATE _cdc_cp SET change_id = ${changeId}`));\n  },\n  restore: async (db) => {\n    const row = await db.get(sql.raw(\"SELECT change_id FROM _cdc_cp\"));\n    return row?.change_id;\n  },\n};\n\nfor await (const event of streamEvents(db, users, {\n  checkpoint,\n  batchSize: 50,\n})) {\n  // process event — saved checkpoint means <50 events replay on crash\n}\n```\n\n`save` errors are silently caught — the stream continues and retries on the next batch.\n\n### Graceful shutdown\n\nPass an `AbortSignal` to stop the stream cleanly. A final checkpoint is saved before exit.\n\n```ts\nconst ac = new AbortController();\nprocess.on('SIGTERM', () => ac.abort());\n\nfor await (const event of streamEvents(db, users, {\n  signal: ac.signal,\n  checkpoint,\n})) {\n  // the last yielded event's changeId is saved before exit\n}\n```\n\n## API\n\n### `enableCdc(db, mode?)`\n\n| Param | Type | Default | Description |\n|-------|------|---------|-------------|\n| `db` | `TursoDatabaseDatabase` | — | Drizzle Turso database |\n| `mode` | `CdcMode` | `'id'` | `'id'` \\| `'before'` \\| `'after'` \\| `'full'` |\n\nRuns `PRAGMA capture_data_changes_conn`. Use `'full'` to capture row data (required for `mode: 'full'` queries). Tracked per-connection — call once per connection.\n\n### `disableCdc(db)`\n\nDisables CDC. No options.\n\n### `getEvents(db, table, opts?)`\n\nReturns `CdcEvent<TTable>[]`. COMMIT rows are filtered out automatically.\n\n| Param | Type | Description |\n|-------|------|-------------|\n| `db` | `TursoDatabaseDatabase` | Drizzle Turso database |\n| `table` | `TTable` | A Drizzle table definition |\n| `opts.afterId` | `ChangeId` | Exclusive lower bound (gt) — events after this id |\n| `opts.beforeId` | `ChangeId` | Exclusive upper bound (lt) — events before this id |\n| `opts.kinds` | `CdcChangeKind[]` | `['INSERT']` \\| `['UPDATE']` \\| `['DELETE']` |\n| `opts.mode` | `'id'` \\| `'full'` | `'full'` decodes blob data into `before`/`after` |\n| `opts.deleteAfterRead` | `boolean` | Auto-delete returned events |\n| `opts.limit` | `number` | **Required.** Max events to return |\n\n### `streamEvents(db, table, opts?)`\n\nReturns `AsyncGenerator<CdcEvent<TTable>>`. Polls every `pollIntervalMs` (default 1000). Wrap in `for await`.\n\n| Param | Type | Default | Description |\n|-------|------|---------|-------------|\n| `db` | `TursoDatabaseDatabase` | — | Drizzle Turso database |\n| `table` | `TTable` | — | A Drizzle table definition |\n| `opts.pollIntervalMs` | `number` | `1000` | Poll interval in ms |\n| `opts.signal` | `AbortSignal` | — | Stop the stream via `AbortController` |\n| `opts.afterId` | `ChangeId` | — | Resume from a previous event |\n| `opts.beforeId` | `ChangeId` | — | Exclusive upper bound |\n| `opts.mode` | `'id'` \\| `'full'` | `'id'` | `'full'` includes decoded blob data |\n| `opts.kinds` | `CdcChangeKind[]` | — | `['INSERT']` \\| `['UPDATE']` \\| `['DELETE']` |\n| `opts.deleteAfterRead` | `boolean` | — | Auto-delete events after yielding |\n| `opts.deleteBatchSize` | `number` | — | Batch delete every N events (requires `deleteAfterRead`) |\n| `opts.deleteBatchWaitMs` | `number` | — | Max wait before flushing a partial batch (requires `deleteBatchSize`) |\n| `opts.batchSize` | `number` | `100` | Max events per poll cycle. Also drives checkpoint cadence — checkpoint is saved after each batch. |\n| `opts.checkpoint` | `CheckpointStrategy` | — | Persistence strategy for crash recovery. See [Resilience](#resilience). |\n\n### `deleteEvents(db, opts)`\n\nDeletes events by `changeId` range, date range, or both (AND). Requires at least one range.\n\n| Param | Type | Description |\n|-------|------|-------------|\n| `opts.changeId.from` | `number` | Inclusive lower bound |\n| `opts.changeId.to` | `number` | Inclusive upper bound |\n| `opts.date.from` | `number` | Unix timestamp, inclusive lower bound |\n| `opts.date.to` | `number` | Unix timestamp, inclusive upper bound |\n| `opts.tableName` | `string` | Scope deletion to a specific table |\n\n## Types\n\n### `CdcEvent<TTable>`\n\n```ts\ninterface CdcEvent<TTable extends AnySQLiteTable = AnySQLiteTable> {\n  changeId:   ChangeId;\n  changeType: CdcChangeKind;\n  changeTime: number | null;\n  changeTxnId: number | null;\n  tableName:  string;\n  rowId:      number | null;\n  before:     InferSelectModel<TTable> | null;\n  after:      InferSelectModel<TTable> | null;\n  updates:    Record<string, unknown> | null;\n}\n```\n\nPass a Drizzle table as the type parameter. `before`/`after` resolve to the table's row shape.\n\n### `CdcChangeKind`\n\n```ts\nexport const CdcChangeKind = {\n  INSERT: 'INSERT',\n  UPDATE: 'UPDATE',\n  DELETE: 'DELETE',\n  COMMIT: 'COMMIT',\n} as const;\n\nexport type CdcChangeKind = (typeof CdcChangeKind)[keyof typeof CdcChangeKind];\n```\n\n`'COMMIT'` is an internal marker — `getEvents` never returns COMMIT rows.\n\n### `ChangeId`\n\nBranded `number` — use as-is from event fields, or cast yours with `id as ChangeId`.\n\n### `CdcMode`\n\n```ts\ntype CdcMode = 'id' | 'before' | 'after' | 'full';\n```\n\nControls what data Turso captures at the PRAGMA level.\n\n## Development\n\n```sh\nnpm test\nnpm run build\nnpm run format\n```\n","readmeFilename":"README.md"}