tuwunel_database/txn.rs
1//! Atomic database writes backed by one RocksDB write batch.
2//!
3//! A transaction queues operations for maps owned by one database engine and
4//! commits them only when [`Txn::execute`] consumes it. Typed operations use
5//! the database codec, while raw operations preserve caller-provided bytes.
6
7use std::{fmt::Debug, iter::once, sync::Arc};
8
9use rocksdb::WriteBatch;
10use serde::Serialize;
11use tuwunel_core::{error, implement};
12
13use crate::{
14 Engine, Map,
15 keyval::{Key, Slice, serialize_key, serialize_val},
16 util::or_else,
17};
18
19/// Atomic write batch spanning one or more column families from one database.
20///
21/// Every queued map must belong to the captured engine because column family
22/// identifiers are interpreted within that database. Dropping an unexecuted
23/// transaction leaves the database unchanged.
24#[must_use = "does nothing until execute()"]
25pub struct Txn {
26 batch: WriteBatch,
27 engine: Arc<Engine>,
28}
29
30/// Record parser yielding each queued key with its resolved map.
31struct Keys<'a> {
32 engine: &'a Engine,
33 data: Key<'a>,
34}
35
36/// Batch representation header: a fixed64 sequence then a fixed32 count.
37const HEADER: usize = 12;
38
39/// Worst-case per-record overhead: a type tag and three varint32s.
40const PER_OP: usize = 16;
41
42/// Record tags per rocksdb `write_batch.cc`; puts and deletes against
43/// column family id 0 encode as the legacy untagged types.
44#[derive(Clone, Copy)]
45enum Tag {
46 Deletion = 0x0,
47 Value = 0x1,
48 CfDeletion = 0x4,
49 CfValue = 0x5,
50}
51
52impl TryFrom<u8> for Tag {
53 type Error = u8;
54
55 fn try_from(byte: u8) -> Result<Self, Self::Error> {
56 match byte {
57 | 0x0 => Ok(Self::Deletion),
58 | 0x1 => Ok(Self::Value),
59 | 0x4 => Ok(Self::CfDeletion),
60 | 0x5 => Ok(Self::CfValue),
61 | unrecognized => Err(unrecognized),
62 }
63 }
64}
65
66/// Creates an empty transaction for one database engine.
67///
68/// Operations can be appended through the typed or raw queueing methods. The
69/// transaction remains inert until [`Txn::execute`] consumes it.
70#[implement(Txn)]
71pub fn new(engine: &Arc<Engine>) -> Self {
72 Self {
73 batch: WriteBatch::default(),
74 engine: engine.clone(),
75 }
76}
77
78/// Creates an empty transaction with reserved batch capacity.
79///
80/// `capacity_bytes` reserves storage for the serialized RocksDB batch
81/// representation. The reservation affects allocation only and does not queue
82/// an operation.
83#[implement(Txn)]
84pub fn with_capacity_bytes(engine: &Arc<Engine>, capacity_bytes: usize) -> Self {
85 Self {
86 batch: WriteBatch::with_capacity_bytes(capacity_bytes),
87 engine: engine.clone(),
88 }
89}
90
91/// Queues raw key and value pairs for one map from a single pass.
92///
93/// The database codec is not applied, and the write batch copies each supplied
94/// byte sequence. Empty input produces an empty transaction whose execution is
95/// a no-op.
96#[implement(Txn)]
97pub fn insert<I, K, V>(map: &Map, items: I) -> Self
98where
99 I: IntoIterator<Item = (K, V)>,
100 K: AsRef<Slice>,
101 V: AsRef<Slice>,
102{
103 items
104 .into_iter()
105 .fold(Self::new(map.engine()), |mut txn, (key, val)| {
106 txn.insert_raw(map, key, val);
107 txn
108 })
109}
110
111/// Queues a raw slice for one map with a precomputed capacity estimate.
112///
113/// The estimate includes payload lengths and worst-case record overhead before
114/// the items are copied into the write batch. Empty input produces an empty
115/// transaction.
116#[implement(Txn)]
117pub fn insert_slice<K, V>(map: &Map, items: &[(K, V)]) -> Self
118where
119 K: AsRef<Slice>,
120 V: AsRef<Slice>,
121{
122 let capacity_bytes = size_hint(items.iter().map(|(key, val)| (key, val)));
123
124 items.iter().fold(
125 Self::with_capacity_bytes(map.engine(), capacity_bytes),
126 |mut txn, (key, val)| {
127 txn.insert_raw(map, key, val);
128 txn
129 },
130 )
131}
132
133/// Queues raw entries across maps from a nonempty single pass.
134///
135/// The first item selects the database engine, and every subsequent map must
136/// belong to that same engine. The database codec is not applied to keys or
137/// values.
138///
139/// # Panics
140///
141/// Panics when `items` is empty or when any map belongs to a different database
142/// engine.
143#[implement(Txn)]
144pub fn insert_each<'a, I, K, V>(items: I) -> Self
145where
146 I: IntoIterator<Item = (&'a Map, K, V)>,
147 K: AsRef<Slice>,
148 V: AsRef<Slice>,
149{
150 let mut items = items.into_iter();
151 let (map, key, val) = items
152 .next()
153 .expect("insert_each: at least one item");
154
155 let mut txn = Self::new(map.engine());
156
157 txn.insert_raw(map, key, val);
158 txn.extend(items);
159
160 txn
161}
162
163/// Queues a nonempty raw slice across maps with a capacity estimate.
164///
165/// The first item selects the database engine, and every map must belong to
166/// that same engine. The database codec is not applied to keys or values.
167///
168/// # Panics
169///
170/// Panics when `items` is empty or when any map belongs to a different database
171/// engine.
172#[implement(Txn)]
173pub fn insert_each_slice<K, V>(items: &[(&Map, K, V)]) -> Self
174where
175 K: AsRef<[u8]>,
176 V: AsRef<[u8]>,
177{
178 let map = items
179 .first()
180 .expect("insert_each_slice: at least one item")
181 .0;
182
183 let capacity_bytes = size_hint(items.iter().map(|(_, key, val)| (key, val)));
184
185 let mut txn = Self::with_capacity_bytes(map.engine(), capacity_bytes);
186
187 txn.extend(
188 items
189 .iter()
190 .map(|(map, key, val)| (*map, key, val)),
191 );
192
193 txn
194}
195
196/// Serializes and queues entries across maps from a nonempty pass.
197///
198/// The first item selects the database engine, and every map must belong to
199/// that same engine. All keys and values are encoded with the database record
200/// codec before being copied into the batch.
201///
202/// # Panics
203///
204/// Panics when `items` is empty, a map belongs to another database engine, or
205/// serialization of a key or value fails.
206#[implement(Txn)]
207pub fn put_each<'a, I, K, V>(items: I) -> Self
208where
209 I: IntoIterator<Item = (&'a Map, K, V)>,
210 K: Serialize + Debug,
211 V: Serialize,
212{
213 let mut items = items.into_iter();
214 let (map, key, val) = items.next().expect("put_each: at least one item");
215 let txn = Self::new(map.engine());
216
217 once((map, key, val))
218 .chain(items)
219 .fold(txn, |mut txn, (map, key, val)| {
220 txn.put(map, key, val);
221 txn
222 })
223}
224
225/// Serializes and queues one insertion.
226///
227/// The key and value use the database record codec, and the operation remains
228/// pending until [`Txn::execute`]. The map must belong to the transaction's
229/// database engine.
230///
231/// # Panics
232///
233/// Panics when the map belongs to another database engine or serialization of
234/// the key or value fails.
235#[implement(Txn)]
236pub fn put<K, V>(&mut self, map: &Map, key: K, val: V)
237where
238 K: Serialize + Debug,
239 V: Serialize,
240{
241 self.assert_map(map);
242
243 let key = serialize_key(key).expect("failed to serialize batch key");
244 let val = serialize_val(val).expect("failed to serialize batch val");
245
246 self.batch.put_cf(&map.cf(), key, val);
247}
248
249/// Serializes the key and queues one raw-value insertion.
250///
251/// The key uses the database record codec, while the value bytes are copied
252/// unchanged into the batch. The operation remains pending until
253/// [`Txn::execute`], and the map must belong to the transaction's database
254/// engine.
255///
256/// # Panics
257///
258/// Panics when the map belongs to another database engine or serialization of
259/// the key fails.
260#[implement(Txn)]
261pub fn put_raw<K, V>(&mut self, map: &Map, key: K, val: V)
262where
263 K: Serialize + Debug,
264 V: AsRef<Slice>,
265{
266 self.assert_map(map);
267
268 let key = serialize_key(key).expect("failed to serialize batch key");
269
270 self.batch.put_cf(&map.cf(), key, val);
271}
272
273/// Queues one raw-key insertion after serializing the value.
274///
275/// The key bytes are copied unchanged into the batch, while the value uses the
276/// database record codec. The operation remains pending until [`Txn::execute`],
277/// and the map must belong to the transaction's database engine.
278///
279/// # Panics
280///
281/// Panics when the map belongs to another database engine or serialization of
282/// the value fails.
283#[implement(Txn)]
284pub fn raw_put<K, V>(&mut self, map: &Map, key: K, val: V)
285where
286 K: AsRef<Slice>,
287 V: Serialize,
288{
289 self.assert_map(map);
290
291 let val = serialize_val(val).expect("failed to serialize batch val");
292
293 self.batch.put_cf(&map.cf(), key, val);
294}
295
296/// Serializes and queues one deletion.
297///
298/// The key uses the database record codec, and the operation remains pending
299/// until [`Txn::execute`]. The map must belong to the transaction's database
300/// engine.
301///
302/// # Panics
303///
304/// Panics when the map belongs to another database engine or serialization of
305/// the key fails.
306#[implement(Txn)]
307pub fn del<K>(&mut self, map: &Map, key: K)
308where
309 K: Serialize + Debug,
310{
311 self.assert_map(map);
312
313 let key = serialize_key(key).expect("failed to serialize batch key");
314
315 self.batch.delete_cf(&map.cf(), key);
316}
317
318/// Queues one deletion for an already serialized key.
319///
320/// The key bytes are copied into the write batch without invoking the database
321/// codec. The map must belong to the transaction's database engine.
322///
323/// # Panics
324///
325/// Panics when the map belongs to another database engine.
326#[implement(Txn)]
327pub fn del_raw<K>(&mut self, map: &Map, key: K)
328where
329 K: AsRef<Slice>,
330{
331 self.assert_map(map);
332 self.batch.delete_cf(&map.cf(), key);
333}
334
335/// Commits the batch atomically, flushes unless corked, and notifies matching
336/// watchers.
337///
338/// An empty transaction returns without touching the engine. For a nonempty
339/// batch, notifications occur only after the write.
340///
341/// # Panics
342///
343/// Panics when RocksDB rejects the batch write or when the required database
344/// flush fails.
345#[implement(Txn)]
346#[inline]
347pub fn execute(self) {
348 if self.is_empty() {
349 return;
350 }
351
352 self.commit();
353}
354
355/// Commits the batch atomically, flushes unless corked, and notifies matchers.
356///
357/// Batch must not be empty. Notifications occur only after the write.
358///
359/// # Panics
360///
361/// Panics when RocksDB rejects the batch write or when the required database
362/// flush fails.
363#[implement(Txn)]
364#[tracing::instrument(
365 level = "trace",
366 skip_all,
367 fields(
368 ops = self.len(),
369 bytes = self.size_in_bytes(),
370 )
371)]
372fn commit(self) {
373 debug_assert!(!self.is_empty(), "Txn must not be empty.");
374
375 self.engine
376 .db
377 .write_opt(&self.batch, &self.engine.write_options)
378 .or_else(or_else)
379 .expect("database transaction execute error");
380
381 if !self.engine.corked() {
382 self.engine
383 .flush()
384 .inspect_err(|e| error!(?e, "database flush error"))
385 .ok();
386 }
387
388 self.notify();
389}
390
391/// Notifies watchers after a successful commit for queued keys that resolve to
392/// catalog maps.
393///
394/// Keys are parsed lazily from the batch representation and consumed in queue
395/// order. Operations without a live map in the engine's startup catalog are
396/// skipped.
397#[implement(Txn)]
398fn notify(&self) {
399 for (map, key) in self.keys() {
400 map.notify(key);
401 }
402}
403
404/// Iterate queued put and delete keys in insertion order.
405///
406/// The iterator borrows keys directly from the serialized write batch without
407/// materializing a container. Keys whose column families are outside the
408/// startup map catalog are omitted.
409///
410/// # Panics
411///
412/// Iteration panics if a record has an unsupported operation tag, is truncated,
413/// or contains a varint whose fifth byte retains its continuation bit.
414#[implement(Txn)]
415pub fn keys(&self) -> impl Iterator<Item = (Arc<Map>, &Slice)> + '_ {
416 let data = self.batch.data();
417
418 Keys {
419 engine: &self.engine,
420 data: data.get(HEADER..).unwrap_or_default(),
421 }
422}
423
424/// Returns the number of operations queued in the batch.
425///
426/// Both insertions and deletions count as one operation. Inspecting the count
427/// does not execute the transaction.
428#[implement(Txn)]
429#[inline]
430#[must_use]
431pub fn len(&self) -> usize { self.batch.len() }
432
433/// Reports whether the batch contains no queued operations.
434///
435/// A newly created or cleared transaction is empty. Executing an empty
436/// transaction performs no database work.
437#[implement(Txn)]
438#[inline]
439#[must_use]
440pub fn is_empty(&self) -> bool { self.batch.is_empty() }
441
442/// Returns the encoded size of the RocksDB write batch in bytes.
443///
444/// The size includes batch metadata and queued record data. Inspecting it does
445/// not execute the transaction.
446#[implement(Txn)]
447#[inline]
448#[must_use]
449pub fn size_in_bytes(&self) -> usize { self.batch.size_in_bytes() }
450
451/// Removes every queued operation from the transaction.
452///
453/// The captured database engine remains attached, so the transaction can be
454/// populated again. Executing it before another operation is queued is a no-op.
455#[implement(Txn)]
456#[inline]
457pub fn clear(&mut self) { self.batch.clear(); }
458
459/// Queue one unencoded key and value after enforcing map ownership.
460///
461/// Both byte sequences are copied into the write batch without invoking the
462/// database codec. The operation remains pending until [`Txn::execute`].
463///
464/// # Panics
465///
466/// Panics when the map belongs to another database engine.
467#[implement(Txn)]
468pub fn insert_raw<K, V>(&mut self, map: &Map, key: K, val: V)
469where
470 K: AsRef<Slice>,
471 V: AsRef<Slice>,
472{
473 self.assert_map(map);
474 self.batch.put_cf(&map.cf(), key, val);
475}
476
477/// Verifies that a map belongs to the transaction's database engine.
478///
479/// RocksDB identifies column families numerically within one database, so
480/// accepting a foreign map could target a same-numbered column family in the
481/// captured engine.
482///
483/// # Panics
484///
485/// Panics when `map` belongs to a different database engine.
486#[implement(Txn)]
487#[inline]
488fn assert_map(&self, map: &Map) {
489 assert!(
490 Arc::ptr_eq(&self.engine, map.engine()),
491 "transaction map belongs to a different database"
492 );
493}
494
495impl<'a> Iterator for Keys<'a> {
496 type Item = (Arc<Map>, Key<'a>);
497
498 fn next(&mut self) -> Option<Self::Item> {
499 while !self.data.is_empty() {
500 let (cf_id, key) =
501 next_record(&mut self.data).expect("malformed write batch representation");
502
503 if let Some(map) = self.engine.map_by_cf_id(cf_id) {
504 return Some((map, key));
505 }
506 }
507
508 None
509 }
510}
511
512/// Extends this transaction with raw insertions across maps.
513///
514/// Each tuple queues its raw key and value through [`Txn::insert_raw`]. Use
515/// [`Txn::put_each`] when the keys and values need serialization. Every map
516/// must belong to the transaction's database engine.
517///
518/// # Panics
519///
520/// Panics when any map belongs to another database engine.
521impl<'a, K, V> Extend<(&'a Map, K, V)> for Txn
522where
523 K: AsRef<Slice>,
524 V: AsRef<Slice>,
525{
526 fn extend<I>(&mut self, items: I)
527 where
528 I: IntoIterator<Item = (&'a Map, K, V)>,
529 {
530 for (map, key, val) in items {
531 self.insert_raw(map, key, val);
532 }
533 }
534}
535
536/// Extends this transaction with raw-key deletions across maps.
537///
538/// Each tuple queues its raw key through [`Txn::del_raw`]. Every map must
539/// belong to the transaction's database engine.
540///
541/// # Panics
542///
543/// Panics when any map belongs to another database engine.
544impl<'a, K> Extend<(&'a Map, K)> for Txn
545where
546 K: AsRef<Slice>,
547{
548 fn extend<I>(&mut self, items: I)
549 where
550 I: IntoIterator<Item = (&'a Map, K)>,
551 {
552 for (map, key) in items {
553 self.del_raw(map, key);
554 }
555 }
556}
557
558/// Decodes one record into its column family identifier and borrowed key.
559///
560/// Value payloads are skipped after their lengths are consumed. Unsupported
561/// tags, truncated fields, and varints whose fifth byte retains its
562/// continuation bit return `None` and may leave the input advanced through the
563/// parsed prefix.
564pub(crate) fn next_record<'a>(data: &mut Key<'a>) -> Option<(u32, Key<'a>)> {
565 let (&tag, rest) = data.split_first()?;
566 *data = rest;
567
568 let tag = Tag::try_from(tag).ok()?;
569
570 let cf_id = match tag {
571 | Tag::Value | Tag::Deletion => 0,
572 | Tag::CfValue | Tag::CfDeletion => take_varint32(data)?,
573 };
574
575 let key = take_varstring(data)?;
576
577 if matches!(tag, Tag::Value | Tag::CfValue) {
578 take_varstring(data)?;
579 }
580
581 Some((cf_id, key))
582}
583
584/// Takes one length-prefixed byte string from the front of a batch record.
585///
586/// The returned slice borrows the original batch representation, and `data`
587/// advances past it. Invalid lengths or truncated input return `None`.
588fn take_varstring<'a>(data: &mut Key<'a>) -> Option<Key<'a>> {
589 let len = take_varint32(data)?.try_into().ok()?;
590
591 let (string, rest) = data.split_at_checked(len)?;
592 *data = rest;
593
594 Some(string)
595}
596
597/// Takes one RocksDB varint32 from the front of a batch record.
598///
599/// The parser consumes at most five bytes and advances `data` as bytes are
600/// read. A missing byte or a continuation bit on the fifth byte returns `None`.
601fn take_varint32(data: &mut &Slice) -> Option<u32> {
602 let mut result = 0_u32;
603
604 for shift in (0_u32..32).step_by(7) {
605 let (&byte, rest) = data.split_first()?;
606 *data = rest;
607 result |= u32::from(byte & 0x7F).checked_shl(shift)?;
608
609 if byte & 0x80 == 0 {
610 return Some(result);
611 }
612 }
613
614 None
615}
616
617/// Estimates write-batch capacity for a reusable sequence of raw pairs.
618///
619/// The estimate includes the fixed header, worst-case per-operation metadata,
620/// and payload lengths. Saturating arithmetic prevents an oversized input from
621/// wrapping the reservation.
622fn size_hint<'a, K, V, I>(items: I) -> usize
623where
624 I: Iterator<Item = (&'a K, &'a V)>,
625 K: AsRef<Slice> + 'a,
626 V: AsRef<Slice> + 'a,
627{
628 items.fold(HEADER, |capacity_bytes, (key, val)| {
629 capacity_bytes
630 .saturating_add(PER_OP)
631 .saturating_add(key.as_ref().len())
632 .saturating_add(val.as_ref().len())
633 })
634}