eidetica/transaction/mod.rs
1//! Transaction system for atomic database modifications
2//!
3//! This module provides the transaction API for making atomic changes to an Eidetica database.
4//! Transactions ensure that all changes within a transaction are applied atomically and maintain
5//! proper parent-child relationships in the Merkle-CRDT DAG structure.
6//!
7//! # Subtree Parent Management
8//!
9//! One of the critical responsibilities of the transaction system is establishing proper
10//! subtree parent relationships. When a store (subtree) is accessed for the first time
11//! in a transaction, the system must determine the correct parent entries for that subtree.
12//! This involves:
13//!
14//! 1. Checking for existing subtree tips (leaf nodes)
15//! 2. If no tips exist, traversing the DAG to find reachable subtree entries
16//! 3. Setting appropriate parent relationships (empty for first entry, or proper parents)
17
18pub mod errors;
19
20#[cfg(test)]
21mod tests;
22
23use std::{
24 collections::{BTreeMap, HashMap},
25 sync::{
26 Arc, Mutex,
27 atomic::{AtomicBool, Ordering},
28 },
29};
30
31pub use errors::TransactionError;
32use serde::{Deserialize, Serialize};
33
34use crate::{
35 Database, Result, Snapshot, Store,
36 auth::{
37 AuthSettings,
38 crypto::{PrivateKey, sign_entry},
39 types::{AuthInfo, SigKey},
40 validation::AuthValidator,
41 },
42 backend::{RecordMutations, RecordRange, RecordView, VerificationStatus},
43 constants::{INDEX, ROOT, SETTINGS},
44 crdt::{CRDT, Data, Doc, doc::Value},
45 entry::{Entry, EntryBuilder, ID},
46 height::HeightStrategy,
47 instance::WriteSource,
48 store::{ProjectionDescriptor, RecordProjection, Registry, SettingsStore, StoreError, state},
49};
50
51/// Creates a synthetic entry ID for multi-tip merged CRDT state caching.
52///
53/// Tips are sorted to ensure deterministic keys regardless of input order.
54/// The resulting ID has format `merge:{tip1}:{tip2}:...` which is distinct
55/// from real content-addressed entry IDs.
56fn create_merge_cache_id(tip_ids: &[ID]) -> ID {
57 let mut sorted_tips = tip_ids.to_vec();
58 sorted_tips.sort();
59
60 // Create a deterministic cache key by hashing the sorted tip IDs
61 let mut key = String::from("merge");
62 for tip in &sorted_tips {
63 key.push(':');
64 key.push_str(&tip.to_string());
65 }
66 ID::from_bytes(key.as_bytes())
67}
68
69/// Trait for encrypting/decrypting subtree data transparently
70///
71/// Encryptors are registered with a Transaction for specific subtrees, allowing
72/// transparent encryption/decryption at the transaction boundary. When an encryptor
73/// is registered:
74///
75/// - `get_full_state()` decrypts each historical entry before CRDT merging
76/// - `get_local_data()` returns plaintext (cached in EntryBuilder)
77/// - `update_subtree()` stores plaintext in cache, encrypted on commit
78///
79/// This ensures proper CRDT merge semantics while keeping data encrypted at rest.
80///
81/// # Wire Format
82///
83/// The trait operates on raw bytes, allowing implementations to define their own
84/// wire format. For example, AES-GCM implementations typically use `nonce || ciphertext`.
85/// Entry subtree payloads are opaque bytes, so ciphertext is stored verbatim with no
86/// additional encoding.
87///
88/// # Example
89///
90/// ```rust,ignore
91/// struct PasswordEncryptor { /* ... */ }
92///
93/// impl Encryptor for PasswordEncryptor {
94/// fn decrypt(&self, ciphertext: &[u8]) -> Result<Vec<u8>> {
95/// let (nonce, ct) = ciphertext.split_at(12);
96/// // decrypt with nonce and ciphertext...
97/// }
98///
99/// fn encrypt(&self, plaintext: &[u8]) -> Result<Vec<u8>> {
100/// let nonce = generate_nonce();
101/// let ct = encrypt(plaintext, &nonce);
102/// // return nonce || ciphertext
103/// }
104/// }
105/// ```
106pub(crate) trait Encryptor: Send + Sync {
107 /// Decrypt ciphertext bytes to plaintext bytes
108 ///
109 /// # Arguments
110 /// * `ciphertext` - Encrypted data in implementation-defined format
111 ///
112 /// # Returns
113 /// Plaintext bytes in whatever format the wrapped store produces (e.g.
114 /// JSON for `DocStore`/`Table`, binary Yrs updates for `YDoc`).
115 fn decrypt(&self, ciphertext: &[u8]) -> Result<Vec<u8>>;
116
117 /// Encrypt plaintext bytes to ciphertext bytes
118 ///
119 /// # Arguments
120 /// * `plaintext` - Bytes to encrypt; format is whatever the wrapped store
121 /// produces (JSON, binary CRDT update, etc.). The Encryptor itself does
122 /// not interpret the bytes.
123 ///
124 /// # Returns
125 /// Encrypted data in implementation-defined format
126 fn encrypt(&self, plaintext: &[u8]) -> Result<Vec<u8>>;
127
128 fn physical_record_key(&self, logical_key: &[u8]) -> Result<Vec<u8>> {
129 Ok(logical_key.to_vec())
130 }
131
132 fn encrypt_record(&self, _logical_key: &[u8], plaintext: &[u8]) -> Result<Vec<u8>> {
133 self.encrypt(plaintext)
134 }
135
136 fn decrypt_record(&self, physical_key: &[u8], ciphertext: &[u8]) -> Result<(Vec<u8>, Vec<u8>)> {
137 Ok((physical_key.to_vec(), self.decrypt(ciphertext)?))
138 }
139
140 fn projection_descriptor(&self, descriptor: ProjectionDescriptor) -> ProjectionDescriptor {
141 descriptor
142 }
143}
144
145/// Metadata structure for entries
146#[derive(Debug, Clone, Serialize, Deserialize)]
147pub(crate) struct EntryMetadata {
148 /// Snapshot of the `_settings` subtree at the time this entry was created.
149 /// This is the entry's **pin**: the exact `_settings` state its
150 /// signature must be validated against. Used for sync performance,
151 /// sparse-checkout validation, and deferred re-verification.
152 ///
153 /// Wire name remains `settings_tips` for on-disk stability — `Snapshot`
154 /// serializes as a bare ID array, identical to the legacy `Vec<ID>` shape.
155 #[serde(rename = "settings_tips")]
156 pub(crate) settings_snapshot: Snapshot,
157 /// Random entropy for ensuring unique IDs for root entries
158 pub(crate) entropy: Option<u64>,
159}
160
161/// Represents a single, atomic transaction for modifying a `Database`.
162///
163/// An `Transaction` encapsulates a mutable `EntryBuilder` being constructed. Users interact with
164/// specific `Store` instances obtained via `Transaction::get_store` to stage changes.
165/// All staged changes across different subtrees within the transaction are recorded
166/// in the internal `EntryBuilder`.
167///
168/// When `commit()` is called, the transaction:
169/// 1. Finalizes the `EntryBuilder` by building an immutable `Entry`
170/// 2. Calculates the entry's content-addressable ID
171/// 3. Ensures the correct parent links are set based on the tree's state
172/// 4. Removes any empty subtrees that didn't have data staged
173/// 5. Signs the entry if authentication is configured
174/// 6. Persists the resulting immutable `Entry` to the backend
175///
176/// `Transaction` instances are typically created via `Database::new_transaction()`.
177#[derive(Clone)]
178pub struct Transaction {
179 /// The entry builder being modified, wrapped in Option to support consuming on commit
180 entry_builder: Arc<Mutex<Option<EntryBuilder>>>,
181 /// The database this transaction belongs to
182 db: Database,
183 /// Provided signing key paired with its auth identity
184 provided_signing_key: Option<(PrivateKey, SigKey)>,
185 /// Registered encryptors for transparent encryption/decryption of specific subtrees
186 /// Maps subtree name -> encryptor implementation
187 /// When an encryptor is registered, the transaction automatically encrypts writes
188 /// and decrypts reads for that subtree
189 encryptors: Arc<Mutex<HashMap<String, Box<dyn Encryptor>>>>,
190 /// When true, `get_store` rejects any `_`-prefixed subtree name. Used by
191 /// `Database::create_with_init` to keep its init callback from opening the
192 /// system subtrees (`_settings`, `_root`, `_index`) that `create_with_init`
193 /// itself manages. Legitimate internal paths — `get_index()` (via
194 /// `Registry::new` → `DocStore::load`) and `Store::register`'s own
195 /// `_index` updates — bypass `get_store` and remain unaffected.
196 system_subtrees_locked: Arc<AtomicBool>,
197 record_mutations: Arc<Mutex<HashMap<String, RecordMutations>>>,
198 logical_record_mutations: Arc<Mutex<HashMap<String, RecordMutations>>>,
199 record_views: Arc<Mutex<HashMap<String, RecordView>>>,
200}
201
202/// RAII guard returned by [`Transaction::lock_system_subtrees`]. Releases the
203/// lock on drop, covering early-return-via-`?` and panic-unwind alike.
204pub(crate) struct SystemSubtreeLockGuard {
205 flag: Arc<AtomicBool>,
206}
207
208impl Drop for SystemSubtreeLockGuard {
209 fn drop(&mut self) {
210 self.flag.store(false, Ordering::Release);
211 }
212}
213
214impl Transaction {
215 /// Creates a new atomic transaction for a specific `Database` anchored at a snapshot.
216 ///
217 /// Initializes an internal `EntryBuilder` with its main parent pointers set to the
218 /// snapshot's tips instead of the database's current state. This allows creating
219 /// transactions that branch from specific points in the database history (e.g.
220 /// diamond patterns).
221 ///
222 /// # Arguments
223 /// * `database` - The `Database` this transaction will modify.
224 /// * `snapshot` - The snapshot to anchor the transaction at. Must contain at least one tip,
225 /// unless this transaction is creating the database's root entry.
226 pub(crate) async fn new_at(database: &Database, snapshot: &Snapshot) -> Result<Self> {
227 let tips = snapshot.tips();
228 // Validate that tips are not empty, unless we're creating the root entry
229 if tips.is_empty() {
230 // Check if this is a root entry creation by seeing if the database root exists in backend
231 let root_exists = database.ops().get(database.root_id()).await.is_ok();
232
233 if root_exists {
234 return Err(TransactionError::EmptyTipsNotAllowed.into());
235 }
236 // If root doesn't exist, this is valid (creating the root entry)
237 }
238
239 // Validate that all tips belong to the same tree
240 let backend = database.ops();
241 for tip_id in tips {
242 let entry = backend.get(tip_id).await?;
243 if !entry.in_tree(database.root_id()) {
244 return Err(TransactionError::InvalidTip {
245 tip_id: tip_id.clone(),
246 }
247 .into());
248 }
249 }
250
251 // Start with a basic entry linked to the database's root.
252 // Data and parents will be filled based on the transaction type.
253 let mut builder = Entry::builder(database.root_id().clone());
254
255 // Use the provided tips as parents (only if not empty)
256 if !tips.is_empty() {
257 builder.set_parents_mut(tips.to_vec());
258 }
259
260 Ok(Self {
261 entry_builder: Arc::new(Mutex::new(Some(builder))),
262 db: database.clone(),
263 provided_signing_key: None,
264 encryptors: Arc::new(Mutex::new(HashMap::new())),
265 system_subtrees_locked: Arc::new(AtomicBool::new(false)),
266 record_mutations: Arc::new(Mutex::new(HashMap::new())),
267 logical_record_mutations: Arc::new(Mutex::new(HashMap::new())),
268 record_views: Arc::new(Mutex::new(HashMap::new())),
269 })
270 }
271
272 /// Lock the system subtrees (`_settings`, `_root`, `_index`) for the
273 /// returned guard's lifetime.
274 ///
275 /// While the guard is live, [`Self::get_store`] rejects any subtree name
276 /// beginning with `_` (the system-subtree prefix). Used by
277 /// [`Database::create_with_init`] to guard its init callback against
278 /// clobbering `_settings`/`_root`/`_index`, all of which `create_with_init`
279 /// manages itself or via a dedicated accessor. The lock releases on the
280 /// guard's `Drop`, so it covers both early-return-via-`?` and panic-unwind
281 /// paths.
282 ///
283 /// The lock is process-local state on the (cloned-by-Arc) transaction;
284 /// clones share the same flag. Legitimate internal paths — `get_index()`
285 /// (via `Registry::new` → `DocStore::load`) and `Store::register`'s own
286 /// `_index` writes — bypass `get_store` and are unaffected.
287 pub(crate) fn lock_system_subtrees(&self) -> SystemSubtreeLockGuard {
288 self.system_subtrees_locked.store(true, Ordering::Release);
289 SystemSubtreeLockGuard {
290 flag: self.system_subtrees_locked.clone(),
291 }
292 }
293
294 /// Set signing key directly for user context (internal API).
295 ///
296 /// This method is used when a Database has a key attached
297 /// (via `Database::open().with_key()`). The provided SigningKey is already
298 /// decrypted and ready to use, eliminating the need for backend key lookup.
299 ///
300 /// # Arguments
301 /// * `signing_key` - The decrypted signing key from UserKeyManager
302 /// * `identity` - The SigKey identity used in database auth settings
303 pub(crate) fn set_provided_key(&mut self, signing_key: PrivateKey, identity: SigKey) {
304 self.provided_signing_key = Some((signing_key, identity));
305 }
306
307 /// Get current time as RFC3339 string.
308 ///
309 /// Delegates to the underlying instance's clock.
310 pub(crate) fn now_rfc3339(&self) -> Result<String> {
311 Ok(self.db.instance()?.clock().now_rfc3339())
312 }
313
314 /// Register an encryptor for transparent encryption/decryption of a specific subtree.
315 ///
316 /// Once registered, the transaction will automatically:
317 /// - Decrypt each historical entry before CRDT merging in `get_full_state()`
318 /// - Return plaintext data from `get_local_data()` (cached in EntryBuilder)
319 /// - Encrypt plaintext data before persisting in `commit()`
320 ///
321 /// This ensures proper CRDT merge semantics while keeping data encrypted at rest.
322 ///
323 /// # Arguments
324 /// * `subtree` - The name of the subtree to encrypt/decrypt
325 /// * `encryptor` - The encryptor implementation to use
326 ///
327 /// # Example
328 ///
329 /// For password-based encryption, use [`PasswordStore`] which handles
330 /// encryptor registration automatically:
331 ///
332 /// ```rust,ignore
333 /// let mut encrypted = tx.get_store::<PasswordStore<DocStore>>("secrets")?;
334 /// encrypted.initialize("my_password", Doc::new())?;
335 ///
336 /// // PasswordStore registers the encryptor internally
337 /// let docstore = encrypted.inner()?;
338 /// ```
339 ///
340 /// For custom encryption, implement the [`Encryptor`] trait:
341 ///
342 /// ```rust,ignore
343 /// struct MyEncryptor { /* ... */ }
344 /// impl Encryptor for MyEncryptor {
345 /// fn encrypt(&self, plaintext: &[u8]) -> Result<Vec<u8>> { /* ... */ }
346 /// fn decrypt(&self, ciphertext: &[u8]) -> Result<Vec<u8>> { /* ... */ }
347 /// }
348 ///
349 /// transaction.register_encryptor("secrets", Box::new(MyEncryptor::new()))?;
350 /// ```
351 ///
352 /// [`PasswordStore`]: crate::store::PasswordStore
353 /// [`Encryptor`]: crate::Encryptor
354 pub(crate) fn register_encryptor(
355 &self,
356 subtree: impl Into<String>,
357 encryptor: Box<dyn Encryptor>,
358 ) -> Result<()> {
359 self.encryptors
360 .lock()
361 .unwrap()
362 .insert(subtree.into(), encryptor);
363 Ok(())
364 }
365
366 fn physical_record_key(&self, store: &str, logical_key: &[u8]) -> Result<Vec<u8>> {
367 self.encryptors.lock().unwrap().get(store).map_or_else(
368 || Ok(logical_key.to_vec()),
369 |encryptor| encryptor.physical_record_key(logical_key),
370 )
371 }
372
373 fn encrypt_record(&self, store: &str, logical_key: &[u8], value: &[u8]) -> Result<Vec<u8>> {
374 self.encryptors.lock().unwrap().get(store).map_or_else(
375 || Ok(value.to_vec()),
376 |encryptor| encryptor.encrypt_record(logical_key, value),
377 )
378 }
379
380 fn decrypt_record(&self, store: &str, key: &[u8], value: &[u8]) -> Result<(Vec<u8>, Vec<u8>)> {
381 self.encryptors.lock().unwrap().get(store).map_or_else(
382 || Ok((key.to_vec(), value.to_vec())),
383 |encryptor| encryptor.decrypt_record(key, value),
384 )
385 }
386
387 pub(crate) fn database_id(&self) -> &ID {
388 self.db.root_id()
389 }
390
391 /// Decrypt bytes if an encryptor is registered, otherwise return them unchanged.
392 ///
393 /// This is used throughout Transaction to transparently decrypt encrypted payloads
394 /// before deserializing into CRDT types.
395 fn decrypt_if_needed(&self, subtree: &str, data: &[u8]) -> Result<Vec<u8>> {
396 if let Some(encryptor) = self.encryptors.lock().unwrap().get(subtree) {
397 encryptor.decrypt(data)
398 } else {
399 Ok(data.to_vec())
400 }
401 }
402
403 /// Encrypt bytes if an encryptor is registered for the subtree, otherwise return them unchanged.
404 fn encrypt_if_needed(&self, subtree: &str, plaintext: &[u8]) -> Result<Vec<u8>> {
405 if let Some(encryptor) = self.encryptors.lock().unwrap().get(subtree) {
406 encryptor.encrypt(plaintext)
407 } else {
408 Ok(plaintext.to_vec())
409 }
410 }
411
412 /// Get a SettingsStore handle for the settings subtree within this transaction.
413 ///
414 /// This method returns a `SettingsStore` that provides specialized access to the `_settings` subtree,
415 /// allowing you to read and modify settings data within this atomic transaction.
416 /// The DocStore automatically merges historical settings from the database with any
417 /// staged changes in this transaction.
418 ///
419 /// # Returns
420 ///
421 /// Returns a `Result<SettingsStore>` that can be used to:
422 /// - Read current settings values (including both historical and staged data)
423 /// - Stage new settings changes within this transaction
424 /// - Access nested settings structures
425 ///
426 /// # Example
427 ///
428 /// ```rust,no_run
429 /// # use eidetica::Database;
430 /// # async fn example(database: Database) -> eidetica::Result<()> {
431 /// let txn = database.new_transaction().await?;
432 /// let settings = txn.get_settings()?;
433 ///
434 /// // Read a setting
435 /// if let Ok(name) = settings.get_name().await {
436 /// println!("Database name: {}", name);
437 /// }
438 ///
439 /// // Modify a setting
440 /// settings.set_name("Updated Database Name").await?;
441 /// # Ok(())
442 /// # }
443 /// ```
444 ///
445 /// # Errors
446 ///
447 /// Returns an error if:
448 /// - Unable to create the SettingsStore for the settings subtree
449 /// - Operation has already been committed
450 pub fn get_settings(&self) -> Result<SettingsStore> {
451 // Create a SettingsStore for the settings subtree
452 SettingsStore::new(self)
453 }
454
455 /// Gets a handle to the Index for managing subtree registry and metadata.
456 ///
457 /// The Index provides access to the `_index` subtree, which stores metadata
458 /// about all subtrees in the database including their type identifiers and configurations.
459 ///
460 /// # Returns
461 ///
462 /// A `Result<Registry>` containing the handle for managing the index.
463 ///
464 /// # Errors
465 ///
466 /// Returns an error if:
467 /// - Unable to create the Registry for the _index subtree
468 /// - Operation has already been committed
469 pub async fn get_index(&self) -> Result<Registry> {
470 Registry::new(self, INDEX).await
471 }
472
473 /// Set the tree root field for the entry being built.
474 ///
475 /// Called by `Database::create()` to override the placeholder root with
476 /// `ID::default()`, making the entry a proper top-level root.
477 ///
478 /// # Arguments
479 /// * `root` - The tree root ID to set (use `ID::default()` for top-level roots)
480 pub(crate) fn set_entry_root(&self, root: ID) -> Result<()> {
481 let mut builder_ref = self.entry_builder.lock().unwrap();
482 let builder = builder_ref
483 .as_mut()
484 .ok_or(TransactionError::TransactionAlreadyCommitted)?;
485 builder.set_root_mut(root);
486 Ok(())
487 }
488
489 /// Set entropy in the entry metadata.
490 ///
491 /// This is used during database creation to ensure unique IDs for databases
492 /// even when they have identical settings.
493 ///
494 /// # Arguments
495 /// * `entropy` - Random entropy value
496 pub(crate) fn set_metadata_entropy(&self, entropy: u64) -> Result<()> {
497 let mut builder_ref = self.entry_builder.lock().unwrap();
498 let builder = builder_ref
499 .as_mut()
500 .ok_or(TransactionError::TransactionAlreadyCommitted)?;
501
502 // Parse existing metadata if present, or create new
503 let mut metadata = builder
504 .metadata()
505 .and_then(|m| serde_json::from_slice::<EntryMetadata>(m).ok())
506 .unwrap_or(EntryMetadata {
507 settings_snapshot: Snapshot::EMPTY,
508 entropy: None,
509 });
510
511 // Set entropy
512 metadata.entropy = Some(entropy);
513
514 // Serialize and set metadata
515 let metadata_json = serde_json::to_vec(&metadata)?;
516 builder.set_metadata_mut(metadata_json);
517
518 Ok(())
519 }
520
521 /// Stages an update for a specific subtree within this atomic transaction.
522 ///
523 /// This method is primarily intended for internal use by `Store` implementations
524 /// (like `DocStore::set`). It records the serialized `data` for the given `subtree`
525 /// name within the transaction's internal `EntryBuilder`.
526 ///
527 /// If this is the first modification to the named subtree within this transaction,
528 /// it also fetches and records the current tips of that subtree from the backend
529 /// to set the correct `subtree_parents` for the new entry.
530 ///
531 /// # Arguments
532 /// * `subtree` - The name of the subtree to update.
533 /// * `data` - The serialized CRDT data to stage for the subtree.
534 ///
535 /// # Returns
536 /// A `Result<()>` indicating success or an error.
537 pub(crate) async fn update_subtree(
538 &self,
539 subtree: impl AsRef<str>,
540 data: impl Into<Vec<u8>>,
541 ) -> Result<()> {
542 let subtree = subtree.as_ref();
543 let data = data.into();
544
545 // Check if we need to fetch tips (check without holding borrow across await)
546 let needs_tips = {
547 let builder_ref = self.entry_builder.lock().unwrap();
548 let builder = builder_ref
549 .as_ref()
550 .ok_or(TransactionError::TransactionAlreadyCommitted)?;
551 !builder.subtrees().contains(&subtree.to_string())
552 };
553
554 // Fetch tips if needed (no borrow held across this await)
555 let tips = if needs_tips {
556 let backend = self.db.ops();
557 // FIXME: we should get the subtree snapshot while still using the parent pointers
558 Some(
559 backend
560 .store_snapshot(self.db.root_id(), subtree)
561 .await?
562 .into_tips(),
563 )
564 } else {
565 None
566 };
567
568 // Now update the builder
569 let mut builder_ref = self.entry_builder.lock().unwrap();
570 let builder = builder_ref
571 .as_mut()
572 .ok_or(TransactionError::TransactionAlreadyCommitted)?;
573
574 builder.set_subtree_data_mut(subtree.to_string(), data);
575 if let Some(tips) = tips {
576 builder.set_subtree_parents_mut(subtree, tips);
577 }
578
579 Ok(())
580 }
581
582 pub(crate) fn stage_record(
583 &self,
584 store: &str,
585 projection: &dyn RecordProjection<Doc>,
586 key: Vec<u8>,
587 value: Option<Vec<u8>>,
588 ) -> Result<()> {
589 let Some(key) = projection.normalize_record_key(&key)? else {
590 return Ok(());
591 };
592 let conflicting_keys = {
593 let mut staged = self.logical_record_mutations.lock().unwrap();
594 let mutations = staged.entry(store.to_string()).or_default();
595 let conflicting_keys = mutations
596 .keys()
597 .filter(|staged_key| projection.staged_keys_conflict(staged_key, &key))
598 .cloned()
599 .collect::<Vec<_>>();
600 mutations.retain(|staged_key, _| !conflicting_keys.contains(staged_key));
601 mutations.insert(key.clone(), value.clone());
602 conflicting_keys
603 };
604 let physical_key = self.physical_record_key(store, &key)?;
605 let conflicting_physical_keys = conflicting_keys
606 .iter()
607 .map(|key| self.physical_record_key(store, key))
608 .collect::<Result<Vec<_>>>()?;
609 let value = value
610 .map(|value| self.encrypt_record(store, &key, &value))
611 .transpose()?;
612 let mut record_mutations = self.record_mutations.lock().unwrap();
613 let mutations = record_mutations.entry(store.to_string()).or_default();
614 mutations.retain(|staged_key, _| !conflicting_physical_keys.contains(staged_key));
615 mutations.insert(physical_key, value);
616 Ok(())
617 }
618
619 pub(crate) async fn record_get(
620 &self,
621 store: &str,
622 projection: &dyn RecordProjection<Doc>,
623 key: &[u8],
624 ) -> Result<Option<Vec<u8>>> {
625 let Some(key) = projection.normalize_record_key(key)? else {
626 return Ok(None);
627 };
628 if let Some(mutations) = self.logical_record_mutations.lock().unwrap().get(store) {
629 if let Some(value) = mutations.get(&key) {
630 return Ok(value.clone());
631 }
632 if mutations
633 .keys()
634 .any(|staged_key| projection.staged_key_shadows_cached(staged_key, &key))
635 {
636 return Ok(None);
637 }
638 }
639 let physical_key = self.physical_record_key(store, &key)?;
640 let view = match self.record_view(store, projection).await {
641 Err(err) if err.is_unsupported_store_state() => {
642 return self.record_get_from_history(store, &key).await;
643 }
644 result => result?,
645 };
646 match self
647 .db
648 .ops()
649 .store_state_record_get(&view, &physical_key)
650 .await
651 {
652 Err(err) if err.is_unsupported_store_state() => {
653 self.record_get_from_history(store, &key).await
654 }
655 Err(err) if err.is_invalid_store_state_view() => {
656 self.record_views.lock().unwrap().remove(store);
657 let view = match self.record_view(store, projection).await {
658 Err(err) if err.is_unsupported_store_state() => {
659 return self.record_get_from_history(store, &key).await;
660 }
661 result => result?,
662 };
663 match self
664 .db
665 .ops()
666 .store_state_record_get(&view, &physical_key)
667 .await
668 {
669 Err(err) if err.is_unsupported_store_state() => {
670 self.record_get_from_history(store, &key).await
671 }
672 result => result?.map_or(Ok(None), |value| {
673 self.decrypt_record(store, &physical_key, &value)
674 .map(|(_, value)| Some(value))
675 }),
676 }
677 }
678 result => result?.map_or(Ok(None), |value| {
679 self.decrypt_record(store, &physical_key, &value)
680 .map(|(_, value)| Some(value))
681 }),
682 }
683 }
684
685 pub(crate) fn record_has_staged_descendant(
686 &self,
687 store: &str,
688 projection: &dyn RecordProjection<Doc>,
689 key: &[u8],
690 ) -> Result<bool> {
691 let Some(key) = projection.normalize_record_key(key)? else {
692 return Ok(false);
693 };
694 Ok(self
695 .logical_record_mutations
696 .lock()
697 .unwrap()
698 .get(store)
699 .is_some_and(|mutations| {
700 mutations.iter().any(|(staged_key, value)| {
701 value.is_some() && projection.staged_key_descends_from(staged_key, &key)
702 })
703 }))
704 }
705
706 pub(crate) async fn record_scan(
707 &self,
708 store: &str,
709 projection: &dyn RecordProjection<Doc>,
710 after: Option<&[u8]>,
711 limit: usize,
712 ) -> Result<crate::backend::RecordPage> {
713 if limit == 0 {
714 return Ok(crate::backend::RecordPage::default());
715 }
716 let view = match self.record_view(store, projection).await {
717 Err(err) if err.is_unsupported_store_state() => {
718 return self
719 .record_scan_from_history(store, projection, after, limit)
720 .await;
721 }
722 result => result?,
723 };
724 let logical_mutations = self
725 .logical_record_mutations
726 .lock()
727 .unwrap()
728 .get(store)
729 .cloned()
730 .unwrap_or_default();
731 let mutations = self
732 .record_mutations
733 .lock()
734 .unwrap()
735 .get(store)
736 .cloned()
737 .unwrap_or_default();
738 let mut merged = mutations
739 .iter()
740 .filter(|(key, value)| {
741 value.is_some() && after.is_none_or(|after| key.as_slice() > after)
742 })
743 .map(|(key, value)| (key.clone(), value.clone().unwrap()))
744 .collect::<BTreeMap<_, _>>();
745 let mut backend_after = after.map(ToOwned::to_owned);
746 let mut backend_has_more = true;
747 while backend_has_more {
748 let page = match self
749 .db
750 .ops()
751 .store_state_record_scan(
752 &view,
753 &RecordRange::default(),
754 backend_after.as_deref(),
755 limit.max(1),
756 )
757 .await
758 {
759 Err(err) if err.is_unsupported_store_state() => {
760 return self
761 .record_scan_from_history(store, projection, after, limit)
762 .await;
763 }
764 Err(err) if err.is_invalid_store_state_view() => {
765 self.record_views.lock().unwrap().remove(store);
766 match self.record_view(store, projection).await {
767 Err(err) if err.is_unsupported_store_state() => {
768 return self
769 .record_scan_from_history(store, projection, after, limit)
770 .await;
771 }
772 Ok(view) => match self
773 .db
774 .ops()
775 .store_state_record_scan(
776 &view,
777 &RecordRange::default(),
778 backend_after.as_deref(),
779 limit.max(1),
780 )
781 .await
782 {
783 Err(err) if err.is_unsupported_store_state() => {
784 return self
785 .record_scan_from_history(store, projection, after, limit)
786 .await;
787 }
788 result => result?,
789 },
790 Err(err) => return Err(err),
791 }
792 }
793 result => result?,
794 };
795 for (physical_key, value) in page.records {
796 let (logical_key, _) = self.decrypt_record(store, &physical_key, &value)?;
797 if !logical_mutations.keys().any(|staged_key| {
798 projection.staged_key_shadows_cached(staged_key, &logical_key)
799 }) {
800 merged.insert(physical_key, value);
801 }
802 }
803 backend_after = page.next;
804 backend_has_more = backend_after.is_some();
805 let backend_reached_page_end = merged.keys().nth(limit).is_some_and(|page_end| {
806 backend_after
807 .as_ref()
808 .is_some_and(|backend_after| backend_after >= page_end)
809 });
810 if backend_reached_page_end {
811 break;
812 }
813 }
814 let mut records = merged
815 .into_iter()
816 .take(limit.saturating_add(1))
817 .collect::<Vec<_>>();
818 let has_more = records.len() > limit || backend_has_more;
819 records.truncate(limit);
820 let next = has_more.then(|| records.last().unwrap().0.clone());
821 for (key, value) in &mut records {
822 (*key, *value) = self.decrypt_record(store, key, value)?;
823 }
824 Ok(crate::backend::RecordPage { records, next })
825 }
826
827 async fn record_get_from_history(&self, store: &str, key: &[u8]) -> Result<Option<Vec<u8>>> {
828 let state: Doc = self.get_full_state(store).await?;
829 Ok(state
830 .get(std::str::from_utf8(key).map_err(|error| {
831 TransactionError::StoreDeserializationFailed {
832 store: store.to_string(),
833 reason: error.to_string(),
834 }
835 })?)
836 .and_then(Value::as_text)
837 .map(|value| value.as_bytes().to_vec()))
838 }
839
840 async fn record_scan_from_history(
841 &self,
842 store: &str,
843 projection: &dyn RecordProjection<Doc>,
844 after: Option<&[u8]>,
845 limit: usize,
846 ) -> Result<crate::backend::RecordPage> {
847 let state: Doc = self.get_full_state(store).await?;
848 let mut projected = RecordMutations::new();
849 projection.project_delta(&state, &mut projected)?;
850 let mutations = self
851 .logical_record_mutations
852 .lock()
853 .unwrap()
854 .get(store)
855 .cloned()
856 .unwrap_or_default();
857 let mut records = projected
858 .into_iter()
859 .filter_map(|(key, value)| value.map(|value| (key, value)))
860 .filter(|(cached_key, _)| {
861 !mutations
862 .keys()
863 .any(|staged_key| projection.staged_key_shadows_cached(staged_key, cached_key))
864 })
865 .collect::<BTreeMap<_, _>>();
866 for (key, value) in mutations {
867 match value {
868 Some(value) => {
869 records.insert(key, value);
870 }
871 None => {
872 records.remove(&key);
873 }
874 }
875 }
876 let mut records = records
877 .into_iter()
878 .map(|(logical_key, value)| {
879 Ok((
880 self.physical_record_key(store, &logical_key)?,
881 (logical_key, value),
882 ))
883 })
884 .collect::<Result<BTreeMap<_, _>>>()?
885 .into_iter()
886 .filter(|(physical_key, _)| after.is_none_or(|after| physical_key.as_slice() > after))
887 .take(limit.saturating_add(1))
888 .collect::<Vec<_>>();
889 let has_more = records.len() > limit;
890 records.truncate(limit);
891 let next = has_more.then(|| records.last().unwrap().0.clone());
892 let records = records
893 .into_iter()
894 .map(|(_, (logical_key, value))| (logical_key, value))
895 .collect();
896 Ok(crate::backend::RecordPage { records, next })
897 }
898
899 fn encrypted_projection_descriptor(
900 &self,
901 store: &str,
902 descriptor: ProjectionDescriptor,
903 ) -> ProjectionDescriptor {
904 self.encryptors
905 .lock()
906 .unwrap()
907 .get(store)
908 .map_or(descriptor.clone(), |encryptor| {
909 encryptor.projection_descriptor(descriptor)
910 })
911 }
912
913 async fn publish_record_view(
914 &self,
915 store: &str,
916 projection: &dyn RecordProjection<Doc>,
917 request: crate::backend::StoreStateRequest,
918 entries: &[Entry],
919 ) -> Result<RecordView> {
920 if !self.encryptors.lock().unwrap().contains_key(store) {
921 return state::publish_records(
922 self.db.ops(),
923 request,
924 entries
925 .iter()
926 .filter_map(|entry| entry.data(store).ok().map(|bytes| bytes.as_slice())),
927 projection,
928 )
929 .await;
930 }
931
932 let token = self.db.ops().begin_store_state_staging(request).await?;
933 let result = async {
934 let mut logical = RecordMutations::new();
935 for entry in entries {
936 if let Ok(bytes) = entry.data(store) {
937 let plaintext = self.decrypt_if_needed(store, bytes)?;
938 let delta: Doc = serde_json::from_slice(&plaintext)?;
939 projection.project_delta(&delta, &mut logical)?;
940 }
941 }
942 logical.retain(|_, value| value.is_some());
943 let mut physical = RecordMutations::new();
944 for (key, value) in logical {
945 if let Some(value) = value {
946 physical.insert(
947 self.physical_record_key(store, &key)?,
948 Some(self.encrypt_record(store, &key, &value)?),
949 );
950 }
951 }
952 if !physical.is_empty() {
953 self.db
954 .ops()
955 .stage_store_state_records(&token, physical)
956 .await?;
957 }
958 self.db.ops().publish_store_state(token.clone()).await
959 }
960 .await;
961 if result.is_err() {
962 let _ = self.db.ops().abort_store_state(token).await;
963 }
964 result
965 }
966
967 async fn record_view(
968 &self,
969 store: &str,
970 projection: &dyn RecordProjection<Doc>,
971 ) -> Result<RecordView> {
972 if let Some(view) = self.record_views.lock().unwrap().get(store).cloned() {
973 return Ok(view);
974 }
975 self.init_subtree_parents(store).await?;
976 let parents = self
977 .entry_builder
978 .lock()
979 .unwrap()
980 .as_ref()
981 .ok_or(TransactionError::TransactionAlreadyCommitted)?
982 .subtree_parents(store)
983 .unwrap_or_default();
984 let source_key = create_merge_cache_id(&parents).to_string().into_bytes();
985 let descriptor = self.encrypted_projection_descriptor(store, projection.descriptor());
986 let request = state::records_request(
987 self.db.root_id(),
988 store,
989 descriptor,
990 source_key,
991 crate::backend::CacheScope::Shared,
992 );
993 let view = if let Some(view) = self.db.ops().resolve_store_state(&request).await? {
994 view
995 } else {
996 let boundary = Snapshot::from(parents);
997 let entries = self
998 .db
999 .ops()
1000 .store_at(self.db.root_id(), store, &boundary)
1001 .await?;
1002 self.publish_record_view(store, projection, request, &entries)
1003 .await?
1004 };
1005 self.record_views
1006 .lock()
1007 .unwrap()
1008 .insert(store.to_string(), view.clone());
1009 Ok(view)
1010 }
1011
1012 /// Gets a handle to a specific `Store` for modification within this transaction.
1013 ///
1014 /// This method creates and returns an instance of the specified `Store` type `T`,
1015 /// associated with this `Transaction`. The returned `Store` handle can be used to
1016 /// stage changes (e.g., using `DocStore::set`).
1017 /// These changes are recorded within this `Transaction`.
1018 ///
1019 /// If this is the first time this subtree is accessed within the transaction,
1020 /// its parent tips will be fetched and stored.
1021 ///
1022 /// # Type Parameters
1023 /// * `T` - The concrete `Store` implementation type to create.
1024 ///
1025 /// # Arguments
1026 /// * `subtree_name` - The name of the subtree to get a modification handle for.
1027 ///
1028 /// # Returns
1029 /// A `Result<T>` containing the `Store` handle.
1030 pub async fn get_store<T>(&self, subtree_name: impl Into<String> + Send) -> Result<T>
1031 where
1032 T: Store + Send,
1033 {
1034 let subtree_name = subtree_name.into();
1035
1036 // Skip special system subtrees to avoid circular dependencies
1037 let is_system_subtree =
1038 subtree_name == INDEX || subtree_name == SETTINGS || subtree_name == ROOT;
1039
1040 if is_system_subtree && self.system_subtrees_locked.load(Ordering::Acquire) {
1041 return Err(TransactionError::SystemSubtreeLocked { name: subtree_name }.into());
1042 }
1043
1044 // Initialize subtree parents before checking _index
1045 self.init_subtree_parents(&subtree_name).await?;
1046
1047 if is_system_subtree {
1048 // System subtrees don't use _index registration
1049 return T::load(self, subtree_name).await;
1050 }
1051
1052 // Check _index to determine if this is a new or existing subtree
1053 let index_store = self.get_index().await?;
1054 if index_store.contains(&subtree_name).await {
1055 // Type validation for existing subtree
1056 let subtree_info = index_store.get_entry(&subtree_name).await?;
1057
1058 if !T::supports_type_id(&subtree_info.type_id) {
1059 return Err(StoreError::TypeMismatch {
1060 store: subtree_name,
1061 expected: T::type_id().to_string(),
1062 actual: subtree_info.type_id,
1063 }
1064 .into());
1065 }
1066
1067 // Type supported - create the Store
1068 T::load(self, subtree_name).await
1069 } else {
1070 // New subtree - register adds it to _index
1071 T::register(self, subtree_name).await
1072 }
1073 }
1074
1075 /// Get the subtree tips reachable from the given main tree entries.
1076 async fn get_subtree_tips(&self, subtree_name: &str, main_parents: &[ID]) -> Result<Vec<ID>> {
1077 let boundary = Snapshot::from(main_parents.to_vec());
1078 self.db
1079 .ops()
1080 .store_snapshot_at(self.db.root_id(), subtree_name, &boundary)
1081 .await
1082 .map(Snapshot::into_tips)
1083 }
1084
1085 /// Initialize subtree parents if this is the first time accessing this subtree
1086 /// in this transaction.
1087 pub(crate) async fn init_subtree_parents(&self, subtree_name: &str) -> Result<()> {
1088 let main_parents = {
1089 let builder_ref = self.entry_builder.lock().unwrap();
1090 let builder = builder_ref
1091 .as_ref()
1092 .ok_or(TransactionError::TransactionAlreadyCommitted)?;
1093
1094 let subtrees = builder.subtrees();
1095 if subtrees.contains(&subtree_name.to_string()) {
1096 return Ok(()); // Already initialized
1097 }
1098 builder.parents().unwrap_or_default()
1099 };
1100
1101 let tips = self.get_subtree_tips(subtree_name, &main_parents).await?;
1102
1103 let mut builder_ref = self.entry_builder.lock().unwrap();
1104 let builder = builder_ref
1105 .as_mut()
1106 .ok_or(TransactionError::TransactionAlreadyCommitted)?;
1107
1108 // Initialize the subtree with proper parent relationships
1109 // set_subtree_parents_mut creates the subtree with data=None if it doesn't exist
1110 builder.set_subtree_parents_mut(subtree_name, tips);
1111
1112 Ok(())
1113 }
1114
1115 /// Gets the currently staged data for a specific subtree within this transaction.
1116 ///
1117 /// This is intended for use by `Store` implementations to retrieve the data
1118 /// they have staged locally within the `Transaction` before potentially merging
1119 /// it with historical data.
1120 ///
1121 /// # Type Parameters
1122 /// * `T` - The data type (expected to be a CRDT) to deserialize the staged data into.
1123 ///
1124 /// # Arguments
1125 /// * `subtree_name` - The name of the subtree whose staged data is needed.
1126 ///
1127 /// # Returns
1128 /// A `Result<Option<T>>`:
1129 ///
1130 /// # Behavior
1131 /// - If the subtree doesn't exist or has no data, returns `Ok(None)`
1132 /// - If the subtree exists but has empty data (empty string or whitespace), returns `Ok(None)`
1133 /// - Otherwise deserializes the JSON data to type `T` and returns `Ok(Some(T))`
1134 ///
1135 /// # Errors
1136 /// Returns an error if the transaction has already been committed or if the
1137 /// subtree data exists but cannot be deserialized to type `T`.
1138 pub fn get_local_data<T>(&self, subtree_name: impl AsRef<str>) -> Result<Option<T>>
1139 where
1140 T: Data,
1141 {
1142 let subtree_name = subtree_name.as_ref();
1143 let builder_ref = self.entry_builder.lock().unwrap();
1144 let builder = builder_ref
1145 .as_ref()
1146 .ok_or(TransactionError::TransactionAlreadyCommitted)?;
1147
1148 if let Ok(data) = builder.data(subtree_name) {
1149 if data.is_empty() {
1150 Ok(None)
1151 } else {
1152 serde_json::from_slice(data).map(Some).map_err(|e| {
1153 TransactionError::StoreDeserializationFailed {
1154 store: subtree_name.to_string(),
1155 reason: e.to_string(),
1156 }
1157 .into()
1158 })
1159 }
1160 } else {
1161 Ok(None)
1162 }
1163 }
1164
1165 /// Gets the fully merged historical state of a subtree up to the point this transaction began.
1166 ///
1167 /// This retrieves all relevant historical entries for the `subtree_name` from the backend,
1168 /// considering the parent tips recorded when this `Transaction` was created (or when the
1169 /// subtree was first accessed within the transaction). It deserializes the data from each
1170 /// relevant entry into the CRDT type `T` and merges them according to `T`'s `CRDT::merge`
1171 /// implementation.
1172 ///
1173 /// This is intended for use by `Store` implementations (e.g., in their `get` or `get_all` methods)
1174 /// to provide the historical context against which staged changes might be applied or compared.
1175 ///
1176 /// # Type Parameters
1177 /// * `T` - The CRDT type to deserialize and merge the historical subtree data into.
1178 ///
1179 /// # Arguments
1180 /// * `subtree_name` - The name of the subtree.
1181 ///
1182 /// # Returns
1183 /// A `Result<T>` containing the merged historical data of type `T`. Returns `Ok(T::default())`
1184 /// if the subtree has no history prior to this transaction.
1185 pub(crate) async fn get_full_state<T>(&self, subtree_name: impl AsRef<str> + Send) -> Result<T>
1186 where
1187 T: CRDT + Send,
1188 {
1189 self.get_full_state_with_descriptor::<T>(
1190 subtree_name,
1191 ProjectionDescriptor {
1192 name: "eidetica/opaque".to_string(),
1193 version: 0,
1194 },
1195 )
1196 .await
1197 }
1198
1199 async fn get_full_state_with_descriptor<T>(
1200 &self,
1201 subtree_name: impl AsRef<str> + Send,
1202 descriptor: ProjectionDescriptor,
1203 ) -> Result<T>
1204 where
1205 T: CRDT + Send,
1206 {
1207 let subtree_name = subtree_name.as_ref();
1208
1209 // Check if we need to initialize subtree tips (get data from RefCell before await)
1210 let (needs_init, main_parents) = {
1211 let builder_ref = self.entry_builder.lock().unwrap();
1212 let builder = builder_ref
1213 .as_ref()
1214 .ok_or(TransactionError::TransactionAlreadyCommitted)?;
1215
1216 let subtrees = builder.subtrees();
1217 if subtrees.contains(&subtree_name.to_string()) {
1218 (false, Vec::new())
1219 } else {
1220 (true, builder.parents().unwrap_or_default())
1221 }
1222 };
1223
1224 // Initialize subtree tips if needed (async operations)
1225 if needs_init {
1226 let current_database_snapshot = self.db.ops().snapshot(self.db.root_id()).await?;
1227
1228 // Set-equal comparison via Snapshot canonical form.
1229 let parents_snapshot = Snapshot::from(main_parents.clone());
1230 let tips = if parents_snapshot == current_database_snapshot {
1231 let backend = self.db.ops();
1232 backend
1233 .store_snapshot(self.db.root_id(), subtree_name)
1234 .await?
1235 .into_tips()
1236 } else {
1237 // This transaction uses custom tips - use special handler
1238 self.db
1239 .ops()
1240 .store_snapshot_at(self.db.root_id(), subtree_name, &parents_snapshot)
1241 .await?
1242 .into_tips()
1243 };
1244
1245 // Update RefCell after async operations
1246 let mut builder_ref = self.entry_builder.lock().unwrap();
1247 let builder = builder_ref
1248 .as_mut()
1249 .ok_or(TransactionError::TransactionAlreadyCommitted)?;
1250 builder.set_subtree_parents_mut(subtree_name, tips);
1251 }
1252
1253 // Get the parent pointers for this subtree
1254 let parents = {
1255 let builder_ref = self.entry_builder.lock().unwrap();
1256 let builder = builder_ref
1257 .as_ref()
1258 .ok_or(TransactionError::TransactionAlreadyCommitted)?;
1259 builder.subtree_parents(subtree_name).unwrap_or_default()
1260 };
1261
1262 // If there are no parents, return a default
1263 if parents.is_empty() {
1264 return Ok(T::default());
1265 }
1266
1267 // Compute the CRDT state using merge-base ROOT-to-target computation
1268 self.compute_subtree_state_merge_based(subtree_name, &parents, &descriptor)
1269 .await
1270 }
1271
1272 /// Computes the CRDT state for a subtree using correct recursive merge-base algorithm.
1273 ///
1274 /// Algorithm:
1275 /// 1. If no entries, return default state
1276 /// 2. If single entry, compute its state recursively
1277 /// 3. If multiple entries, find their merge base and compute state from there
1278 ///
1279 /// # Type Parameters
1280 /// * `T` - The CRDT type to compute the state for
1281 ///
1282 /// # Arguments
1283 /// * `subtree_name` - The name of the subtree
1284 /// * `entry_ids` - The entry IDs to compute the merged state for (tips)
1285 ///
1286 /// # Returns
1287 /// A `Result<T>` containing the computed CRDT state
1288 async fn compute_subtree_state_merge_based<T>(
1289 &self,
1290 subtree_name: impl AsRef<str> + Send,
1291 entry_ids: &[ID],
1292 descriptor: &ProjectionDescriptor,
1293 ) -> Result<T>
1294 where
1295 T: CRDT + Send,
1296 {
1297 // Base case: no entries
1298 if entry_ids.is_empty() {
1299 return Ok(T::default());
1300 }
1301
1302 let subtree_name = subtree_name.as_ref();
1303
1304 // If we have a single entry, compute its state recursively
1305 if entry_ids.len() == 1 {
1306 return self
1307 .compute_single_entry_state_recursive(subtree_name, &entry_ids[0], descriptor)
1308 .await;
1309 }
1310
1311 // Multiple entries: check multi-tip cache first
1312 let cache_id = create_merge_cache_id(entry_ids);
1313
1314 let cache_request = state::opaque_request(
1315 self.db.root_id(),
1316 subtree_name,
1317 descriptor.clone(),
1318 cache_id.to_string().into_bytes(),
1319 crate::backend::CacheScope::Shared,
1320 );
1321 if let Some(bytes) = state::load_cached(self.db.ops(), &cache_request).await? {
1322 let decrypted = self.decrypt_if_needed(subtree_name, &bytes)?;
1323 let result: T = serde_json::from_slice(&decrypted)?;
1324 return Ok(result);
1325 }
1326
1327 // Cache miss: resolve the merge base and the path to fold in a
1328 // single call, so both come from one view of the store.
1329 let merge = self
1330 .db
1331 .ops()
1332 .compute_merge_state(self.db.root_id(), subtree_name, entry_ids)
1333 .await?;
1334
1335 let result = match &merge.merge_base {
1336 Some(base) => {
1337 // Compute the base state recursively, then fold the path
1338 // entries (deduplicated, height/ID sorted) on top of it.
1339 let state: T = self
1340 .compute_single_entry_state_recursive(subtree_name, base, descriptor)
1341 .await?;
1342 self.merge_path_entries(subtree_name, state, &merge.path)
1343 .await?
1344 }
1345 // With no merge base the histories are disjoint — a store
1346 // created independently on both sides of a fork — so fold the
1347 // full ancestry from a default state. Convergence must not
1348 // depend on which peer created the store first. `store_at`
1349 // batch-fetches the whole entries in fold order: one query
1350 // instead of a per-ID fetch of the entire history.
1351 None => {
1352 let boundary = Snapshot::from(entry_ids.to_vec());
1353 let entries = self
1354 .db
1355 .ops()
1356 .store_at(self.db.root_id(), subtree_name, &boundary)
1357 .await?;
1358 self.fold_store_entries(subtree_name, &entries)?
1359 }
1360 };
1361
1362 // Cache the computed merge result
1363 let bytes = self.encrypt_if_needed(subtree_name, &serde_json::to_vec(&result)?)?;
1364 state::store_cached(self.db.ops(), cache_request, bytes).await?;
1365
1366 Ok(result)
1367 }
1368
1369 /// Computes the CRDT state for a single entry using batch fetching.
1370 ///
1371 /// Algorithm:
1372 /// 1. Check if entry state is cached → return it
1373 /// 2. Fetch all ancestors in one batch query (sorted by height)
1374 /// 3. Merge all entries in order from root to target
1375 /// 4. Cache only the final result
1376 ///
1377 /// # Type Parameters
1378 /// * `T` - The CRDT type to compute the state for
1379 ///
1380 /// # Arguments
1381 /// * `subtree_name` - The name of the subtree
1382 /// * `entry_id` - The entry ID to compute the state for
1383 ///
1384 /// # Returns
1385 /// A `Result<T>` containing the computed CRDT state for the entry
1386 fn compute_single_entry_state_recursive<'a, T>(
1387 &'a self,
1388 subtree_name: &'a str,
1389 entry_id: &'a ID,
1390 descriptor: &'a ProjectionDescriptor,
1391 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<T>> + Send + 'a>>
1392 where
1393 T: CRDT + Send + 'a,
1394 {
1395 Box::pin(async move {
1396 let request = state::opaque_request(
1397 self.db.root_id(),
1398 subtree_name,
1399 descriptor.clone(),
1400 entry_id.to_string().into_bytes(),
1401 crate::backend::CacheScope::Shared,
1402 );
1403 if let Some(bytes) = state::load_cached(self.db.ops(), &request).await? {
1404 let decrypted = self.decrypt_if_needed(subtree_name, &bytes)?;
1405 let result: T = serde_json::from_slice(&decrypted)?;
1406 return Ok(result);
1407 }
1408
1409 // Step 2: Batch fetch all ancestors sorted by height (root first)
1410 // This single query replaces N recursive queries
1411 let boundary = Snapshot::from([entry_id.clone()]);
1412 let entries = self
1413 .db
1414 .ops()
1415 .store_at(self.db.root_id(), subtree_name, &boundary)
1416 .await?;
1417
1418 // Step 3: Merge all entries in order (already sorted by height, root first)
1419 let result: T = self.fold_store_entries(subtree_name, &entries)?;
1420
1421 // Step 4: Cache only the final result (encrypted if encryptor is registered)
1422 let bytes = self.encrypt_if_needed(subtree_name, &serde_json::to_vec(&result)?)?;
1423 state::store_cached(self.db.ops(), request, bytes).await?;
1424
1425 Ok(result)
1426 })
1427 }
1428
1429 /// Folds already-fetched entries into a CRDT state from the default.
1430 ///
1431 /// The entries must be in fold order (height then ID, root first), as
1432 /// `store_at` returns them.
1433 fn fold_store_entries<T>(&self, subtree_name: &str, entries: &[Entry]) -> Result<T>
1434 where
1435 T: CRDT,
1436 {
1437 let mut result = T::default();
1438 for entry in entries {
1439 let local_data = if let Ok(data) = entry.data(subtree_name) {
1440 // Decrypt before deserializing
1441 let plaintext = self.decrypt_if_needed(subtree_name, data)?;
1442 serde_json::from_slice::<T>(&plaintext)?
1443 } else {
1444 T::default()
1445 };
1446 result = result.merge(&local_data)?;
1447 }
1448 Ok(result)
1449 }
1450
1451 /// Merges a sequence of entries into a CRDT state.
1452 ///
1453 /// # Arguments
1454 /// * `subtree_name` - The name of the subtree
1455 /// * `initial_state` - The initial CRDT state to merge into
1456 /// * `entry_ids` - The entry IDs to merge in order
1457 ///
1458 /// # Returns
1459 /// A `Result<T>` containing the merged CRDT state
1460 async fn merge_path_entries<T>(
1461 &self,
1462 subtree_name: &str,
1463 mut state: T,
1464 entry_ids: &[ID],
1465 ) -> Result<T>
1466 where
1467 T: CRDT,
1468 {
1469 for entry_id in entry_ids {
1470 let entry = self.db.ops().get(entry_id).await?;
1471
1472 // Get local data for this entry in the subtree
1473 let local_data = if let Ok(data) = entry.data(subtree_name) {
1474 // Decrypt before deserializing
1475 let plaintext = self.decrypt_if_needed(subtree_name, data)?;
1476 serde_json::from_slice::<T>(&plaintext)?
1477 } else {
1478 T::default()
1479 };
1480
1481 state = state.merge(&local_data)?;
1482 }
1483
1484 Ok(state)
1485 }
1486
1487 /// Commits the transaction, finalizing and persisting the entry to the backend.
1488 ///
1489 /// This method:
1490 /// 1. Takes ownership of the `EntryBuilder` from the internal `Option`
1491 /// 2. Removes any empty subtrees
1492 /// 3. Adds metadata if appropriate
1493 /// 4. Sets authentication if configured
1494 /// 5. Builds the immutable `Entry` using `EntryBuilder::build()`
1495 /// 6. Signs the entry if authentication is configured
1496 /// 7. Validates authentication if present
1497 /// 8. Calculates the entry's content-addressable ID
1498 /// 9. Persists the entry to the backend
1499 /// 10. Returns the ID of the newly created entry
1500 ///
1501 /// After commit, the transaction cannot be used again, as the internal
1502 /// `EntryBuilder` has been consumed.
1503 ///
1504 /// # Returns
1505 /// A `Result<ID>` containing the ID of the committed entry.
1506 pub async fn commit(self) -> Result<ID> {
1507 self.commit_inner(None).await
1508 }
1509
1510 /// Commit while the caller holds this database's per-tree write lock.
1511 ///
1512 /// This is used by operations whose decision depends on a read made before
1513 /// the write. Ordinary transactions acquire the same lock when persisting.
1514 pub(crate) async fn commit_under_tree_lock(
1515 self,
1516 guard: tokio::sync::OwnedMutexGuard<()>,
1517 ) -> Result<ID> {
1518 self.commit_inner(Some(guard)).await
1519 }
1520
1521 async fn commit_inner(self, guard: Option<tokio::sync::OwnedMutexGuard<()>>) -> Result<ID> {
1522 {
1523 let staged = self.logical_record_mutations.lock().unwrap().clone();
1524 for (store, mutations) in staged {
1525 let delta = crate::store::table::encode_entry_delta(&mutations)?;
1526 self.update_subtree(store, serde_json::to_vec(&delta)?)
1527 .await?;
1528 }
1529 }
1530
1531 // Check if this is a settings subtree update and get the effective settings before any borrowing
1532 let has_settings_update = {
1533 let builder_cell = self.entry_builder.lock().unwrap();
1534 let builder = builder_cell
1535 .as_ref()
1536 .ok_or(TransactionError::TransactionAlreadyCommitted)?;
1537 builder.subtrees().contains(&SETTINGS.to_string())
1538 };
1539
1540 // Get settings using full CRDT state computation
1541 let historical_settings = self.get_full_state::<Doc>(SETTINGS).await?;
1542
1543 // However, if this is a settings update and there's no historical auth but staged auth exists,
1544 // use the staged settings for validation (this handles initial database creation with auth)
1545 let effective_settings_for_validation = if has_settings_update {
1546 let historical_has_auth = matches!(historical_settings.get("auth"), Some(Value::Doc(auth_map)) if !auth_map.is_empty());
1547 if !historical_has_auth {
1548 let staged_settings = self.get_local_data::<Doc>(SETTINGS)?.unwrap_or_default();
1549 let staged_has_auth = matches!(staged_settings.get("auth"), Some(Value::Doc(auth_map)) if !auth_map.is_empty());
1550 if staged_has_auth {
1551 staged_settings
1552 } else {
1553 historical_settings
1554 }
1555 } else {
1556 historical_settings
1557 }
1558 } else {
1559 historical_settings
1560 };
1561
1562 // VALIDATION: Ensure that the new settings state (after this transaction) doesn't corrupt auth
1563 // This prevents committing entries that would corrupt the database's auth configuration
1564 if has_settings_update {
1565 // Compute what the new settings state will be after merging local changes
1566 let local_settings = self.get_local_data::<Doc>(SETTINGS)?.unwrap_or_default();
1567 let new_settings = effective_settings_for_validation.merge(&local_settings)?;
1568
1569 // Check if the new settings would have corrupted auth
1570 if new_settings.is_tombstone("auth") {
1571 // Auth was explicitly deleted - this would corrupt the database
1572 return Err(TransactionError::CorruptedAuthConfiguration.into());
1573 } else if let Some(auth_value) = new_settings.get("auth") {
1574 // Auth exists in new settings - check if it's the right type
1575 if !matches!(auth_value, Value::Doc(_)) {
1576 // Auth exists but has wrong type (not a Doc) - this would corrupt the database
1577 return Err(TransactionError::CorruptedAuthConfiguration.into());
1578 }
1579 }
1580 // If auth is None (not configured), that's fine - we allow empty auth
1581 }
1582
1583 // Ensure _index constraint: subtrees referenced in _index must appear in Entry.
1584 // This adds subtrees with None data if they're referenced in _index but not yet in builder.
1585 // First, get the data we need before any async operations
1586 let (_index_data_opt, main_parents, missing_subtrees) = {
1587 let builder_ref = self.entry_builder.lock().unwrap();
1588 let builder = builder_ref
1589 .as_ref()
1590 .ok_or(TransactionError::TransactionAlreadyCommitted)?;
1591
1592 let index_data_opt = builder.data(INDEX).ok().cloned();
1593 let main_parents = builder.parents().unwrap_or_default();
1594 let existing_subtrees = builder.subtrees();
1595
1596 // Find missing subtrees
1597 let missing = if let Some(ref index_data) = index_data_opt
1598 && let Ok(index_doc) = serde_json::from_slice::<Doc>(index_data)
1599 {
1600 index_doc
1601 .keys()
1602 .filter(|name| !existing_subtrees.contains(&name.to_string()))
1603 .cloned()
1604 .collect::<Vec<_>>()
1605 } else {
1606 Vec::new()
1607 };
1608
1609 (index_data_opt, main_parents, missing)
1610 };
1611
1612 // Get tips for missing subtrees (async)
1613 let mut subtree_tips: Vec<(String, Vec<ID>)> = Vec::new();
1614 for subtree_name in missing_subtrees {
1615 let tips = self.get_subtree_tips(&subtree_name, &main_parents).await?;
1616 subtree_tips.push((subtree_name, tips));
1617 }
1618
1619 // Now update the builder with the tips
1620 {
1621 let mut builder_ref = self.entry_builder.lock().unwrap();
1622 let builder = builder_ref
1623 .as_mut()
1624 .ok_or(TransactionError::TransactionAlreadyCommitted)?;
1625
1626 for (subtree_name, tips) in subtree_tips {
1627 builder.set_subtree_parents_mut(&subtree_name, tips);
1628 }
1629
1630 builder.remove_empty_subtrees_mut()?;
1631 }
1632
1633 // Add metadata with settings snapshot for all entries
1634 // Get the backend to access the settings snapshot (do async ops before RefCell borrow)
1635 let db_snapshot = self.db.snapshot().await?;
1636 let settings_snapshot = self
1637 .db
1638 .ops()
1639 .store_snapshot_at(self.db.root_id(), SETTINGS, &db_snapshot)
1640 .await?;
1641
1642 // Clone the builder from RefCell (limit borrow scope to avoid holding across await)
1643 let mut builder = {
1644 let builder_cell = self.entry_builder.lock().unwrap();
1645 let builder_from_cell = builder_cell
1646 .as_ref()
1647 .ok_or(TransactionError::TransactionAlreadyCommitted)?;
1648 builder_from_cell.clone()
1649 };
1650
1651 // Parse existing metadata if present, or create new
1652 let mut metadata = builder
1653 .metadata()
1654 .and_then(|m| serde_json::from_slice::<EntryMetadata>(m).ok())
1655 .unwrap_or(EntryMetadata {
1656 settings_snapshot: Snapshot::EMPTY,
1657 entropy: None,
1658 });
1659
1660 // Update settings snapshot
1661 metadata.settings_snapshot = settings_snapshot;
1662
1663 // Serialize the metadata
1664 let metadata_json = serde_json::to_vec(&metadata)?;
1665
1666 // Add metadata to the entry builder
1667 builder.set_metadata_mut(metadata_json);
1668
1669 // Handle authentication configuration before building
1670 // All entries must now be authenticated - fail if no auth key is configured
1671
1672 // Use provided signing key
1673 let signing_key = if let Some((ref provided_key, ref identity)) = self.provided_signing_key
1674 {
1675 // Use provided signing key directly (already decrypted from UserKeyManager or device key)
1676 let key_clone = provided_key.clone();
1677
1678 // Build AuthInfo from the already-typed SigKey identity
1679 let sig_builder = AuthInfo::builder().key(identity.clone());
1680
1681 // Set auth ID on the entry builder (without signature initially)
1682 builder.set_auth_mut(sig_builder.build());
1683
1684 Some(key_clone)
1685 } else {
1686 // No authentication key configured
1687 return Err(TransactionError::AuthenticationRequired.into());
1688 };
1689 // Encrypt subtree data if encryptors are registered
1690 // This must happen before building the entry to ensure encrypted data is persisted
1691 {
1692 let encryptors = self.encryptors.lock().unwrap();
1693 for subtree_name in builder.subtrees() {
1694 if let Some(encryptor) = encryptors.get(&subtree_name)
1695 && let Ok(plaintext_data) = builder.data(&subtree_name)
1696 && !plaintext_data.is_empty()
1697 {
1698 let ciphertext = encryptor.encrypt(plaintext_data)?;
1699 builder.set_subtree_data_mut(subtree_name.clone(), ciphertext);
1700 }
1701 }
1702 }
1703
1704 // Extract height strategy from settings (defaults to Incremental)
1705 // If this transaction includes settings updates, merge them to get the effective strategy
1706 let settings_for_height = if has_settings_update {
1707 let local_settings = self.get_local_data::<Doc>(SETTINGS)?.unwrap_or_default();
1708 effective_settings_for_validation.merge(&local_settings)?
1709 } else {
1710 effective_settings_for_validation.clone()
1711 };
1712 let height_strategy: HeightStrategy = settings_for_height
1713 .get_json("height_strategy")
1714 .unwrap_or_default();
1715
1716 // Compute heights from parent entries using the configured strategy
1717 {
1718 let backend = self.db.ops();
1719 let instance = self.db.instance()?;
1720 let calculator = height_strategy.into_calculator(instance.clock_arc());
1721
1722 // Compute main tree height using the height strategy
1723 let main_parents = builder.parents().unwrap_or_default();
1724 let max_parent_height = if main_parents.is_empty() {
1725 None
1726 } else {
1727 let mut max_height = 0u64;
1728 for parent_id in &main_parents {
1729 if let Ok(parent) = backend.get(parent_id).await {
1730 max_height = max_height.max(parent.height());
1731 }
1732 }
1733 Some(max_height)
1734 };
1735 let tree_height = calculator.calculate_height(max_parent_height);
1736 builder.set_height_mut(tree_height);
1737
1738 // Compute subtree heights based on per-subtree settings from _index
1739 // System subtrees (prefixed with _) always inherit from tree.
1740 // Regular subtrees check _index for a height_strategy override.
1741 //
1742 // If a subtree has no override, its height is left as None, which means
1743 // Entry.subtree_height() will return the tree height (inheritance).
1744 let index = self.get_index().await.ok();
1745
1746 for subtree_name in builder.subtrees() {
1747 // Determine the effective strategy for this subtree:
1748 // - System subtrees (_settings, _index, etc.): inherit (None)
1749 // - User subtrees: look up in _index, default to inherit (None)
1750 let subtree_strategy: Option<HeightStrategy> = if subtree_name.starts_with('_') {
1751 // System subtrees always inherit from tree
1752 None
1753 } else if let Some(ref idx) = index {
1754 idx.get_subtree_settings(&subtree_name)
1755 .await
1756 .ok()
1757 .and_then(|s| s.height_strategy)
1758 } else {
1759 None
1760 };
1761
1762 match subtree_strategy {
1763 None => {
1764 // Inherit from tree - height stays None (default)
1765 // Entry.subtree_height() will return tree height
1766 }
1767 Some(strategy) => {
1768 // Calculate independent height from subtree parents
1769 let subtree_calculator = strategy.into_calculator(instance.clock_arc());
1770 let subtree_parents =
1771 builder.subtree_parents(&subtree_name).unwrap_or_default();
1772 let max_subtree_parent_height = if subtree_parents.is_empty() {
1773 None
1774 } else {
1775 let mut max_height = 0u64;
1776 for parent_id in &subtree_parents {
1777 if let Ok(parent) = backend.get(parent_id).await
1778 && let Ok(height) = parent.subtree_height(&subtree_name)
1779 {
1780 max_height = max_height.max(height);
1781 }
1782 }
1783 Some(max_height)
1784 };
1785 let subtree_height =
1786 subtree_calculator.calculate_height(max_subtree_parent_height);
1787 builder.set_subtree_height_mut(&subtree_name, Some(subtree_height));
1788 }
1789 }
1790 }
1791 }
1792
1793 // Build the final immutable Entry
1794 let mut entry = builder.build()?;
1795
1796 // CRITICAL VALIDATION: Ensure entry structural integrity before commit
1797 //
1798 // This validation is crucial because the transaction layer has already:
1799 // 1. Discovered proper parent relationships through DAG traversal
1800 // 2. Set up correct subtree parents via find_subtree_parents_from_main_parents()
1801 // 3. Ensured all references point to valid entries in the backend
1802 //
1803 // The validate() call here ensures that:
1804 // - Non-root entries have main tree parents (preventing orphaned nodes)
1805 // - Parent IDs are not empty strings (preventing reference errors)
1806 // - The entry structure is valid before signing and storage
1807 //
1808 // This catches any issues early in the transaction, providing clear error
1809 // messages before the entry is signed or reaches the backend storage layer.
1810 entry.validate()?;
1811
1812 // Sign the entry if we have a signing key
1813 if let Some(signing_key) = signing_key {
1814 let signature = sign_entry(&entry, &signing_key)?;
1815 entry = entry.with_auth(|auth| auth.signature = Some(signature));
1816 }
1817
1818 // Validate authentication (all entries must be authenticated)
1819 let mut validator = AuthValidator::new();
1820
1821 // Get the final settings state for validation
1822 // IMPORTANT: For permission checking, we must use the historical auth configuration
1823 // (before this transaction), not the auth configuration from the current entry.
1824 // This prevents operations from modifying their own permission requirements.
1825
1826 // Extract AuthSettings from effective settings for validation
1827 // IMPORTANT: Distinguish between empty auth vs corrupted/deleted auth:
1828 // - None: No auth ever configured → Allow unsigned operations (empty AuthSettings)
1829 // - Some(Doc): Normal auth configuration → Use it for validation
1830 // - Tombstone (deleted): Auth was configured then deleted → CORRUPTED (fail-safe)
1831 // - Some(other types): Wrong type in auth field → CORRUPTED (fail-safe)
1832 //
1833 // NOTE: Doc::get() hides tombstones (returns None for deleted values), so we need
1834 // to check for tombstones explicitly using is_tombstone() before using get().
1835 let auth_settings_for_validation = if effective_settings_for_validation.is_tombstone("auth")
1836 {
1837 // Auth was configured then explicitly deleted - this is corrupted
1838 return Err(TransactionError::CorruptedAuthConfiguration.into());
1839 } else {
1840 match effective_settings_for_validation.get("auth") {
1841 Some(Value::Doc(auth_doc)) => auth_doc.clone().into(),
1842 None => AuthSettings::new(), // Empty auth - never configured
1843 Some(_) => {
1844 // Auth exists but has wrong type (not a Doc) - this is corrupted
1845 return Err(TransactionError::CorruptedAuthConfiguration.into());
1846 }
1847 }
1848 };
1849
1850 let instance = self.db.instance()?;
1851
1852 // Validate entry (signature + permissions)
1853 let is_valid = validator
1854 .validate_entry(&entry, &auth_settings_for_validation, Some(&instance))
1855 .await?;
1856
1857 if !is_valid {
1858 return Err(TransactionError::EntryValidationFailed.into());
1859 }
1860
1861 let verification_status = VerificationStatus::Verified;
1862
1863 // Get the entry's ID
1864 let id = entry.id();
1865
1866 // Write entry through Instance which handles backend storage and callback dispatch
1867 let instance = self.db.instance()?;
1868 if let Some(guard) = guard {
1869 instance
1870 .put_entry_under_tree_lock(
1871 guard,
1872 self.db.root_id(),
1873 verification_status,
1874 entry.clone(),
1875 WriteSource::Local,
1876 )
1877 .await?;
1878 } else {
1879 instance
1880 .put_entry(
1881 self.db.root_id(),
1882 verification_status,
1883 entry.clone(),
1884 WriteSource::Local,
1885 )
1886 .await?;
1887 }
1888
1889 Ok(id)
1890 }
1891}