Skip to main content

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}