Skip to content

Latest commit

 

History

History
647 lines (479 loc) · 31 KB

File metadata and controls

647 lines (479 loc) · 31 KB

KVS: @coderbuzz/kvs

Multi-backend key-value store for Bun. Synchronous SQLite, asynchronous SQLite, and PostgreSQL. Atomic transactions, TTL expiry, persistent queue, real-time watch. AI agents: see AI_KNOWLEDGE.md for expert context.

npm version npm downloads MIT License GitHub Stars CI Codecov

KVS is an embeddable key-value store backed by SQLite (sync or async) or PostgreSQL (async). Use it directly in your code. No HTTP server required. Pair with @coderbuzz/kvs-server for HTTP/WS, or @coderbuzz/kvs-client for the client SDK.


Why KVS?

Need KVS Redis Upstash
Infrastructure SQLite file or PostgreSQL Server required Managed
Embeddable Yes, just new KVStore() No (separate process) No
Backends SQLite (sync), SQLite + PostgreSQL (async) - -
Package code 112 KB index.js (24 KB gzip), no dependencies ioredis 6.0.0: 1.29 MB unpacked, 6 dependencies N/A
Transactions Version-based checks + atomic commit MULTI/EXEC/WATCH Conditional checks
Queue Built-in with retries Redis lists + pub/sub Add-on
Watch Push-based (via server) Keyspace notifications Polling

Benchmarks

Measured with bench/throughput.bench.ts in this repository (median of 5 runs; kvs 0.6, Bun 1.4.2, Intel Xeon 2.1 GHz with 4 vCPU, :memory: SQLite, PostgreSQL 16 in Docker on the same host). Numbers from another machine are not comparable; rerun the script there.

Backend set('k','v') get() hit get() miss delete() increment() getMany(10 keys) atomic().sum()
KVStore (bun:sqlite) 388,424 ops/s 618,121 ops/s 1,078,434 ops/s 451,968 ops/s 125,737 ops/s 42,916 ops/s 103,971 ops/s
AsyncKVStore (SQLite) 56,011 ops/s 69,605 ops/s 68,773 ops/s 74,343 ops/s 21,620 ops/s 16,220 ops/s 20,846 ops/s
AsyncKVStore (PostgreSQL) 1,485 ops/s 5,555 ops/s 6,167 ops/s 1,415 ops/s 1,125 ops/s 3,682 ops/s 609 ops/s

Every PostgreSQL call is a network round trip (and each write a commit), which is what the lower numbers pay for multi-process access. On AsyncKVStore, one getMany() of 10 keys takes about the time of 2.5 separate get() calls on SQLite; on KVStore it costs about the same as 10 get() calls and buys one consistent snapshot. kvs makes no speed claim against other stores: the suite at github.com/coderbuzz/benchmarks measures kvs alone, and its current results are for kvs 0.2.11.


Features

  • Hierarchical keys: ["users", "alice"], prefix/range queries (a range inside a prefix too), deterministic sort
  • Typed values: any JSON value, plus Date, bigint, Uint8Array, Map, Set, undefined, NaN, Infinity; anything else is refused, never silently changed
  • Atomic transactions: version checks + set/delete/sum/enqueue in one commit, with store-wide versionstamps (a key never gets an old version back); sum adds to counters inside the commit
  • Batch reads: getMany() reads up to 1000 keys in one statement
  • TTL expiry: millisecond precision, background cleanup every 60 s
  • Built-in queue: leases with ack tokens, retries with backoff, dead letters, awaited listeners with concurrency, per-topic stats
  • Versioned schema: migrations run on open; tablePrefix (and schema on PostgreSQL) keep kvs tables apart from yours
  • Real-time watch: in-process key-change callbacks; over the network via @coderbuzz/kvs-server
  • getAsync: cache-with-compute with singleflight deduplication
  • Multi-backend: SQLite (sync), SQLite + PostgreSQL (async) via unified AsyncKVStore
  • Zero dependencies: no external libs beyond bun:sqlite / bun:sql
  • Bun only (≥ 1.2.21): on Node or Deno, run @coderbuzz/kvs-server under Bun and use @coderbuzz/kvs-client

Installation

npm install @coderbuzz/kvs

Runtime: Bun 1.2.21 or newer, only (engines.bun). KVStore uses bun:sqlite; AsyncKVStore uses Bun.SQL (built-in, no extra deps) with SQLite or PostgreSQL, and Bun 1.2.20 and older have no SQLite in Bun.SQL. The package does not load on Node or Deno; from there, talk to @coderbuzz/kvs-server with @coderbuzz/kvs-client.

Engine minimums. The queue picker uses WITH ... AS MATERIALIZED, which needs PostgreSQL 12+ and SQLite 3.35+. Bun 1.4 bundles SQLite 3.53, so the SQLite backends are always fine; only an older PostgreSQL server is a problem.

Upgrading from 0.5

Nothing to migrate (the schema stays at version 3). New: atomic().sum(), getMany(), and list() with prefix plus start/end, which 0.5 refused with a TypeError and now narrows the range inside the prefix. On PostgreSQL, increment() on a non-numeric value throws the same TypeError as on SQLite (0.5: PostgreSQL's own invalid input syntax for type numeric error).

Upgrading from 0.4

Opening a 0.4 database migrates it to schema version 3, once:

  • Versions are store-wide versionstamps. Every write gets a number no write had before, so a key deleted and created again (or expired) never matches an old version check. Versions no longer go 1, 2, 3 per key; they start above the highest 0.4 version.
  • list({ prefix }) returns the children of the prefix only: not the key ["users"] itself for prefix: ["users"], and not ["users\0x"]. prefix: [] lists every key. prefix together with start or end throws instead of ignoring them (0.6 supports the combination).
  • Values keep their type. A Date comes back as a Date (0.4: an ISO string), a Map/Set as itself (0.4: {}), NaN as NaN (0.4: null), a bigint is stored (0.4: threw). A function, symbol, class instance or circular value throws a TypeError naming where it is. Plain JSON values are stored exactly as in 0.4.
  • Key encoding: -0 is the key 0, NaN is refused, negative bigints sort by value, a bigint over 255 bytes is refused. Keys stored with -0 or a negative bigint are rewritten during the migration; where the rewritten key already exists, the existing row is kept.
  • Keys over 2 KiB (encoded) are refused; maxKeySize changes the limit.
  • atomic().commit() returns one version for the whole commit (0.4: the last set's, or 0), and a builder can be committed once.
  • A 0.4 process refuses a database a 0.5 process has migrated.

Upgrading from 0.3

Opening a 0.3 database migrates it in place to schema version 2 (recorded in the meta table). Nothing to run by hand, but note:

  • PostgreSQL rewrites kv and queue once (version and the queue id become BIGINT), holding an exclusive lock on each while it runs. Plan it for a quiet moment on a large table.
  • acknowledge(id) is now acknowledge(id, token): pass msg.token from dequeue().
  • Listeners are awaited and acknowledge for you: a handler that resolves acks its message, one that throws nacks it. Remove your own acknowledge() call from listeners, or pass autoAck: false.
  • Acknowledged messages are deleted (0.3 kept them as done forever). Keep them for a while with queue.doneRetention.
  • A message a 0.3 worker was holding counts as an expired lease and is delivered again.
  • A database that a newer kvs has migrated is refused with an error rather than misread.

Quick Start

import { KVStore, AsyncKVStore } from "@coderbuzz/kvs";

// Sync (SQLite via bun:sqlite)
const store = new KVStore("kv.db");
store.set(["users", "alice"], { name: "Alice" });
console.log(store.get(["users", "alice"])?.value);

// Async (SQLite via bun:sql)
const asyncStore = new AsyncKVStore("sqlite://kv.db");
await asyncStore.set(["users", "alice"], { name: "Alice" });
console.log(await asyncStore.get(["users", "alice"]));

// Async (PostgreSQL)
const pgStore = new AsyncKVStore("postgres://user:pass@localhost:5432/kvdb");
await pgStore.set(["key"], "value");
await pgStore.delete(["key"]);

KVStore API

new KVStore(path?: string, options?: KVStoreOptions)

Creates or opens a SQLite database. Default path: "kv.db".

Opens with WAL mode, 64 MB cache, 256 MB mmap, busy_timeout = 5000, then creates or migrates the tables.

Option Default Meaning
tablePrefix "" Prefix for the kv, queue and meta tables, e.g. "kvs_"
durability "normal" "full" syncs every commit to disk, see Durability
queue.visibilityTimeout 30_000 Lease length of a dequeued message, ms
queue.backoff 1 s, 2 s, 4 s … max 5 min Retry delay after a failure: number[] (last entry repeats) or (attempt) => ms
queue.doneRetention 0 Keep acknowledged messages this long, ms (0 deletes on ack)
queue.deadRetention Infinity Keep dead messages this long, ms
maxKeySize 2048 Longest encoded key, bytes (Infinity: no limit)

Sharing a database with your application? kvs refuses to open when a table named kv, queue or meta exists without the kvs columns, and never alters it. The check only looks at column names, so give kvs its own names with tablePrefix:

const store = new KVStore("app.db", { tablePrefix: "kvs_" }); // kvs_kv, kvs_queue, kvs_meta

get(key: KvKey): KvEntry | null

const entry = store.get(["users", "alice"]);
// { key: ["users", "alice"], value: { name: "Alice" }, version: 1843 }
// null if missing or expired

getMany(keys: KvKey[]): (KvEntry | null)[]

const [user, settings, missing] = store.getMany([["users", "alice"], ["settings", "alice"], ["nope"]]);
// one slot per key, in order; null when missing or expired

Up to 1000 keys (RangeError beyond), all validated before anything is read. They are read together (one read transaction on KVStore, one statement on the async backends), so the entries come from one snapshot. A key given twice gets two slots.

set(key: KvKey, value: unknown, options?: { ttl?: number }): KvCommitResult

const result = store.set(["users", "alice"], { name: "Alice" });
// { ok: true, version: 1843 }

store.set(["cache", "key"], value, { ttl: 60_000 }); // expires in 60 s
store.set(["event", 1], { at: new Date(), amount: 1999n, raw: new Uint8Array([1, 2]) }); // types kept

Every write gets a new versionstamp: a number unique across the whole store, increasing within a process. Use it with atomic().check(). ttl must be a finite number ≥ 0 (RangeError otherwise).

Values: JSON values plus Date, bigint, Uint8Array (and Buffer, returned as Uint8Array), Map, Set, undefined, NaN and ±Infinity, nested anywhere. -0 is stored as 0. Anything else (functions, symbols, class instances, circular references) throws a TypeError naming the path, e.g. value.items[3].price is a Money. Shared references are stored as copies.

delete(key: KvKey): void

store.delete(["users", "alice"]);

increment(key: KvKey, delta?: number): number

Atomically increment a numeric value. Creates the key with delta if it doesn't exist or has expired. Default delta: 1.

  • A live key keeps its TTL, so set(key, 0, { ttl }) + increment(key) makes a fixed-window rate limiter.
  • A key that holds anything but a number throws a TypeError and is left unchanged, on every backend.
  • delta must be a finite number (RangeError otherwise).
  • To increment inside a transaction, with version checks or several counters at once, use atomic().sum().
// Basic increment
store.increment(["counters", "pageviews"]); // 1 (first call)
store.increment(["counters", "pageviews"]); // 2

// Custom delta (decrement with negative)
store.increment(["users", "alice", "balance"], 500);  // 500
store.increment(["users", "alice", "balance"], -100); // 400

// Rate limiting pattern
const attempts = store.increment(["ratelimit", ip], 1);
if (attempts > 10) throw new Error("Rate limit exceeded");

list(selector: KvListSelector, options?: KvListOptions): KvListResult

// Children of a prefix (not ["users"] itself)
store.list({ prefix: ["users"] });
// Every key
store.list({ prefix: [] });
// Range query
store.list({ start: ["events", 1000], end: ["events", 2000] });
// A range inside a prefix: start inclusive, end exclusive
store.list({ prefix: ["orders"], start: ["orders", "2026-09"], end: ["orders", "2026-10"] });
// Paginated
store.list({ prefix: ["logs"] }, { limit: 20, cursor: cursor });
// Reverse
store.list({ prefix: ["logs"] }, { limit: 5, reverse: true });

With prefix, start and end must be keys inside the prefix (its children) and narrow the range within it, as in Deno KV; a bound outside the prefix, or the prefix key itself, throws a TypeError, so the range never leaves the prefix. Defaults: limit: 100, max 1000 (larger values are capped), ascending. limit must be an integer ≥ 1. cursor is opaque base64 and only valid for the selector it came from: a cursor outside the selector's range throws a RangeError, so a cursor cannot page past a prefix.

KvListResult: { entries: KvEntry[], cursor: string | null }

getAsync<T>(key: KvKey, fn: () => T | Promise<T>, ttl?: number): Promise<T>

Cache-with-compute pattern with singleflight deduplication:

// 100 concurrent callers: fn() runs once, result cached for 30 s
const ad = await store.getAsync(["ads", "venue", 42], () => fetchNextAd(42), 30_000);

Algorithm:

  1. Check SQLite: return immediately on cache hit
  2. Singleflight dedup within process
  3. Call fn() exactly once
  4. Store result in SQLite with TTL
  5. Return to all concurrent callers

atomic(): AtomicOperation

Fluent builder for version-checked transactions:

const seen = store.get(["counter"])!;       // { value: 3, version: 1843, ... }
const result = store
  .atomic()
  .check({ key: ["counter"], version: seen.version }) // fail if written since
  .check({ key: ["new-key"], version: null }) // fail if exists
  .set(["counter"], seen.value as number + 1)
  .set(["meta"], { updatedAt: Date.now() })
  .sum(["stats", "updates"], 1)
  .delete(["old-key"])
  .enqueue({ task: "notify" }, { topic: "jobs" })
  .commit();

if (result.ok) {
  console.log("Version:", result.version); // both sets carry this one version
} else {
  console.log("Check failed, retry");
}
Method Signature Description
check (...checks: KvCheck[]): this Assert key versions
set (key, value, options?): this options: { ttl?: number }
delete (key): this
sum (key, delta: number): this Add delta to the number at key
enqueue (payload, options?): this options: QueueOptions
commit (): KvCommitResult | KvCommitError Execute transaction

check(version: null) = "key must not exist". check(version: N) = "key must be at version N", i.e. unchanged since it was read at N. Because versions are store-wide and never reused, a key deleted and written again (or expired) no longer matches.

A commit gets one versionstamp, returned even when it has no set. A builder can be committed once; a second commit() throws.

sum(key, delta) is increment() inside the commit: a missing or expired key becomes delta with no TTL, a live key keeps its TTL, and it sees the commit's earlier mutations (set(k, 10).sum(k, 5) leaves 15). A non-numeric value makes commit() throw a TypeError and write nothing; delta must be a finite number (RangeError). It is atomic against concurrent writers on every backend (PostgreSQL adds under the row lock), so counters need no read, no version check and no retry:

// A stock movement: the ledger entry and both counters commit together or not at all
store.atomic()
  .set(["ledger", movementId], { sku, qty: -2 })
  .sum(["stock", sku], -2)
  .sum(["stats", "movements"], 1)
  .commit();

Queue

A message is leased to one consumer at a time. dequeue() hands out the message with a token; finish it with acknowledge(id, token) or give it back with nack(id, token). If the consumer dies, the lease runs out after visibilityTimeout (30 s) and the message is delivered again. After maxAttempts deliveries it is dead-lettered instead of retried, where you can inspect, retry or delete it.

enqueue → pending ──dequeue──▶ processing ──acknowledge──▶ deleted (or done, with doneRetention)
             ▲                    │   │
             └── nack / lease ────┘   └── nack or lease expiry on the last attempt ──▶ dead
                 expiry (backoff)                                           retryDead ─┘

enqueue(payload: unknown, options?: QueueOptions): { ok: true, id: number }

store.enqueue(
  { to: "user@example.com", subject: "Welcome" },
  { topic: "emails", delay: 5_000, maxAttempts: 5 },
);

Defaults: topic: "default", delay: 0, maxAttempts: 3 (an integer >= 1).

dequeue(topic?: string, limit?: number, options?: { visibilityTimeout?: number }): QueueMessage[]

Lease up to limit (default 1, max 1000) due messages, oldest first. Each message carries token, lockedUntil, attempts and the lastError of the previous attempt.

for (const msg of store.dequeue("emails", 10)) {
  try {
    await sendEmail(msg.payload);
    store.acknowledge(msg.id, msg.token);
  } catch (error) {
    store.nack(msg.id, msg.token, { error: String(error) }); // retried after the backoff
  }
}

acknowledge(id: number, token: string): boolean

Finish a message: it is deleted (kept as done under queue.doneRetention). Returns false if the lease is no longer yours: it expired and the message went to someone else.

nack(id: number, token: string, options?: { error?: string, delay?: number }): boolean

Give a message back after a failure. It is retried after the backoff (or delay ms), or dead-lettered once it has used maxAttempts. error is stored as lastError.

extendLease(id: number, token: string, visibilityTimeout?: number): boolean

Keep a long job's lease: the lease now ends visibilityTimeout ms from now (default: the store's). 0 hands the message back at once.

listDead(topic?, { limit?, after? }?), retryDead(topic?, id?), deleteDead(topic?, id?)

for (const msg of store.listDead("emails")) console.log(msg.id, msg.lastError, msg.failedAt);
store.retryDead("emails", 42); // one message, attempts start again at 0
store.retryDead("emails");     // every dead message of the topic
store.deleteDead("emails");

queueStats(topic?: string): QueueStats[]

Counts per topic: pending (due now), delayed, processing, dead, done, and oldestPendingAt (lag = Date.now() - oldestPendingAt).

cleanQueue(): number

Queue maintenance: dead-letters messages whose last lease expired, and deletes done/dead messages past their retention. Runs every 60 s on its own; call it from a short-lived script, whose timers never fire.

watch(keys: KvKey[], callback: WatchCallback): { cancel: () => void }

Subscribe to key changes. Fires immediately with current values. The optional second callback argument carries a process-local ordered sequence:

const { cancel } = store.watch(
  [["config", "theme"], ["config", "lang"]],
  (entries, event) => {
    // entries[0] = KvEntry | null for ["config", "theme"]
    console.log(event?.sequence, event?.initial, event?.changedKeys);
  },
);
cancel(); // stop watching

One committed mutation batch invokes each matching watcher once. Committed entries are reused directly, and unchanged keys in multi-key watches are read at most once per batch regardless of subscriber count. Atomic commits changing multiple watched keys therefore produce one coherent callback.

getWatchDiagnostics(): KvWatchDiagnostics

Returns bounded counters for active watchers, committed batches, callbacks, shared reads, callback errors, and dispatcher errors. It does not expose raw keys or values.

addQueueListener(topic, handler, options?): { cancel: () => Promise<void> }

Run handler for each message of a topic. The handler is awaited; when it resolves the message is acknowledged, when it throws it is nacked (retried with backoff, dead-lettered after maxAttempts). While it runs, the lease is renewed.

const listener = store.addQueueListener("emails", async (msg) => {
  await sendEmail(msg.payload); // throw to retry
}, { concurrency: 4 });

await listener.cancel(); // stops taking messages, waits for running handlers
Option Default Meaning
concurrency 1 Handlers of this listener running at once
visibilityTimeout store's Lease length
autoAck true false: the handler calls acknowledge/nack itself

Handlers never run inside enqueue(). Several listeners of a topic (in this or other processes) share its messages. New messages wake the listener at once, including those from atomic().enqueue(); delayed messages, retries and other processes' messages are picked up by a 1 s poll. getQueueDiagnostics() returns counters: delivered, acked, nacked, handler errors, lost leases.

cleanExpired(): number

Manually trigger cleanup of expired entries. Auto-runs every 60s. Returns number of deleted rows and emits null tombstones to active watchers.

store.set(["cache", "a"], "x", { ttl: 1000 });
store.set(["cache", "b"], "y", { ttl: 1000 });
// After 2s, entries are expired: cleanExpired() removes them immediately
store.cleanExpired(); // returns 2

reset(): void

Delete ALL data from kv + queue tables. Active watchers receive a reset tombstone and remain registered, so later writes continue to be delivered.

store.set(["users", "alice"], { name: "Alice" });
store.enqueue("test");
store.reset();
store.get(["users", "alice"]); // null

Durability

SQLite runs in WAL mode with synchronous = NORMAL by default: a power cut or OS crash can lose the last transactions, but never corrupts the file. A process crash loses nothing. Pass durability: "full" to sync every commit, at the cost of write throughput. On PostgreSQL durability is the server's setting (synchronous_commit).

close(): void

Close database, stop cleanup/dispatch timers, cancel watchers/listeners. No operations work after close. The timers are unref'd, so a script that never calls close() still exits.

// Graceful shutdown
process.on("SIGINT", () => {
  store.close();
  process.exit(0);
});

AsyncKVStore API

new AsyncKVStore(connection: string | { adapter: SqlAdapter, queue? }, settings?)

Creates an async KV store backed by SQLite or PostgreSQL. Adapter auto-detected from connection string:

// SQLite file
new AsyncKVStore("sqlite://kv.db");
// SQLite in-memory
new AsyncKVStore(":memory:");
// SQLite via file:// URL
new AsyncKVStore("file://kv.db");
// PostgreSQL, tables kvs_kv/kvs_queue/kvs_meta in schema "infra"
new AsyncKVStore("postgres://user:pass@localhost:5432/app", { tablePrefix: "kvs_", schema: "infra" });
// Pre-built adapter
new AsyncKVStore({ adapter: new PostgresAdapter("postgres://...", { tablePrefix: "kvs_" }) });

settings takes tablePrefix, schema (PostgreSQL only, created if missing), durability (SQLite only) and queue (as for KVStore). With a pre-built adapter, give the table options to the adapter.

Connection string rules:

  • postgres://... or postgresql://... → PostgresAdapter
  • anything else → SQLiteAsyncAdapter, which passes the string to Bun's SQL. Use sqlite://..., file://..., or :memory:. A plain filename such as "kv.db" is parsed by Bun as a PostgreSQL connection and fails to connect.

Methods

All methods return Promise<T> (same signatures as KVStore but async):

await store.get(key);             // Promise<KvEntry | null>
await store.getMany(keys);        // Promise<(KvEntry | null)[]>
await store.set(key, val, opts?); // Promise<KvCommitResult>
await store.delete(key);          // Promise<void>
await store.list(sel, opts?);     // Promise<KvListResult>
await store.increment(key, n?);   // Promise<number>
await store.enqueue(payload, opts?); // Promise<{ ok, id }>
await store.dequeue(topic?, n?, opts?); // Promise<QueueMessage[]>
await store.acknowledge(id, token);    // Promise<boolean>
await store.nack(id, token, opts?);    // Promise<boolean>
await store.extendLease(id, token, ms?); // Promise<boolean>
await store.listDead(topic?, opts?);   // Promise<QueueDeadMessage[]>
await store.retryDead(topic?, id?);    // Promise<number>
await store.deleteDead(topic?, id?);   // Promise<number>
await store.queueStats(topic?);        // Promise<QueueStats[]>
await store.cleanQueue();              // Promise<number>
await store.cleanExpired();       // Promise<number>
await store.reset();              // Promise<void>
await store.close();              // Promise<void>
await store.getAsync(key, fn, ttl?); // Promise<T> (already async)

watch(), addQueueListener(), getWatchDiagnostics() and getQueueDiagnostics() remain sync (in-process callbacks). On AsyncKVStore the initial watch snapshot is delivered asynchronously.

On SQLite, AsyncKVStore runs one statement at a time (Bun's SQLite client has a single connection), so an atomic() never picks up or rolls back another call's write. On PostgreSQL, atomic() locks the keys it checks and writes, so two concurrent commits against the same version cannot both succeed.

atomic(): AsyncAtomicOperation

Same fluent builder as AtomicOperation but commit() is async:

const result = await store
  .atomic()
  .check({ key: ["counter"], version: seen.version })
  .set(["counter"], 4)
  .sum(["stats", "writes"], 1)
  .enqueue({ task: "notify" }, { topic: "jobs" })
  .commit();

Adapters

The AsyncKVStore uses an internal SqlAdapter interface. You can build custom adapters or use the built-in ones:

Adapter Class Backend
SQLite SQLiteAsyncAdapter SQLite via bun:sql
PostgreSQL PostgresAdapter PostgreSQL via bun:sql
import { PostgresAdapter } from "@coderbuzz/kvs";

const adapter = new PostgresAdapter("postgres://user:pass@localhost:5432/kvdb", { tablePrefix: "kvs_", schema: "infra" });
const store = new AsyncKVStore({ adapter });

SQLiteAsyncAdapter(connection, { tablePrefix?, durability? }) and PostgresAdapter(connection, { tablePrefix?, schema? }). A custom adapter implements SqlAdapter, whose queue methods changed in 0.4 (see AI_KNOWLEDGE.md).

SQL Dialect Differences

Feature SQLite PostgreSQL
Key column BLOB BYTEA
Queue ID INTEGER PRIMARY KEY AUTOINCREMENT BIGINT from a sequence (0.3 SERIAL migrated)
Entry version INTEGER (64-bit) BIGINT (0.3 INTEGER migrated)
Versionstamp one-row counter table; each store reserves blocks of 1024 (unique store-wide, increasing per process) a sequence, nextval() per write (unique and increasing store-wide)
Timestamp INTEGER BIGINT
increment(), atomic().sum() read-modify-write in one transaction (JS arithmetic) one upsert with NUMERIC arithmetic (exact decimals) under the row lock
getMany() one IN (...) statement (KVStore: reads in one transaction) one IN (...) statement
atomic() concurrency transactions run one at a time advisory lock per key, FOR UPDATE on checked rows
Concurrent dequeue MATERIALIZED CTE picker (writers serialized) MATERIALIZED CTE picker with FOR UPDATE SKIP LOCKED
Concurrent migration BEGIN IMMEDIATE transaction-scoped advisory lock
Partial indexes WHERE expires_at IS NOT NULL same

Types

import type {
  KvKey,           // KvKeyPart[]
  KvKeyPart,       // string | number | bigint | boolean | Uint8Array (NaN refused, -0 = 0)
  KvEntry,         // { key, value, version }
  KvWatchEvent,    // { sequence, initial, changedKeys, coalesced?, reset? }
  KvWatchDiagnostics,
  KvCommitResult,  // { ok: true, version }
  KvCommitError,   // { ok: false }
  KvCheck,         // { key, version }
  KvMutation,      // { type: "set"|"delete"|"sum", key, value?, ttl? }
  KvListSelector,  // { prefix?, start?, end? }
  KvListOptions,   // { limit?, cursor?, reverse? }
  KvListResult,    // { entries, cursor }
  QueueMessage,    // { id, topic, payload, enqueuedAt, deliverAt, attempts, maxAttempts, token, lockedUntil, lastError }
  QueueDeadMessage,// { id, topic, payload, …, lastError, failedAt }
  QueueOptions,    // { topic?, delay?, maxAttempts? }
  QueueStats,      // { topic, pending, delayed, processing, dead, done, oldestPendingAt }
  KvQueueConfig,   // { visibilityTimeout?, backoff?, doneRetention?, deadRetention? }
} from "@coderbuzz/kvs";

Key Encoding

Keys are encoded to bytes with deterministic sort order:

Uint8Array < string < number < bigint < false < true
["a"] < ["b"]
["users", 1] < ["users", 2]
["ledger", -1000n] < ["ledger", -1n] < ["ledger", 0n]
["items", true] > ["items", false]

-0 and 0 are the same key; NaN is not a valid key part; a bigint part holds at most 255 bytes. An encoded key is at most maxKeySize bytes (2048 by default, as in Deno KV).

Encoding Utilities

Low-level functions for key serialization:

import { encodeKey, decodeKey, encodeKeyPrefix, prefixSuccessor } from "@coderbuzz/kvs";

// Round-trip: KvKey → bytes → KvKey
const encoded = encodeKey(["users", "alice"]);
const decoded = decodeKey(encoded); // ["users", "alice"]

// Low-level byte range for custom queries. It includes the key ["events"] itself
// and string siblings such as ["events\0x"]; list({ prefix }) does not.
const prefix = encodeKeyPrefix(["events"]);
const upper = prefixSuccessor(prefix);
// Range: key >= prefix AND key < upper

Singleflight

Exported standalone for deduplicating concurrent async work:

import { Singleflight } from "@coderbuzz/kvs";

const sf = new Singleflight<User>();
const user = await sf.do("user:42", () => fetchUser(42));
sf.clear();

Server & Client

  • Server: @coderbuzz/kvs-server wraps KVStore (sync) or AsyncKVStore (async) into HTTP REST + WebSocket server
    • createServer(store, opts) for sync KVStore
    • createAsyncServer(store, opts) for async AsyncKVStore
  • Client: @coderbuzz/kvs-client is the TypeScript SDK for the server

License

MIT © 2026 Indra Gunawan