Package: @coderbuzz/kvs
Purpose: Multi-backend key-value store. Sync KVStore (bun:sqlite) and async AsyncKVStore (bun:sql, SQLite + PostgreSQL).
Distribution: ESM only (dist/index.js + one .d.ts per source file). dist/index.js is 115,093 bytes (24,433 gzip -9), unminified ES2022, no dependencies (kvs 0.6).
Runtime: Bun ≥ 1.2.21 only (engines.bun). Measured: the full suite (415 tests, PostgreSQL included) passes on Bun 1.2.21, 1.3.0, 1.3.11 and 1.4.2; Bun 1.2.20 fails 137 tests with Unsupported adapter: sqlite. Only "postgres" is supported for now (Bun.SQL gained SQLite in 1.2.21). dist/index.js imports bun:sqlite and bun at top level, so it does not load on Node (ERR_UNSUPPORTED_ESM_URL_SCHEME) or Deno. From Node/Deno, run kvs-server on Bun and use @coderbuzz/kvs-client.
KVStore("kv.db") : sync, bun:sqlite, embedded SQLite only
AsyncKVStore("sqlite://kv.db") : async, bun:sql, SQLite
AsyncKVStore("postgres://...") : async, bun:sql, PostgreSQL
KVStore("kv.db")
├── get/set/delete : CRUD (sync)
├── getMany : up to 1000 keys, one snapshot (sync, 0.6)
├── increment : atomic counter (sync)
├── list : prefix/range queries, range inside a prefix (sync)
├── atomic() : version-checked transactions: check/set/delete/sum/enqueue (sync)
├── enqueue/dequeue : persistent queue with leases (sync)
├── acknowledge/nack/extendLease(id, token)
├── listDead/retryDead/deleteDead, queueStats, cleanQueue
├── watch() : in-process callbacks (sync)
├── addQueueListener() : awaited handlers, auto-ack/nack, concurrency
├── getAsync() : cache-with-compute (singleflight, async)
├── cleanExpired() / reset()
└── close()
AsyncKVStore("sqlite://kv.db" | "postgres://...")
├── get/set/delete : CRUD (async)
├── getMany : up to 1000 keys in one statement (async, 0.6)
├── increment : atomic counter (async)
├── list : prefix/range queries, range inside a prefix (async)
├── atomic() : version-checked transactions incl. sum (async commit)
├── enqueue/dequeue/ack/nack, dead letters, stats : persistent queue (async)
├── watch() : same as KVStore (in-process; initial snapshot delivered async)
├── addQueueListener() : same as KVStore
├── getAsync() : same as KVStore
├── cleanExpired() / reset() : async
└── close() : async
bench/throughput.bench.ts (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; PostgreSQL only when KVS_BENCH_PG names a server you own, tables kvsbench_*). delete() removes an existing key per call; increment() calls increment() (the kvs-throughput suite of github.com/coderbuzz/benchmarks, whose latest results are for kvs 0.2.11 on Apple Silicon, measured get + set under that name).
| 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, every write a commit. PostgreSQL
atomic().sum()also takes an advisory lock andnextval()in its transaction. getMany(10)vs 10 awaitedget()calls: SQLiteAsyncKVStore18K vs 7K calls/s;KVStoreabout equal (seegetMany).- The hot-path micro-benchmark is
bench/core.bench.ts(encode/decode, get, set, list, queue cycle, batch reads), used for before/after numbers of every optimization. - No claim against other stores is made: no benchmark in the benchmarks repo compares kvs with one.
No schema change (version 3), nothing to migrate. Behaviour changes: list({ prefix, start | end }) works (0.5: TypeError); PostgreSQL increment() on a non-number throws TypeError: kvs: increment() needs a number, but the key holds a non-numeric value (0.5: the raw PostgresError 22P02/22021). SqlAdapter gains an optional getMany and an optional fourth version argument to increment; existing adapters keep working (getMany() then reads key by key; an increment that ignores version makes the summed key carry its own versionstamp instead of the commit's).
import {
KVStore, AtomicOperation, // sync
AsyncKVStore, AsyncAtomicOperation, // async
Singleflight,
type WatchCallback,
type KVStoreOptions, type AsyncKVStoreOptions, type AsyncKVStoreSettings,
openDatabase, StmtCache, type OpenDatabaseOptions,
SQLiteAsyncAdapter, PostgresAdapter, // adapters
type SQLiteAdapterOptions, type PostgresAdapterOptions, type TableOptions, type Durability,
type SqlAdapter, type KvRow, type QueueRow, type QueueStatsRow, type LeaseRow,
type SetResult, type IncrementResult, type EnqueueResult,
SCHEMA_VERSION, // 2
type KvQueueConfig, type KvQueueDiagnostics,
type QueueDeadMessage, type QueueDequeueOptions, type QueueNackOptions,
type QueueListenerOptions, type QueueStats,
encodeKey, decodeKey, encodeKeyPrefix, prefixSuccessor,
type KvKey, type KvKeyPart, type KvEntry,
type KvWatchEvent, type KvWatchDiagnostics,
type KvCommitResult, type KvCommitError,
type KvCheck, type KvMutation,
type KvListSelector, type KvListOptions, type KvListResult,
type QueueMessage, type QueueOptions,
} from "@coderbuzz/kvs";type KvKeyPart = string | number | bigint | boolean | Uint8Array
type KvKey = KvKeyPart[]
interface KvEntry {
key: KvKey
value: unknown
version: number
}
interface KvCommitResult { ok: true; version: number }
interface KvCommitError { ok: false }
interface KvCheck {
key: KvKey
version: number | null // number = "key must be at this version"
// null = "key must not exist"
}
// "sum" (0.6): value is the delta, a finite number; ttl applies to "set" only
interface KvMutation { type: "set" | "delete" | "sum"; key: KvKey; value?: unknown; ttl?: number }
// prefix alone: its children. start inclusive, end exclusive. With prefix (0.6),
// start/end must be children of the prefix and narrow the range inside it.
interface KvListSelector { prefix?: KvKey; start?: KvKey; end?: KvKey }
interface KvListOptions { limit?: number; cursor?: string; reverse?: boolean }
interface KvListResult { entries: KvEntry[]; cursor: string | null }
interface QueueMessage {
id: number; topic: string; payload: unknown
enqueuedAt: number
deliverAt: number // due time; a nack with backoff moves it to the retry time
attempts: number // deliveries so far, this one included
maxAttempts: number
token: string // lease token of this delivery (crypto.randomUUID, shared by one dequeue call)
lockedUntil: number // lease end, ms epoch
lastError: string | null // error text of the last failed attempt (nack or "lease expired")
}
interface QueueDeadMessage {
id: number; topic: string; payload: unknown
enqueuedAt: number; deliverAt: number; attempts: number; maxAttempts: number
lastError: string | null
failedAt: number // when it was dead-lettered
}
interface QueueOptions { topic?: string; delay?: number; maxAttempts?: number } // maxAttempts: integer >= 1, default 3
interface QueueDequeueOptions { visibilityTimeout?: number } // ms > 0
interface QueueNackOptions { error?: string | null; delay?: number } // delay ms >= 0 overrides backoff
interface QueueListenerOptions { concurrency?: number; visibilityTimeout?: number; autoAck?: boolean }
interface KvQueueConfig {
visibilityTimeout?: number // default 30_000
backoff?: number[] | ((attempt: number) => number) // default min(1000 * 2^(attempt-1), 300_000)
doneRetention?: number // default 0 (delete on ack)
deadRetention?: number // default Infinity
}
interface QueueStats {
topic: string
pending: number; delayed: number; processing: number; dead: number; done: number
oldestPendingAt: number | null
}
interface KvQueueDiagnostics {
listeners: number; inFlight: number
delivered: number; acked: number; nacked: number
handlerErrors: number
leaseLost: number // ack/nack/renewal refused: the lease moved to another delivery
dispatchErrors: number // failed dequeue/ack/renewal (closed store, lost connection)
}
interface KvWatchEvent {
sequence: number // monotonic only for the current store process
initial: boolean
changedKeys: KvKey[]
coalesced?: boolean
reset?: boolean
}
interface KvWatchDiagnostics {
activeWatchers: number
committedBatches: number
callbacks: number
sharedReads: number
callbackErrors: number
dispatchErrors: number
}
type WatchCallback = (entries: (KvEntry | null)[], event?: KvWatchEvent) => voidKvKeyPart = string | number | bigint | boolean | Uint8Array
KvKey = KvKeyPart[]
Sort: Uint8Array < string < number < bigint < false < true
Encoding: each key part is a type-tag byte (0x01 bytes, 0x02 string, 0x03 number, 0x04 bigint, 0x05 false, 0x06 true) followed by its payload, then a 0x00 separator. String/byte payloads escape 0x00 as 0x00 0xff; numbers are 8-byte big-endian float64 with sign-flip so byte order matches numeric order; bigints are a sign byte (0x00 negative, 0x01 otherwise), a length byte, then magnitude bytes, all inverted (0xff - b, the length too) for negatives so they sort by value. The concatenated bytes sort lexicographically. An empty key encodes to zero bytes.
Rules (KVS-10, kvs 0.5): -0 is normalized to 0 (0.4 encoded it as eight 0x00 bytes and decoded it as NaN); NaN throws TypeError: kvs: NaN cannot be a key part; a bigint needs at most 255 magnitude bytes (RangeError beyond; 0.4 wrapped the length byte and wrote an undecodable key); any other part type (null, undefined, object, array, symbol) throws TypeError (0.4 wrote it as true/false by truthiness); a key that is not an array throws TypeError. The store refuses an encoded key longer than maxKeySize (default 2048) with RangeError: kvs: key is N bytes encoded, over the 2048-byte limit (maxKeySize), on every call that takes a key (reads, writes, checks, list selectors). encodeKey() itself has no size limit.
list({ prefix: p }) reads [enc(p) ‖ 0x01, enc(p) ‖ 0x07): exactly the keys with at least one more part after p (KVS-16). enc(p) ends with the 0x00 separator; the next byte of a child is its type tag 1..6, while the key p itself ends there and a string sibling "p\0x" continues with 0xff. prefix: [] gives [0x01, 0x07), every key. kvs-server's scoped credentials use the same rule (withinPrefix: equal, or the prefix followed by a byte 1..6).
With prefix plus start/end (0.6, as Deno KV), start replaces the lower bound enc(p) ‖ 0x01 with enc(start) and end replaces the upper bound enc(p) ‖ 0x07 with enc(end). Each must be a child of p by the rule above (enc(bound) starts with enc(p) and the next byte is 1..6), otherwise TypeError: kvs: list() start must be a key inside the prefix (or end). So ["v"], the prefix key ["u"] itself, the byte sibling ["u\0x"], ["uX", 1] and [] are all refused for prefix: ["u"], and the range can never leave the prefix: this is what keeps a kvs-server credential scoped to a prefix inside it. start > end gives an empty page.
Decoder (KVS-23, 0.6). decodeKey scans each string/bytes part once to its unescaped 0x00 separator. A string part without escapes is decoded straight from a view of the row's bytes with one shared TextDecoder, or, when it is ASCII and at most 8 bytes, with String.fromCharCode (faster for short parts: 29M vs 8.7M decodes/s of ["k"]; TextDecoder wins from ~13 bytes). A bytes part is always returned as a fresh copy (slice), never a view of the row buffer. Number parts go through one shared 8-byte scratch buffer instead of a new ArrayBuffer + DataView per part; encodeKey uses the same scratch and one shared TextEncoder. Measured on the same machine (Bun 1.4.2, bench/core.bench.ts, median of 3 runs, 0.5 → 0.6): decodeKey(["users","alice",42]) 782K → 5.30M ops/s, encodeKey 2.87M → 4.65M, KVStore.get hit 551K → 679K (+23%), list 100 entries with number parts 4.5K → 10.2K, with string parts 5.2K → 5.9K (+14%). A first version (indexOf + a result object + TextDecoder for every string) made get hit 10% slower: short parts were slower than 0.5's copy loop.
import { encodeKey, decodeKey, encodeKeyPrefix, prefixSuccessor } from "@coderbuzz/kvs";
["a"] < ["b"]
["users", 1] < ["users", 2]
["items", true] > ["items", false]
// Round-trip: KvKey → bytes → KvKey
const encoded = encodeKey(["users", "alice"]);
const decoded = decodeKey(encoded); // ["users", "alice"]
// Low-level byte range for custom queries: includes the key ["events"] itself and
// string siblings like ["events\0x"], unlike list({ prefix })
const prefix = encodeKeyPrefix(["events"]);
const upper = prefixSuccessor(prefix);
// Resulting range: key >= prefix AND key < upper- Default path:
"kv.db" - Opens/creates SQLite database with WAL mode, 64 MB cache, 256 MB mmap,
busy_timeout = 5000, then runs the schema migration (see "Schema versions") insideBEGIN IMMEDIATE. A failed migration closes the database and throws. - Starts one 60 s maintenance timer: TTL cleanup + queue maintenance (
cleanQueue()). options:
| Option | Type | Default | Notes |
|---|---|---|---|
tablePrefix |
string |
"" |
/^[A-Za-z_][A-Za-z0-9_]*$/, max 40 chars; else TypeError. Applies to tables and indexes (idx_<prefix>kv_expires …) |
durability |
"normal" | "full" |
"normal" |
PRAGMA synchronous; other values TypeError |
queue |
KvQueueConfig |
see Types | invalid values RangeError (backoff of the wrong type: TypeError) |
const store = new KVStore("kv.db"); // sync, bun:sqlite
const shared = new KVStore("app.db", { tablePrefix: "kvs_", durability: "full", queue: { visibilityTimeout: 60_000 } });new AsyncKVStore(connection: string | { adapter: SqlAdapter; queue?: KvQueueConfig }, settings?: AsyncKVStoreSettings)
- Auto-detects adapter from connection string:
"postgres://..."or"postgresql://..."→PostgresAdapter(conn, { tablePrefix, schema })- anything else →
SQLiteAsyncAdapter(conn, { tablePrefix, durability }), which hands the string to Bun'snew SQL(...). Use"sqlite://...","file://...", or":memory:". - A plain filename such as
"kv.db"is NOT SQLite: Bun'sSQLparses it as a PostgreSQL connection, so the first operation fails with a connection error.
settings: { tablePrefix?, schema?, durability?, queue? }.schemawith SQLite, ordurabilitywith PostgreSQL, throwsTypeError. With{ adapter }, passingtablePrefix/schema/durabilityinsettingsthrowsTypeError(configure the adapter);queuemay go in either object.- Migration and the 60 s maintenance timer start lazily on the first operation (
ensureInit()), not in the constructor. Ifmigrate()rejects (database not up yet, or a foreign table), that call rejects and the next operation runsmigrate()again; a failure is not cached.
const asyncStore = new AsyncKVStore("sqlite://kv.db");
const pgStore = new AsyncKVStore("postgres://user:pass@localhost:5432/app", { tablePrefix: "kvs_", schema: "infra" });
const customStore = new AsyncKVStore({ adapter: new PostgresAdapter("postgres://...", { tablePrefix: "kvs_" }) });SCHEMA_VERSION = 3, stored as text in<prefix>metaunder keyschema_version. A database without the row is new or was written by kvs <= 0.3 (whose tables are exactly version 1).- Every open:
CREATE TABLE IF NOT EXISTS meta, read the version, run the missing steps, write the version, then verify every table has the current columns, all in ONE transaction (SQLiteBEGIN IMMEDIATE; PostgreSQLsql.begin+pg_advisory_xact_lockon a hash ofmigrate:<meta table>, andCREATE SCHEMA IF NOT EXISTSfirst whenschemais set). Concurrent openers (processes) wait for each other; any failure rolls everything back. - Step 0 → 1:
CREATE TABLE IF NOT EXISTS kv / queue, then verify the version-1 columns, then the indexes. An existing table of the same name that is not a kvs table (e.g. an applicationqueue) fails here withError: kvs: table "queue" exists but is not a kvs table (missing columns: payload, enqueued_at, attempts, max_attempts). Use the tablePrefix option ...before anything is indexed or altered. - Step 1 → 2: queue columns
locked_until,lease_token,last_error,finished_at;processingrows getlocked_until = 0(expired lease: redelivered, or dead-lettered bycleanQueue()if attempts are used up);donerows getfinished_at = now(then aged out bydoneRetention); indexesidx_<p>queue_lease (status, locked_until) WHERE status='processing'andidx_<p>queue_finished (status, finished_at) WHERE finished_at IS NOT NULL. PostgreSQL also runsALTER TABLE kv ALTER COLUMN version TYPE BIGINT,ALTER TABLE queue ALTER COLUMN id TYPE BIGINTandALTER SEQUENCE <serial seq> AS BIGINT(KVS-17): each rewrites its table under anACCESS EXCLUSIVElock, once. - Step 2 → 3 (kvs 0.5): the versionstamp counter (KVS-05) and the key rewrite (KVS-10).
- SQLite:
CREATE TABLE <p>kv_versionstamp (id INTEGER PRIMARY KEY CHECK (id = 1), version INTEGER NOT NULL), one row initialised toMAX(kv.version). - PostgreSQL:
CREATE SEQUENCE <schema>.<p>kv_versionstamp AS BIGINT,setvaltoMAX(kv.version)(or 1 unused whenkvis empty). - Rekey: candidates are the keys holding a
04 00(negative bigint) or03 00×8(0.4's-0) byte run (instr()on SQLite,position()on PostgreSQL). Each is decoded with the 0.4 decoder (src/legacy-keys.ts) and encoded again; unchanged bytes are skipped, a key that cannot be decoded (0.4 bigint over 255 bytes) is left as it is. When the new key already exists (["x", -0]next to["x", 0]) the old row is deleted and the existing one kept; otherwiseUPDATE kv SET key = new. Versions and values are kept. Keys with 0.4'sNaNbytes decode to the same bytes (a subnormal) and stay.
- SQLite:
- A stored version above
SCHEMA_VERSIONthrowsError: kvs: the database schema is version N, newer than this kvs release supports (3). Upgrade @coderbuzz/kvs.A non-integer value throws too. So a 0.4 process refuses a database a 0.5 process has migrated.
- Every write (
set,increment,getAsyncfill, eachatomic()commit) gets a version no write of this store ever got, so a check{ key, version }passes only if the key was not written since it was read at that version, including across delete/recreate and TTL expiry (0.4 restarted a recreated key at 1: ABA). - SQLite (both stores): the store reserves
VERSION_BLOCK = 1024versions withUPDATE <p>kv_versionstamp SET version = version + 1024 RETURNING versionand hands them out from memory (src/versions.ts). Unique store-wide, also across processes on one file (each reservation is its own write; SQLite has one writer at a time), increasing within one store instance; two processes interleave in blocks, so across processes versions are unique but not time-ordered. A crash or close leaves a gap.KVStorereserves only outside transactions;SQLiteAsyncAdapter.transaction()drops the block when the transaction rolls back, since a reservation made inside it rolled back too. - Measured alternative: taking
version + 1from the counter inside each upsert plus triggers kept one global order but cost 35% ofAsyncKVStore.set(44K → 29K ops/s raw) and 24% of the raw sync upsert. The block form also droppedRETURNING:KVStore.setsmall 202K → 309K ops/s,AsyncKVStore.set36K → 48K. - PostgreSQL:
nextval('<p>kv_versionstamp')inline in the upsert, so versions are unique and increasing store-wide (sequence order, not commit order). atomic().commit(): one version for the whole commit, taken before the transaction (KVStore) or withtx.nextVersion()after the checks (AsyncKVStore), written to everysetof it and returned even for a commit without aset(0.4: last set's version or0, KVS-20). A builder commits once: a secondcommit()throwsError: kvs: this atomic operation was already committed, also after{ ok: false }.- Numbers stay below 2^53 for any realistic lifetime (a block per store open, 1024 per block).
- SQLite: WAL +
synchronous = NORMAL(default): durable across a process crash; a power loss or OS crash can lose the last committed transactions (no corruption).durability: "full"setssynchronous = FULL(fsync per commit).PRAGMA synchronousreads 1 / 2. - PostgreSQL: the server's
synchronous_commitdecides; kvs sets nothing.
const entry = store.get(["users", "alice"]);
// { key: ["users", "alice"], value: { name: "Alice" }, version: 1843 }
// null if missing or expiredconst [a, missing, b] = store.getMany([["users", "a"], ["nope"], ["users", "b"]]);
// [{ key, value, version }, null, { key, value, version }]- One slot per key, in input order;
nullwhen absent or expired. A duplicate key gets its own slot and its own decoded copy of the value. keysmust be an array (TypeError: kvs: getMany() takes an array of keys), at most1000keys (RangeError: kvs: getMany() reads at most 1000 keys, got N). Every key is encoded and checked (part types,maxKeySize) before anything is read;[]returns[].KVStore: the reads run in one read transaction (db.transaction(), deferred) with onenow, so they see one snapshot; a single key is a plainget. Measured: about the cost of N separateget()calls (10 keys: 36K vs 42K calls/s of 10 gets), so use it for the snapshot, not for speed. A singleIN (...)statement was measured and was slower in-process (36–40K vs 43–47K), because of the map from rows back to slots.AsyncKVStore: keys are deduplicated and read with the adapter's optionalgetMany(keys, now), oneSELECT ... WHERE key IN (...) AND livestatement on both built-in adapters (one snapshot), then matched back to slots by key bytes. Measured on SQLite: 10 keys 18K calls/s vs 7K for 10 awaitedget()s. An adapter withoutgetManyis read key by key withget()(no snapshot guarantee).
store.set(["users", "alice"], { name: "Alice" });
// { ok: true, version: 1843 }
store.set(["cache", "key"], value, { ttl: 60_000 }); // expires in 60s
store.set(["event"], { at: new Date(), cents: 1999n, blob: new Uint8Array([1]), tags: new Set(["a"]) });Every write gets a new versionstamp (see "Versionstamps").
Values (KVS-09). Stored bytes by first byte: 0x00 null, 0x01 true, 0x02 false, 0x03 typed JSON, anything else plain JSON text (0.4 format, unchanged for plain JSON values).
- Supported: JSON values plus
Date(incl. invalid),bigint,Uint8Array(aBuffercomes back asUint8Array),Map,Set,undefined(top level, in objects — the key is kept — and in arrays/holes),NaN,±Infinity, nested anywhere.-0is stored as0. Shared references are stored as copies. - Typed JSON = JSON where special values are marker objects
{"~": T, "v": …}:DDate (epoch ms or null),Bbigint (decimal string),UUint8Array (base64),MMap ([[k, v], …]),SSet (array),uundefined,Nnumber ("NaN","Infinity","-Infinity"),Oa user object that has its own"~"key (as entries, so it is never read as a marker). Decoding walks the parsed tree (not aJSON.parsereviver, which would dropundefinedproperties). - Refused with
TypeError: kvs: value.a[2].b is a Money, which kvs cannot store (supported: …)and nothing written: functions, symbols, class instances (any prototype other thanObject.prototype/null, except the types above; e.g.RegExp,Error,Int16Array), circular references. Checked before any write, also for queue payloads and everyatomic()mutation. - Cost: every write scans the value (a fast scan without path; the naming scan runs only on a refusal). Measured: encoding a small value 10.2M → 8.4M ops/s, an 11.5 KB object 51 µs → 84 µs;
KVStore.set11.5 KB 14.9K → 9.9K ops/s (−33%). A scan withObject.keys/Object.valuesor inlined primitive checks was not faster. TTL is in milliseconds and must be a finite number ≥ 0;NaN,Infinityand negative values throwRangeError(onAsyncKVStore, the promise rejects).ttl: 0expires immediately.
store.delete(["users", "alice"]);Atomically increment a numeric value and return the new value. Default delta: 1.
| Stored state | Result |
|---|---|
| key missing | created with delta, no TTL, version 1 |
| key expired (row still present until cleanup) | treated as missing: value delta, TTL removed, version continues from the old row |
| live number | value + delta; the key keeps its TTL |
live non-number (object, string, null, boolean, typed value such as Date/bigint) |
throws TypeError: kvs: increment() needs a number, but the key holds <kind> (KVStore, SQLite AsyncKVStore) or TypeError: kvs: increment() needs a number, but the key holds a non-numeric value (PostgreSQL, 0.6: the adapter maps the NUMERIC cast errors 22P02 and 22021; 0.5 rejected with PostgreSQL's own invalid input syntax for type numeric). The value is left unchanged. kvs-server answers 400 on every backend |
delta not a finite number |
RangeError, nothing written |
| result not finite (overflow) | RangeError, nothing written (SQLite backends) |
Implementation per backend:
KVStore(bun:sqlite):BEGIN IMMEDIATE,SELECT value, version, expires_at, JS arithmetic, upsert through the same statement asatomic().set(),COMMIT. The value is re-encoded like anyset(), so it stays a BLOB.SQLiteAsyncAdapter: the same read-modify-write insidesql.begin(), behind the adapter's serial queue.PostgresAdapter: oneINSERT ... ON CONFLICT (key) DO UPDATEstatement withNUMERICarithmetic on the JSON text. The upsert takes the row lock, so concurrent increments of a new key never lose an update. Decimals are exact in storage (0.1 + 0.2stores0.3); the returned number isFLOAT8.
History: kvs ≤ 0.3.1 stored the result with CAST(... AS TEXT) on SQLite, which made every later get()/list() on that key throw TextDecoder.decode expects an ArrayBuffer, overwrote non-numeric values with delta, and resurrected expired keys. On PostgreSQL, concurrent increments of a new key lost updates. Rows written as TEXT are decoded again since 0.3.2 and rewritten as BLOB on the next increment().
Cost (Bun 1.4.2, Xeon 2.1 GHz, bench/core.bench.ts): KVStore.increment ~63K ops/s (was ~112K on the broken path), SQLite AsyncKVStore.increment ~23K ops/s (was ~60K).
// Basic increment (default delta: 1)
store.increment(["counter", "visits"]); // 1 (first call)
store.increment(["counter", "visits"]); // 2
// Custom delta: positive or negative
store.increment(["counter", "visits"], 5); // 7
store.increment(["counter", "visits"], -1); // 6
// Rate limiting pattern
const attempts = store.increment(["ratelimit", "192.168.1.1"], 1);
if (attempts > 10) throw new Error("Rate limit exceeded");
// Returns the new value after incrementDefaults: limit: 100, max 1000, ascending, reverse: false. cursor is opaque base64. Without prefix: start defaults to the empty key, end to [0xff].
prefix with start/end (see "Key Encoding"): 0.4 silently ignored start/end, 0.5 threw TypeError: kvs: list() takes { prefix } or { start, end }, not both, 0.6 narrows the range inside the prefix (start inclusive, end exclusive, either or both) and throws TypeError: kvs: list() start|end must be a key inside the prefix for a bound that is not a child of the prefix. A cursor must fall inside the narrowed range, so a cursor from the whole prefix that points before start is a RangeError.
Validation:
limitmust be an integer ≥ 1, otherwiseRangeError(0, negatives,2.5,NaN,Infinity). Values above 1000 are capped to 1000, not rejected.cursoris the base64 of the last returned key's encoded bytes. It must fall inside the selector's byte range[start, end), otherwiseRangeError: kvs: list cursor is outside the selector range. Forward pagination then starts atcursor ‖ 0x00(exclusive); reverse pagination ends atcursor(exclusive). kvs ≤ 0.3.1 accepted any cursor and replaced the range bound with it, so a caller authorised for one prefix could read every key in the store (AA==below,/w==above). kvs-server maps thisRangeErrorto HTTP 400.- A cursor is not bound to its selector beyond that range check: reusing a cursor with a different selector whose range still contains it is allowed.
// Children of a prefix: not ["users"] itself, not ["users\0x"] (KVS-16)
store.list({ prefix: ["users"] });
// Every key
store.list({ prefix: [] });
// Range query
store.list({ start: ["events", 1000], end: ["events", 2000] });
// A range inside a prefix (0.6): September's orders, paged
store.list({ prefix: ["orders"], start: ["orders", "2026-09"], end: ["orders", "2026-10"] }, { limit: 50 });
store.list({ prefix: ["orders"], start: ["orders", "2026-09"] }); // from start to the end of the prefix
store.list({ prefix: ["orders"], end: ["orders", "2026-01"] }); // from the prefix start to end
store.list({ prefix: ["orders"], start: ["invoices"] }); // TypeError: start outside the prefix
// Paginated
const page1 = store.list({ prefix: ["logs"] }, { limit: 20 });
// page1 = { entries: [...], cursor: "Abc..." }
const page2 = store.list({ prefix: ["logs"] }, { limit: 20, cursor: page1.cursor });
// Reverse
store.list({ prefix: ["logs"] }, { limit: 5, reverse: true });Fluent builder for version-checked transactions. All operations run in a single SQLite transaction.
const seen = store.get(["counter"])!;
const result = store
.atomic()
.check({ key: ["counter"], version: seen.version }) // fail if written since
.check({ key: ["new-key"], version: null }) // fail if key exists
.set(["counter"], 4)
.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);
} else {
console.log("Check failed, retry");
}AtomicOperation methods:
| Method | Signature | Description |
|---|---|---|
check |
(...checks: KvCheck[]): this |
Assert key versions. version: null = "must not exist". version: N = "must be at version N". |
set |
(key, value, options?): this |
options: { ttl?: number } |
delete |
(key): this |
|
sum |
(key, delta: number): this |
Add delta to the number at key when the commit runs (0.6) |
enqueue |
(payload, options?): this |
options: QueueOptions |
commit |
(): KvCommitResult | KvCommitError |
Execute all operations atomically. Returns { ok: false } if any check fails, else { ok: true, version } with the commit's one versionstamp. Once per builder. |
sum(key, delta) (0.6, Deno KV's sum for plain JS numbers). Mutations run in the order they were added, so a sum sees the commit's earlier set/delete/sum of the same key (set(k, 10).sum(k, 5).sum(k, 5) → 20; delete(k).sum(k, 3) → 3).
| Stored state at commit | Result |
|---|---|
| missing or expired | delta, no TTL |
| live number | value + delta, TTL kept |
| live non-number | commit() throws TypeError: kvs: atomic().sum() needs a number, but the key holds <kind> (PostgreSQL: ... holds a non-numeric value); the whole commit is rolled back, enqueues included |
delta not a finite number (NaN, Infinity, a string, missing) |
RangeError: kvs: atomic().sum() delta must be a finite number, got X, thrown before the transaction opens |
| result not finite | RangeError: kvs: atomic().sum() result is not a finite number (SQLite backends) |
key checked version: null in the same commit |
on PostgreSQL written with insertIfAbsent(delta): a row a concurrent plain write created since the check makes the commit { ok: false }, never a sum onto it |
The summed key gets the commit's versionstamp and watchers see the new value. Implementation: KVStore reads value, expires_at inside the IMMEDIATE transaction and writes through the same upsert as set. AsyncKVStore calls the transaction adapter's increment(key, delta, now, version): SQLite does the read-modify-write in the transaction; PostgreSQL runs its one-statement INSERT ... ON CONFLICT DO UPDATE NUMERIC upsert, which waits for and then adds to a concurrent uncommitted write (a read-then-write would have lost it; tested with a second connection holding the row lock). Unlike Deno KV there is no min/max and no KvU64: values are JS numbers, like increment(). Measured (bench/throughput.bench.ts): KVStore 104K commits/s, SQLite AsyncKVStore 21K, PostgreSQL 609.
Statuses: pending → processing (leased) → deleted on ack (or done with doneRetention > 0), or dead.
enqueue ─▶ pending ──dequeue──▶ processing(lease: locked_until, lease_token)
▲ ▲ │ acknowledge(id, token) ─▶ row deleted | done (finished_at)
│ └─ nack, attempts < max: deliver_at = now + backoff(attempts), token cleared
└──── lease expired, attempts < max: reclaimed before a dequeue, token cleared
│ nack on the last attempt, or lease expired on it (cleanQueue)
▼
dead (finished_at, last_error) ── retryDead ─▶ pending, attempts = 0
- A lease is
visibilityTimeoutms from dequeue (default 30 000), not fromdeliver_at(0.3 counted fromdeliver_at, so a message from a backlog older than 30 s was redelivered while still being processed: KVS-06). acknowledge,nackandextendLeaseall matchid AND lease_token = token AND status = 'processing'; a stale token (lease expired and the message reclaimed or redelivered) returnsfalse. The token is 122 random bits (crypto.randomUUID()), one perdequeue()call.- Delivery is at-least-once: a handler that outlives its lease without
extendLeasecan run twice.
Defaults: topic: "default", delay: 0, maxAttempts: 3. maxAttempts must be an integer >= 1 (RangeError); delay finite (negative = already due).
store.enqueue(
{ to: "user@example.com", subject: "Welcome" },
{ topic: "emails", delay: 5_000, maxAttempts: 5 },
);
// { ok: true, id: 1 }A due message wakes this store's listeners of the topic (in a microtask, never on the caller's stack).
Defaults: topic: "default", limit: 1 (integer >= 1, capped at 1000), visibilityTimeout: queue.visibilityTimeout.
- At most once per second per store instance (
RECLAIM_INTERVAL), first reclaim expired leases of every topic:UPDATE queue SET status='pending', locked_until=NULL, lease_token=NULL WHERE status='processing' AND locked_until <= now AND attempts < max_attempts(PostgreSQL: inside aMATERIALIZEDCTE withFOR UPDATE SKIP LOCKED).extendLease(…, 0)resets the interval, so a release is visible to this store's next dequeue at once; other processes see it within a second. - Pick
status = 'pending' AND deliver_at <= nowin aMATERIALIZEDCTE,ORDER BY deliver_at, id LIMIT n(PostgreSQL addsFOR UPDATE SKIP LOCKED), setstatus='processing', attempts+1, locked_until = now + visibilityTimeout, lease_token = token,RETURNINGthe message columns. Rows are sorted by(deliver_at, id)in JS (RETURNING has no order).
Why not one picker with OR (status='processing' AND locked_until <= now): SQLite answers it with MULTI-INDEX OR + USE TEMP B-TREE FOR ORDER BY, sorting the whole backlog of the topic on every dequeue (measured: queue cycle behind 10 000 due messages 17 600 → 556 ops/s). A UNION ALL of two limited branches kept the index walk but cost 40% at an empty backlog. The separate reclaim + pending-only picker walks idx_<p>queue_pending (topic, status, deliver_at) and stops at the LIMIT.
for (const msg of store.dequeue("emails", 10, { visibilityTimeout: 60_000 })) {
try {
await sendEmail(msg.payload);
store.acknowledge(msg.id, msg.token);
} catch (error) {
store.nack(msg.id, msg.token, { error: String(error) });
}
}- Missing/empty token or non-integer id:
TypeError(sync throw; async rejects). doneRetention === 0(default):DELETE ... WHERE id AND lease_token AND status='processing'. OtherwiseUPDATE status='done', finished_at=now, locked_until=NULL, lease_token=NULL.false: not leased under this token any more.
- Reads
attempts, max_attemptsof the lease (getLease), then one conditional UPDATE. No transaction needed: attempts only change on a dequeue, which changes the token. attempts >= maxAttempts→status='dead', finished_at=now, last_error. Elsestatus='pending', deliver_at = now + (options.delay ?? backoff(attempts)), last_error. Lease cleared either way.backoff(attempt): array →schedule[min(attempt, length) - 1]; function → its result, which must be finite >= 0 (RangeErrorotherwise); defaultmin(1000 * 2^(attempt-1), 300000)(1 s, 2 s, 4 s … 5 min), no jitter.erroris stored as given, truncated to 2000 chars. Listener nacks storeerror.message(orString(error)).
locked_until = now + visibilityTimeout(defaultqueue.visibilityTimeout);0releases at once (the next dequeue reclaims it, attempts are not refunded). Negative/NaN →RangeError.- Works after the lease expired as long as the message was not reclaimed yet.
Runs expireLeases first (dead-letters leases that expired on their last attempt, last_error = 'lease expired'), then SELECT ... WHERE topic AND status='dead' AND id > after ORDER BY id LIMIT limit.
retryDead(topic = "default", id?: number): number / deleteDead(topic = "default", id?: number): number
One dead message of the topic (id), or all of them. retryDead sets status='pending', attempts=0, deliver_at=now, finished_at=NULL (keeps last_error) and wakes listeners. Returns the row count.
Runs expireLeases first. One row per topic that has rows, ordered by topic; with topic given and no rows, one all-zero entry. pending = pending and due, delayed = pending not yet due, processing includes expired leases not yet reclaimed, oldestPendingAt = min deliver_at of due pending rows.
Maintenance, also run by the 60 s timer: expireLeases(now) + delete done with finished_at <= now - doneRetention + (if deadRetention finite) delete dead with finished_at <= now - deadRetention. Returns rows changed. Short-lived scripts should call it: timers are unref'd and may never fire.
Fires immediately with current values, then on every mutation to watched keys.
const { cancel } = store.watch(
[["config", "theme"], ["config", "lang"]],
(entries) => {
// entries[0] = KvEntry | null for ["config", "theme"]
// entries[1] = KvEntry | null for ["config", "lang"]
},
);
cancel(); // stop watchingInternal: Uses a watchIndex: Map<hex-encoded-key, Set<Watcher>>. On any set/delete/increment/getAsync write/atomic.commit, every watcher for a changed key fires once per batch. Deleting a key that does not exist emits nothing. One watcher per watch() call can watch multiple keys. The callback receives (entries, event) where event is a KvWatchEvent; the initial call has initial: true, changedKeys: [], and the current sequence. Callback exceptions are caught and counted in callbackErrors.
addQueueListener(topic, handler: (msg) => unknown, options?: QueueListenerOptions): { cancel: () => Promise<void> }
const listener = store.addQueueListener("emails", async (msg) => {
await sendEmail(msg.payload); // resolve → acknowledged, throw → nacked
}, { concurrency: 4, visibilityTimeout: 60_000 });
await listener.cancel(); // resolves once running handlers settled- Options:
concurrencyinteger 1..1000 (default 1,RangeErrorotherwise),visibilityTimeout(default store's),autoAck(defaulttrue). autoAck: true: handler resolves →acknowledge; throws/rejects →nack({ error: message }). The lease is renewed everyvisibilityTimeout / 2(min 10 ms) while the handler runs.autoAck: false: the handler owns the message (callacknowledge/nackwithmsg.token); the slot is held until the handler's promise settles; no renewal.cancel(): stops dequeuing, returns a promise that resolves when in-flight handlers (and their ack/nack) are done. Messages already dequeued after cancel are released (extendLease(…, 0)).close()cancels all listeners without waiting.- See "Queue dispatch internals" for the algorithm.
Cache-with-compute with singleflight deduplication. Returns Promise. It is the only async method on KVStore.
// 100 concurrent callers: fn() runs once, result cached for 30s
const ad = await store.getAsync(["ads", "venue", 42], () => fetchNextAd(42), 30_000);Algorithm:
- Check SQLite: return immediately on cache hit
- Singleflight dedup within process (coalesce concurrent calls for same key)
- Call
fn()exactly once - Store result in SQLite with TTL (if provided)
- Return to all concurrent callers
Manually delete expired entries. Returns count of deleted rows. (Auto-runs every 60s.) Each deleted key is emitted to active watchers as a null tombstone. On AsyncKVStore, tombstones require the adapter to implement the optional cleanExpiredKeys(); both built-in adapters do, a custom adapter without it only returns the count.
store.set(["cache", "a"], "x", { ttl: 1_000 });
store.set(["cache", "b"], "y", { ttl: 1_000 });
// After 2s, entries are expired: cleanExpired() removes them immediately
const deleted = store.cleanExpired(); // 2Delete ALL data from kv and queue tables. Watchers stay registered: every watched key receives a null tombstone in one batch with event.reset: true, and later writes keep being delivered. Queue listeners stay registered.
store.set(["users", "alice"], { name: "Alice" });
store.enqueue("test");
store.reset();
store.get(["users", "alice"]); // nullClose database, stop cleanup/dispatch timers, cancel all watchers/listeners. No operations work after close.
// Graceful shutdown handler
process.on("SIGINT", () => {
store.close();
process.exit(0);
});
// Or in a web framework
server.on("close", () => store.close());Same as KVStore constructor but async. See Constructor section above for connection string rules.
All methods return Promise<T>. Signatures mirror KVStore exactly:
await store.get(key: KvKey): Promise<KvEntry | null>
await store.getMany(keys: KvKey[]): Promise<(KvEntry | null)[]> // max 1000 keys
await store.set(key: KvKey, value: unknown, options?: { ttl?: number }): Promise<KvCommitResult>
await store.delete(key: KvKey): Promise<void>
await store.increment(key: KvKey, delta?: number): Promise<number> // delta default: 1
await store.list(selector: KvListSelector, options?: KvListOptions): Promise<KvListResult>
await store.enqueue(payload: unknown, options?: QueueOptions): Promise<{ ok: true, id: number }>
await store.dequeue(topic?: string, limit?: number, options?: QueueDequeueOptions): Promise<QueueMessage[]>
await store.acknowledge(id: number, token: string): Promise<boolean>
await store.nack(id: number, token: string, options?: QueueNackOptions): Promise<boolean>
await store.extendLease(id: number, token: string, visibilityTimeout?: number): Promise<boolean>
await store.listDead(topic?: string, options?: { limit?: number; after?: number }): Promise<QueueDeadMessage[]>
await store.retryDead(topic?: string, id?: number): Promise<number>
await store.deleteDead(topic?: string, id?: number): Promise<number>
await store.queueStats(topic?: string): 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<T>(key: KvKey, fn: () => T | Promise<T>, ttl?: number): Promise<T>watch() and addQueueListener() remain sync (in-process callbacks):
store.watch(keys: KvKey[], callback: WatchCallback): { cancel: () => void }
store.addQueueListener(topic: string, handler: (msg: QueueMessage) => unknown, options?: QueueListenerOptions): { cancel: () => Promise<void> }
store.getWatchDiagnostics(): KvWatchDiagnostics
store.getQueueDiagnostics(): KvQueueDiagnosticsSame fluent builder as AtomicOperation but commit() is async:
const result = await store
.atomic()
.check({ key: ["counter"], version: 3 })
.set(["counter"], 4)
.enqueue({ task: "notify" }, { topic: "jobs" })
.commit();AsyncAtomicOperation methods: check(), set(), delete(), sum(), enqueue() all return this. commit(): Promise<KvCommitResult | KvCommitError> (a sum on a non-number rejects with the TypeError above).
| Method | Parameter | Default |
|---|---|---|
KVStore(path, options) |
path |
"kv.db" |
options.maxKeySize |
2048 (bytes, encoded) |
|
options.tablePrefix |
"" |
|
options.durability |
"normal" |
|
options.queue.visibilityTimeout |
30_000 |
|
options.queue.backoff |
min(1000 * 2^(n-1), 300_000) |
|
options.queue.doneRetention |
0 |
|
options.queue.deadRetention |
Infinity |
|
set(key, value, options) |
options |
{} (no TTL) |
increment(key, delta) |
delta |
1 |
getMany(keys) |
keys.length |
max 1000 (RangeError beyond) |
atomic().sum(key, delta) |
delta |
required, finite number |
list(selector, options) |
options.limit |
100 |
options.reverse |
false |
|
enqueue(payload, options) |
options.topic |
"default" |
options.delay |
0 |
|
options.maxAttempts |
3 |
|
dequeue(topic, limit, options) |
topic |
"default" |
limit |
1 (max 1000) |
|
options.visibilityTimeout |
queue.visibilityTimeout |
|
extendLease(id, token, ms) |
ms |
queue.visibilityTimeout |
listDead(topic, options) |
options.limit / after |
100 / 0 |
addQueueListener(t, h, options) |
concurrency / autoAck |
1 / true |
openDatabase(path, options) |
path |
"kv.db" |
The AsyncKVStore uses an internal SqlAdapter interface. Built-in adapters:
| Adapter | Class | Backend |
|---|---|---|
| SQLite | SQLiteAsyncAdapter |
SQLite via bun:sql |
| PostgreSQL | PostgresAdapter |
PostgreSQL via bun:sql |
import type { SqlAdapter } from "@coderbuzz/kvs";
interface KvRow { key: Uint8Array; value: Uint8Array; version: number }
interface SqlAdapter {
migrate(): Promise<void>;
get(key: Uint8Array, now: number): Promise<KvRow | null>; // live rows only
// Optional (0.6): the live rows among `keys` (deduplicated by the store), any order,
// in one statement. Without it getMany() calls get() per key.
getMany?(keys: Uint8Array[], now: number): Promise<KvRow[]>;
// Without `version`: take the next versionstamp; with it: the one nextVersion()
// returned for this atomic() commit (KVS-05)
set(key: Uint8Array, value: Uint8Array, expiresAt: number | null, version?: number): Promise<{ version: number }>;
nextVersion(): Promise<number>; // unique store-wide, never reused
delete(key: Uint8Array): Promise<boolean>;
getVersion(key: Uint8Array, now: number): Promise<number | null>;
// Must create a missing/expired key with `delta` (no TTL), keep a live key's TTL,
// and throw on a non-numeric value (a TypeError starting "kvs: increment() " is
// reworded for atomic().sum()). Returning null = legacy "key missing" (store then
// does a non-atomic set(); atomic().sum() throws instead). `version` (0.6): given by
// atomic().sum() on the transaction adapter, the commit's versionstamp to write.
increment(key: Uint8Array, delta: number, now?: number, version?: number): Promise<{ val: number; version: number } | null>;
listAsc(start: Uint8Array, end: Uint8Array, now: number, limit: number): Promise<KvRow[]>;
listDesc(start: Uint8Array, end: Uint8Array, now: number, limit: number): Promise<KvRow[]>;
cleanExpired(now: number): Promise<number>;
cleanExpiredKeys?(now: number): Promise<Uint8Array[]>; // enables expiry tombstones
enqueue(topic: string, payload: Uint8Array, now: number, deliverAt: number, maxAttempts: number): Promise<{ id: number }>;
// Queue v2 (0.4). Numeric columns must come back as JS numbers (PostgreSQL int8 → Number).
reclaimLeases(now: number): Promise<number>; // expired leases with attempts left → pending, token cleared
dequeue(topic: string, now: number, limit: number, lockedUntil: number, token: string): Promise<QueueRow[]>; // pending only
ack(id: number, token: string, now: number, retain: boolean): Promise<boolean>; // delete, or 'done' when retain
getLease(id: number, token: string): Promise<{ attempts: number; max_attempts: number } | null>;
nack(id: number, token: string, now: number, retryAt: number | null, error: string | null): Promise<boolean>; // null → dead
extendLease(id: number, token: string, lockedUntil: number): Promise<boolean>;
expireLeases(now: number): Promise<number>; // expired, attempts >= max → dead, last_error 'lease expired'
purgeFinished(doneBefore: number, deadBefore: number | null): Promise<number>;
listDead(topic: string, limit: number, afterId: number): Promise<QueueRow[]>; // rows include finished_at
retryDead(topic: string, id: number | null, now: number): Promise<number>;
deleteDead(topic: string, id: number | null): Promise<number>;
queueStats(topic: string | null, now: number): Promise<QueueStatsRow[]>;
reset(): Promise<void>; // DELETE both tables (prefix-aware); keeps meta
transaction<T>(fn: (adapter: SqlAdapter) => Promise<T>): Promise<T>;
// Optional, called on the transaction adapter by atomic().commit(). Needed on a
// backend that runs transactions concurrently (PostgresAdapter implements all three):
lockKeys?(keys: Uint8Array[]): Promise<void>; // held until the transaction ends
getVersionForUpdate?(key: Uint8Array, now: number): Promise<number | null>; // locks the row
insertIfAbsent?(key: Uint8Array, value: Uint8Array, expiresAt: number | null, now: number): Promise<{ version: number } | null>;
close(): Promise<void>;
raw(sql: string): Promise<void>; // runs as is: table names depend on tablePrefix/schema
}
interface QueueRow {
id: number; topic: string; payload: Uint8Array; enqueued_at: number; deliver_at: number
attempts: number; max_attempts: number; locked_until: number | null; last_error: string | null
finished_at?: number | null
}migrate() must implement "Schema versions" above (create or upgrade, verify columns, refuse a newer version). 0.3 adapters (dequeue(topic, now, limit), ack(id), requeueFailed) no longer compile against 0.4.
atomic().commit() inside transaction(): lockKeys(checked ∪ mutated keys) → each check via getVersionForUpdate ?? getVersion → nextVersion() → mutations in order: a set via insertIfAbsent when its key was checked version: null (null result = { ok: false }), otherwise set; a sum via insertIfAbsent(delta) in that same case, otherwise increment(key, delta, now, version); a delete via delete → enqueues.
| Feature | SQLite | PostgreSQL |
|---|---|---|
| Key column | BLOB |
BYTEA |
| Queue ID | INTEGER PRIMARY KEY AUTOINCREMENT |
SERIAL in step 1, BIGINT + AS BIGINT sequence in step 2 |
Entry version |
INTEGER (64-bit) |
BIGINT (step 2); Bun.SQL returns int8 as string, the adapter applies Number() |
| Timestamp | INTEGER |
BIGINT |
increment(), atomic().sum() |
read-modify-write in one transaction (JS arithmetic) | one upsert, convert_from(value,'UTF8')::numeric + delta, returned as FLOAT8; errors 22P02/22021 (not a JSON number) become the kvs TypeError |
getMany() |
KVStore: one read transaction of point reads; SQLiteAsyncAdapter: SELECT ... WHERE key IN (?, …) AND (expires_at IS NULL OR expires_at > ?) |
SELECT ... WHERE key IN ($1, …) AND (expires_at IS NULL OR expires_at > $n+1) |
atomic() concurrency |
serial queue: one statement or transaction at a time | pg_advisory_xact_lock per key, SELECT ... FOR UPDATE for checks, insertIfAbsent for version: null keys |
| Concurrent dequeue | MATERIALIZED CTE picker (writers serialized) |
MATERIALIZED CTE picker with FOR UPDATE SKIP LOCKED |
| Lease reclaim | one UPDATE | MATERIALIZED CTE with FOR UPDATE SKIP LOCKED, then UPDATE ... FROM |
| Migration lock | BEGIN IMMEDIATE |
pg_advisory_xact_lock(hash("migrate:" + meta table)) |
| Table names | "<prefix>kv" |
"<schema>"."<prefix>kv" when schema is set |
| Partial indexes | WHERE expires_at IS NOT NULL |
same |
Exported standalone for deduplicating concurrent async work:
import { Singleflight } from "@coderbuzz/kvs";
const sf = new Singleflight<User>();
// 100 concurrent calls for "user:42": fetchUser() runs once
const user = await sf.do("user:42", () => fetchUser(42));
sf.clear(); // clear all in-flight
sf.size; // number of in-flight keys- WAL mode, 64 MB cache, 256 MB mmap,
busy_timeout = 5000 - TTL cleanup + queue maintenance (
cleanQueue) every 60 s - Lease reclaim before a dequeue, at most once per second
- List max: 1000 per page
- Queue dispatch interval: 1 s
- Same performance profile, async API
- Same SQL features (RETURNING, ON CONFLICT, WAL)
- Bun's SQLite
SQLclient has one connection. A statement sent while a transaction is open runs inside that transaction.SQLiteAsyncAdaptertherefore runs everything through a serial queue: one statement or one wholetransaction()at a time. Without it (kvs ≤ 0.3.1) a plainset()issued during a failingatomic()returned{ ok: true }and was rolled back, and a second concurrentatomic()threwcannot start a transaction within a transaction. Measured cost: none within noise onget(~82K vs ~83K ops/s)
- Connection pooling (configurable via connection string)
SKIP LOCKEDfor safe concurrent dequeueNUMERICarithmetic for increment (exact in storage), result returned asFLOAT8atomic()locks: advisory lock ids are FNV-1a 64 ofnamespace ‖ encodedKey, namespace"kvs:"for the default tables (same as 0.3, so 0.3 and 0.4 processes still exclude each other) and"kvs:<schema>.<prefix>:"otherwise, deduplicated and taken in ascending order (no lock-order deadlock). A hash collision only makes two keys share a lock. They can collide with an application's ownpg_advisory_xact_lock(bigint)ids in the same database; that only serializes, it never breaks correctnessBYTEAfor binary key/value storage
- All internal timers are
unref()'d: an open store never keeps the process alive. A script that forgetsclose()exits normally (kvs ≤ 0.3.1 hung forever). - Maintenance: one 60s timer.
KVStorestarts it in the constructor;AsyncKVStorestarts it after the first operation runsmigrate(). Errors are counted inwatchDiagnostics.dispatchErrors.- Cleanup deletes rows where
expires_at IS NOT NULL AND expires_at <= now(DELETE ... RETURNING key) and emits tombstones. cleanQueue(): dead-letter expired last-attempt leases, applydoneRetention/deadRetention.
- Cleanup deletes rows where
- Queue poll: every 1s while at least one listener exists, every listener is notified (picks up delayed messages, retries, reclaimed leases, other processes' enqueues).
- Lease renewal: one interval per listener with in-flight handlers (
autoAckonly), everyvisibilityTimeout / 2.
watchIndex: Map<hex-encoded-key, Set<Watcher>>- Every successful mutation creates a committed mutation holding the encoded
key, the encoded value bytes (or null for a tombstone) and the version. Values
are not re-read from storage. The
KvEntryis decoded from those bytes only for keys some watcher watches, once per batch; with no watchers a write pays nothing for watch support. (0.3.0–0.3.1 decoded every write, costing ~36–56% ofset()throughput.) - Atomic operations collapse repeated keys to their final state and emit one batch after the transaction commits. A watcher matching multiple changed keys is invoked once for that batch.
- A batch-local entry cache is seeded by changed entries. Every unique unchanged key required by legacy multi-key callbacks is fetched at most once, shared by all matching watchers.
- Sync callbacks run sequentially after commit. Async batches enter one ordered promise chain; watchers are captured at enqueue time and checked for active state before callback, preventing post-cancel sends and stale completion order.
WatchCallbackis(entries, event?) => void.event.sequenceis monotonic only within the current process; it is not a durable database revision.watch()fires immediately withinitial: true.cleanExpired()emits tombstones.reset()emitsreset: truetombstones and keeps watchers active.getWatchDiagnostics()returns active watcher, batch, callback, shared-read, callback-error, and dispatch-error counters without high-cardinality labels.
queueWorkers: Map<topic, Set<QueueWorker>>, oneQueueWorker(src/queue.ts) peraddQueueListener(); the same class servesKVStoreandAsyncKVStore(it awaits whatever the store returns).notify()setsagain = trueand, if no pump is running, schedules one withqueueMicrotask. Callers:enqueue()andatomic()commits with a due message on the topic,retryDead(), the 1 s poll, a finished handler, the listener's own start. The handler therefore never runs onenqueue()'s stack: a listener that enqueues follow-up work to its own topic keeps the stack flat (5 001 chained messages: depth 1).pump():while (active && again) { again = false; while (inFlight < concurrency) { msgs = await dequeue(topic, concurrency - inFlight, { visibilityTimeout }); if (!msgs.length) break; start each } }. A dequeue failure (closed store, lost connection) ends the pump and countsdispatchErrors; the 1 s poll retries (no hot loop). A wake-up that lands while the pump finishes is not lost:finallyre-notifies whenagainis set.run(msg):await handler(msg); withautoAck,acknowledgeornack; afalseresult countsleaseLost, a throwdispatchErrors. Then remove from in-flight, stop renewal when idle, resolvecancel()waiters,notify().- Renewal skips deliveries whose handler already settled, so a renewal that crosses the ack is not counted as a lost lease.
- Several listeners of one topic compete: whoever has a free slot dequeues. There is no round-robin (0.3 had one, over a synchronous drain).
getQueueDiagnostics()sums listeners and in-flight handlers over all workers; counters are shared by the store.
KVStore.get()returnsnullfor expired entries (TTL respected).AtomicOperation.check({ version: null })means "key must NOT exist". This is the opposite of checking a version number.watch()fires immediately with current values, not just on future changes.addQueueListener()handlers are awaited and auto-acked (resolve) or nacked (throw) since 0.4. Do not also callacknowledge()from anautoAckhandler (the second ack returnsfalseand countsleaseLost). UseautoAck: falseto own the message.getAsync()uses the hex of the encoded key as the singleflight dedup key, so keys that encode identically share one flight. Dedup is per store instance (in-process only).- Bun ≥ 1.2.21 only (
engines.bun, 0.6).KVStoreusesbun:sqlite,AsyncKVStoreBun.SQL(SQLite there since Bun 1.2.21). The package does not load on Node or Deno; use@coderbuzz/kvs-clientagainst a kvs-server running on Bun. close()stops all timers, cancels all watchers, and closes the database. No operations work after close.- SQLite WAL means concurrent readers are fine, but writers are serialized.
- Engine minimums: the queue picker is a
WITH ... AS MATERIALIZEDCTE, which requires PostgreSQL 12+ and SQLite 3.35+;RETURNINGalso requires SQLite 3.35+. Bun 1.4 bundles SQLite 3.53, so only the PostgreSQL server version needs checking. dequeue(topic, limit)returns at mostlimitrows, and ties indeliver_atbreak onid ASC. Both are load-bearing. The picker must stay inside the materialized CTE: on PostgreSQL an equivalentid IN (SELECT ... LIMIT n FOR UPDATE SKIP LOCKED)sublink can be planned on the inner side of a nested-loop semi join with noMaterializenode, rescanning the picker once per outer row. Each rescan re-runsSKIP LOCKEDagainst the rows the previous iteration locked, returns a different window, and every candidate row ends up markedprocessing. This delivers a whole backlog to one worker while the call reports the limit it was given.deliver_atis a millisecond timestamp, so enqueue bursts tie constantly; without theidtiebreaker FIFO order is whatever the planner produces.atomic().commit()returns{ ok: true, version }: the one versionstamp of the commit, which everysetof it carries, also when it has noset(0.4: the last set's version, or0). A builder commits once.new AsyncKVStore("kv.db")does not open SQLite. Bun'sSQLtreats a protocol-less filename as PostgreSQL. Use"sqlite://kv.db".- Watch
sequencestarts at 0 per store instance and increments per committed batch (including batches no watcher matches). It is not persisted. KvWatchEvent.coalescedis declared but never set by the core store, and@coderbuzz/kvs-serverdoes not send it. Treat it as reserved.- Validation errors are
RangeError/TypeError:list()limit/cursor,ttl, queuedelay(must be finite; negative means already due),increment()/atomic().sum()delta andgetMany()over 1000 keys throwRangeError;increment()/sumon a non-number, alist()bound outside its prefix and a non-arraygetMany()argument throwTypeError. All start withkvs:, which is how kvs-server tells them apart (400). OnAsyncKVStorethey reject the returned promise. Validation happens before any write. atomic().commit()is safe under concurrency on every backend. Two commits that check the sameversion(orversion: null) never both returnok: true, including write skew (A checks B absent and writes A, B checks A absent and writes B: exactly one wins). On PostgreSQL a check waits for an uncommitted write to the checked row and then re-reads it; aversion: nullkey that a concurrent plainset()creates before the commit makes the commit return{ ok: false }.- Custom
SqlAdapteron a concurrent backend must implementlockKeys,getVersionForUpdateandinsertIfAbsentfor 16 to hold; without thematomic()falls back to plaingetVersion()+set(). Itsincrement(key, delta, now, version?)must create missing/expired keys itself and, when called withversion(byatomic().sum()), write that versionstamp; returningnull(the pre-0.3.2 contract) makes the store fall back to a non-atomicset(), and makesatomic().sum()throw. - Do not call a
SQLiteAsyncAdapterfrom inside its owntransaction()callback except through the adapter the callback receives: the outer adapter waits for the transaction to finish, so the call never completes. acknowledge/nack/extendLeaseneedmsg.token. A 0.3-styleacknowledge(id)throwsTypeError.- A message is at-least-once. A handler slower than its lease (without renewal:
autoAck: false, or a plaindequeue()loop) can see its message redelivered to another consumer; its lateacknowledgethen returnsfalse. - Expired leases are reclaimed lazily: by a dequeue (at most once a second per store) or by
cleanQueue()/listDead()/queueStats()(dead-lettering only). Until thenqueueStats().processingstill counts them. - Shared databases: a table named
kv,queueormetathat lacks the kvs columns makesnew KVStore()throw / the firstAsyncKVStoreoperation reject, and is never altered. The check is by column names only: an application table that happens to have them (e.g.meta(key, value)) passes, and kvs then writes itsschema_versionrow into it. UsetablePrefix(andschemaon PostgreSQL) whenever kvs shares a database.reset()deletes only the store's own tables. - Upgrading a big PostgreSQL database from 0.3 rewrites
kvandqueueonce (ALTER COLUMN ... TYPE BIGINT) under an exclusive lock; every process opening the database waits for the migration. dequeue()from a 1-row topic can still return nothing right after a release in another process: other processes reclaim at most once a second.- Versions are not 1, 2, 3 per key (0.5). Compare a version only for equality with one you read (
atomic().check). On SQLite two processes' versions are unique but not time-ordered. list({ prefix })excludes the prefix key itself (0.5). Read it withget(), or use{ start: p, end: [...p, 0] }-style ranges.- A
Date/bigint/Mapvalue read over kvs-server comes back as JSON (ISO string, decimal string,{}or entries asJSON.stringifygives): the HTTP/WS API is JSON. Only the local store API keeps the types. - Do not write through kvs from inside your own
store.db.transaction()(KVStore): a version block reserved there would roll back with your transaction while the store keeps handing it out. atomic().sum()returns no value. The commit result is{ ok, version }, as in Deno KV. Read the counter afterwards, or useincrement()(outside a commit) when you need the new number back.sumandincrement()share semantics: same TTL rules, same errors, same PostgreSQL NUMERIC arithmetic.getMany()is at most 1000 keys and returnsnullslots, not{ value: null, versionstamp: null }entries as Deno KV does. OnKVStoreit is no faster than a loop ofget(); its point is one snapshot. For more keys, split the call (each call is its own snapshot) or uselist()over a range.- In
list({ prefix, start, end })the bounds are keys inside the prefix, not suffixes:{ prefix: ["orders"], start: ["orders", 5] }, notstart: [5](aTypeError, because[5]is outside["orders"]). The prefix key itself is not a valid bound either.
@coderbuzz/kvs-server:createServer(store)for sync,createAsyncServer(store)for async. Exposes the queue v2 routes (/queue/nack,/queue/extend,/queue/dead,/queue/stats, …) and WebSocket listeners that holdconcurrencyslots until the client acks.@coderbuzz/kvs-client: TypeScript SDK for the server