Skip to content

KV state

A transactional key/value store inside the broker: seven operations, an expiry on every write, and the one property a key/value store standing beside Queen cannot give at any price.

Updated View as Markdown

queen.kv is a key/value store that lives in the same PostgreSQL as the log, and its value is not that it stores keys. It is that a write to it can share the transaction with a push and an ack.

The idempotency marker, the effect and the cursor advance commit together or not at all. A key/value store standing beside the broker cannot give you that, at any price and however fast it is, because its commit and Queen’s commit are two commits and something can happen between them.

The one thing a store on the side cannot do

The shape below is the reason this feature exists.

// Illustrative, not extracted from a test.
const tx = queen.transaction()
  .once('charges', `charge:${order.id}`, { ttlSeconds: 86400, required: true })
  .ack(message)

await chargeCard(order)          // the external effect
const res = await tx.commit()    // marker and cursor, one commit

if (!res.success && res.reason === 'kv_precondition') {
  // Somebody already did this work, and NOTHING in the bundle was committed:
  // the ack rolled back together with the marker. So this message is still on
  // the group's cursor, and it has to be taken off separately or it comes back
  // forever. `res.value` carries the winner's value, so the charge id is here
  // without a second round trip.
  await queen.ack(message, 'completed', { group: message.consumerGroup })
}

once is putIfAbsent under the name of the question it answers. With required: true a lost precondition rolls the whole bundle back, so a redelivery of a message whose work was already done finds the marker present and aborts. That verdict arrives as HTTP 200 with {"success": false, "reason": "kv_precondition"} and never as an exception, because a lost race is the most frequent outcome of this product and it must not land in a retry policy or an error metric.

The property this buys is exact and it is worth stating in the negative first, because the promise is easy to inflate. It is not exactly-once end to end: the charge is outside PostgreSQL, so a crash between the charge and the commit repeats the charge on redelivery, and no broker can prevent that.

What the shared commit does give is the pair of guarantees a store on the side cannot:

  • The marker and the cursor advance are one commit, so there is no state in which the work is marked done and the message will be redelivered, or acknowledged and unmarked. With two systems there is a window between two commits, and whichever order you choose, a crash inside it produces one of those two.
  • A handler that fails before the commit leaves no marker, so the retry is not blocked by the attempt that failed. A marker written in its own transaction, before the work, does exactly that.

Seven operations

Operation What it does
get One key. Answers {found, key, value, version, expiresAt, updatedAt}
getMany A list of keys, answering rows and an explicit missing
getPrefix A keyset page under a prefix, with rows, truncated and nextAfter
put Upsert, or a conditional write when expect is present
putIfAbsent put with expect: 0, under the name of the thing
delete Removes a key, optionally under expect
incr Atomic numeric add, with optional min and max

POST /api/v1/kv is the complete surface and the only one that accepts incr and getPrefix. The path routes (GET, PUT and DELETE on /api/v1/kv/:ns/*key) are sugar for the three cases people write by hand, and an SDK is not one of them.

Namespaces are registered nowhere. Like a queue, a namespace exists if and only if a row exists, so an unknown namespace reads empty and is never an error. The charset is validated precisely because nothing registers it: without validation a typo would not fail, it would mint a phantom namespace that reads empty forever.

Two shapes of the answers are deliberate and worth reading once. found is separate from the value, because null is a legal stored value and {found: true, value: null} is not the same fact as {found: false}. And multi-key reads return rows, never a key-to-value map, so the shape itself makes the confusion between “absent” and “present and null” inexpressible.

Every write carries an expiry

put, putIfAbsent and incr each carry exactly one of ttlSeconds (an integer above zero) or forever: true. Zero or two of them is a 400. The rule lives in SQL rather than in seven clients, so it applies identically to the HTTP routes, the transaction wire and the embedded broker.

This is the difference between a store you can leave running and one you cannot. queen_streams.state has no expiry because a closed window deletes its own row; a KV key has no such event, so an optional TTL would mean a table that grows in silence and an operator who pays for it.

Two consequences follow, and both are the kind that bite once.

A put does not inherit the previous expiry. It is not expressible, and a put that silently inherited a TTL is the fastest way to make an idempotency marker immortal.

And an expired key is never returned and never counts as existing, even before the sweeper has pruned it. The truth is the predicate, not the presence of the row. A putIfAbsent wins against an expired row that is still physically there, and an incr sees an expired counter as zero and starts a fresh window, which is what makes a fixed-window limiter a single call.

Where a read-modify-write is safe, and where it is not

Reading a key in one call and writing it in the next is safe only when the key derives from the partition key. Then the lane serialises the writers and no one else in that consumer group touches the key. When it does not derive from the partition key, use the atomics.

Queen has no lease on a KV key. Two workers can hit the same key at the same instant, which is exactly why expect, putIfAbsent and incr are primitives rather than conveniences. The SDK state handle mints the key for you as @p/<queue>/<partition>/<group>/<name>, so the derivation is imposed by the API rather than remembered by you.

expect is how you make your serialisation assumption falsifiable instead of silent. If you believe the lane serialises you, say so anyway: if it never fails it cost nothing, and the day it fails you have just discovered that two consumers are serving the same partition, and you discovered it as a verdict rather than as a wrong total.

expect Meaning
Absent Unconditional upsert. Replaces the value and the expiry
0 Must not exist. Wins against an expired row that has not been pruned yet
A version Optimistic lock. Never creates a row, so a lost expect writes nothing

That last row is load-bearing. A conditional write that matches nothing must create nothing, or a saga would start the very compensation the expect existed to prevent.

Every write answers with the current value and version even when it did not apply, so the loser of a race needs no second round trip. The reason taxonomy is closed: exists, absent, version, limit, type.

Two ordering rules complete the picture. The transaction is the primary fence and expect is the secondary assertion: a state write sharing the transaction with an ack is undone when an expired lease makes that ack raise, which a compare-and-set cannot do, because an expect on a version that still matches succeeds from a zombie as happily as from the rightful holder. And a bundle may touch a given key at most once, so there is no intra-batch evaluation order to reason about.

incr is the way out of the retry loop

incr has no expect, deliberately: it exists to remove the compare-and-set loop, and a precondition would put the loop back. The value is numeric server side, so nothing overflows; typed SDKs expose int64 and fail loudly rather than hand back a number that lost precision.

With max, applied is the admission decision. If the increment would breach the ceiling nothing is written and the answer is applied: false, reason: "limit" with the current value. It does not saturate and it does not truncate, because a limiter that clamps has already spent the budget for the request that broke the ceiling and cannot give it back.

The TTL of incr is create-only. A live row keeps the expiry it was born with. If an increment extended it, a fixed window over a continuously active caller would never close, which is to say the limiter would stop limiting exactly under load.

Reads, and the boundary that is about cost

get and getMany are allowed inside a transaction because the caller fixes their cost. getPrefix is not, and the boundary is the cost rather than the kind of operation: it is unbounded read work inside the transaction that holds the outermost lock space and, downstream, partition locks.

getPrefix also never appears in a query string. ?prefix=quota:acme: would pass through the broker’s access logs, the proxy’s, the metering sample, the tracing span and any ingress in front, and a mitigation living in one component out of four is not a mitigation. It is available on POST /api/v1/kv and nowhere else, and it requires a prefix: a namespace is not a table to enumerate.

Page limits are clamped rather than refused, and truncated tells the truth. There is a byte ceiling as well as a key ceiling, QUEEN_KV_MAX_READ_BYTES (4 MiB), because a thousand keys of 64 KiB is 64 MB and the real resource is the byte. after is an exclusive keyset cursor, not an offset. Ordering is byte order under the C collation, so non-ASCII keys are not in the alphabetical order of a locale, and each page is its own snapshot: good for compacting state, not for an exact count.

queen_streams.state is a different thing with a similar name

Both are state in PostgreSQL and only one of them is atomic with a stream cycle.

queen_streams.state queen.kv
Identity (query_id, partition_id, key) (tenant, namespace, key)
Atomic with The stream cycle: state, sink push and ack in one transaction The transaction wire, when the write rides a bundle
Reachable across partitions No Yes
Expiry None. A closed window deletes its own row Mandatory on every write
Written by Operators, through state_ops Anybody, through the KV surface

Inside a stream the state primitive stays state_ops. It is older, less visible, and its atomicity with the cycle is free; picking the KV instead silently gives that atomicity up. The KV earns its place inside a stream for one job only, the job state_ops cannot do, which is state that crosses partitions or queries. Streams states the rule where an operator author will meet it.

What does not exist

No query by value, no predicates, no secondary indexes, no listing without a prefix, no watch, subscribe or long poll, no merge, no deletePrefix, no compare-and-set on the value, and no value above 64 KiB. Keys are capped at 512 bytes.

Values are never cached by the broker, at any TTL. The first use of this store is an idempotency marker, and a stale read of a marker says “not there” and performs the external effect twice: a cache would make the transactional primitive eventually consistent in precisely the place it exists not to be.

Navigation

Type to search…

↑↓ navigate↵ selectEsc close