eidetica/database/mod.rs
1//! Database module provides functionality for managing collections of related entries.
2//!
3//! A `Database` represents a hierarchical structure of entries, like a traditional database
4//! or a branch in a version control system. Each database has a root entry and maintains
5//! the history and relationships between entries. Database holds a weak reference to its
6//! parent Instance, accessing storage and coordination services through that handle.
7
8use std::{future::Future, sync::Arc};
9
10use rand::{Rng, RngCore, distributions::Alphanumeric};
11use serde_json;
12
13#[cfg(all(unix, feature = "service"))]
14use crate::instance::backend::RemoteBackend;
15#[cfg(all(unix, feature = "service"))]
16use crate::service::client::RemoteConnection;
17use crate::{
18 Error, Instance, Result, Snapshot, Transaction, WeakInstance,
19 auth::{
20 crypto::{PrivateKey, PublicKey},
21 errors::AuthError,
22 settings::AuthSettings,
23 types::{AuthKey, Permission, SigKey},
24 validation::AuthValidator,
25 },
26 backend::VerificationStatus,
27 constants::{ROOT, SETTINGS},
28 crdt::{CRDT, Doc},
29 entry::{Entry, ID},
30 instance::{WriteCallback, WriteEvent, WriteSource, backend::Backend, errors::InstanceError},
31 store::{SettingsStore, Store, Table},
32 sync::DatabaseTicket,
33 user::{SyncSettings, TrackedDatabase, UserError},
34};
35
36#[cfg(test)]
37mod tests;
38
39tokio::task_local! {
40 /// Set while a `verify()`/validation pass is on the call stack.
41 ///
42 /// Verification reads the database (delegation resolution opens trees,
43 /// reads settings → tips), and the access-time auto-verify hook in
44 /// [`Database::snapshot`] would otherwise re-enter verification
45 /// unboundedly. While this is set, the hook is suppressed and reads
46 /// return raw (still `Failed`-filtered) tips.
47 static IN_VERIFY: bool;
48}
49
50fn auto_verify_suppressed() -> bool {
51 IN_VERIFY.try_with(|v| *v).unwrap_or(false)
52}
53
54/// Outcome of reconstructing the `_settings` state an entry pins.
55///
56/// An entry records, in its signed metadata, the `_settings` tips its
57/// signature must be validated against. We can only verify it if this node
58/// holds that full pinned `_settings` ancestor set.
59enum PinnedSettings {
60 /// The pinned `_settings` set is fully present; here is its auth config.
61 Complete(AuthSettings),
62 /// This node does not yet hold the full pinned `_settings` set, so the
63 /// entry cannot be verified yet (it stays `Unverified`).
64 Incomplete,
65}
66
67/// Summary of a [`Database::verify`] pass.
68#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
69pub struct VerifyReport {
70 /// Entries promoted `Unverified` → `Verified` this pass.
71 pub verified: usize,
72 /// Entries marked `Unverified` → `Failed` (definitively bad) this pass.
73 pub failed: usize,
74 /// Entries left `Unverified` (pinned `_settings` not yet held locally).
75 pub still_unverified: usize,
76}
77
78/// A signing key bound to its identity in a database's auth settings.
79///
80/// Pairs the cryptographic signing key with information about how to look up
81/// permissions in the database's auth configuration. The identity determines
82/// which entry in `_settings.auth` this key maps to.
83#[derive(Clone, Debug)]
84pub struct DatabaseKey {
85 signing_key: Box<PrivateKey>,
86 identity: SigKey,
87}
88
89impl DatabaseKey {
90 /// Identity = pubkey derived from signing key. Most common case.
91 pub fn new(signing_key: PrivateKey) -> Self {
92 let pubkey = signing_key.public_key();
93 Self {
94 signing_key: Box::new(signing_key),
95 identity: SigKey::from_pubkey(&pubkey),
96 }
97 }
98
99 /// Identity = explicit SigKey (name, global, delegation, etc.)
100 pub fn with_identity(signing_key: PrivateKey, identity: SigKey) -> Self {
101 Self {
102 signing_key: Box::new(signing_key),
103 identity,
104 }
105 }
106
107 /// Identity = global permission with actual pubkey embedded for verification.
108 pub fn global(signing_key: PrivateKey) -> Self {
109 let pubkey = signing_key.public_key();
110 Self {
111 signing_key: Box::new(signing_key),
112 identity: SigKey::global(&pubkey),
113 }
114 }
115
116 /// Identity = key name lookup.
117 pub fn with_name(signing_key: PrivateKey, name: impl Into<String>) -> Self {
118 Self {
119 signing_key: Box::new(signing_key),
120 identity: SigKey::from_name(name),
121 }
122 }
123
124 /// Get the signing key.
125 pub fn signing_key(&self) -> &PrivateKey {
126 &self.signing_key
127 }
128
129 /// Get the public key.
130 pub fn public_key(&self) -> PublicKey {
131 self.signing_key.public_key()
132 }
133
134 /// Get the identity used for auth settings lookup.
135 pub fn identity(&self) -> &SigKey {
136 &self.identity
137 }
138
139 /// Consume self and return the parts.
140 pub fn into_parts(self) -> (PrivateKey, SigKey) {
141 (*self.signing_key, self.identity)
142 }
143}
144
145impl From<PrivateKey> for DatabaseKey {
146 /// Convert a `PrivateKey` into a `DatabaseKey` with pubkey-derived identity.
147 ///
148 /// This is equivalent to [`DatabaseKey::new`] and covers the most common case
149 /// where the key's identity in auth settings is its own public key.
150 fn from(signing_key: PrivateKey) -> Self {
151 Self::new(signing_key)
152 }
153}
154
155/// Represents a collection of related entries, like a traditional database or a branch in a version control system.
156///
157/// Each `Database` is identified by the ID of its root `Entry` and manages the history of data
158/// associated with that root. It interacts with the underlying storage through the Instance handle.
159#[derive(Clone, Debug)]
160pub struct Database {
161 root: ID,
162 instance: WeakInstance,
163 /// Storage seam `Transaction`/`Store` reads flow through. On a local
164 /// instance this is a clone of the instance's own [`Backend`] (forwarding
165 /// to the backing engine); on a connected instance a per-handle
166 /// [`RemoteBackend`] bound to this database's acting identity. Derived from
167 /// `instance`/construction only — carrying it across
168 /// `with_key`/`allow_unverified` (`..self`) rebuilds is correct.
169 ops: Arc<dyn Backend>,
170 /// Signing key bound to its auth identity for this database
171 key: Option<DatabaseKey>,
172 /// Private user context carried only by handles opened or created through
173 /// a [`User`](crate::user::User) session. This owns caller-local sharing
174 /// preferences; ordinary database handles remain context-free.
175 user_database: Option<Box<Database>>,
176 /// When `false` (default), reads expose only the maximal all-`Verified`
177 /// prefix of the DAG (the "Verified frontier"). When `true`, reads also
178 /// include `Unverified` entries. `Failed` entries are dropped regardless.
179 /// Set via [`Database::allow_unverified`].
180 allow_unverified: bool,
181}
182
183impl Database {
184 pub(crate) fn from_parts(root: ID, instance: WeakInstance, ops: Arc<dyn Backend>) -> Self {
185 Self {
186 root,
187 instance,
188 ops,
189 key: None,
190 user_database: None,
191 allow_unverified: false,
192 }
193 }
194
195 /// Creates a new `Database` instance with a user-provided signing key.
196 ///
197 /// This constructor creates a new database using a signing key that's already in memory
198 /// (e.g., from UserKeyManager), without requiring the key to be stored in the backend.
199 /// This is the preferred method for creating databases in a User context where keys
200 /// are managed separately from the backend.
201 ///
202 /// The created database will use a `DatabaseKey` for all subsequent operations,
203 /// meaning transactions will use the provided key directly rather than looking it up
204 /// from backend storage.
205 ///
206 /// # Auth Bootstrapping
207 ///
208 /// Auth is always bootstrapped with the signing key as `Admin(0)`. Passing auth
209 /// configuration in `initial_settings` is an error — additional keys must be added
210 /// via follow-up transactions after creation.
211 ///
212 /// # Arguments
213 /// * `instance` - Instance handle for storage and coordination
214 /// * `signing_key` - The signing key to use for the initial commit and subsequent operations.
215 /// This key should already be decrypted and ready to use. The public key is derived
216 /// automatically and used as the key identifier in auth settings.
217 /// * `initial_settings` - `Doc` CRDT containing the initial settings for the database.
218 /// Use `Doc::new()` for an empty settings document.
219 ///
220 /// # Returns
221 /// A `Result` containing the new `Database` instance configured with a `DatabaseKey`.
222 ///
223 /// # Example
224 /// ```rust,no_run
225 /// # use eidetica::*;
226 /// # use eidetica::backend::database::InMemory;
227 /// # use eidetica::auth::crypto::generate_keypair;
228 /// # use eidetica::crdt::Doc;
229 /// # #[tokio::main]
230 /// # async fn main() -> Result<()> {
231 /// let instance = Instance::open_backend(Box::new(InMemory::new())).await?;
232 /// let (signing_key, _public_key) = generate_keypair();
233 ///
234 /// let mut settings = Doc::new();
235 /// settings.set("name", "my_database");
236 ///
237 /// // Create database with user-managed key (no backend storage needed)
238 /// let database = Database::create(&instance, signing_key, settings).await?;
239 ///
240 /// // All transactions automatically use the provided key
241 /// let tx = database.new_transaction().await?;
242 /// # Ok(())
243 /// # }
244 /// ```
245 pub async fn create(
246 instance: &Instance,
247 signing_key: PrivateKey,
248 initial_settings: Doc,
249 ) -> Result<Self> {
250 Self::create_with_init(instance, signing_key, initial_settings, async |_| Ok(())).await
251 }
252
253 /// Creates a new database with an initialization callback that runs inside
254 /// the genesis transaction.
255 ///
256 /// This is the underlying constructor used by [`Self::create`] and by
257 /// `User::new_database()` (the builder API). The callback receives the
258 /// genesis transaction after `_settings` and `_root` have been staged but
259 /// before commit, allowing additional subtrees — stores, initial doc data,
260 /// records — to be written into the same entry that establishes the
261 /// database root.
262 ///
263 /// All writes performed in the callback become part of the single genesis
264 /// entry: one signed entry, one backend write, atomic. If the callback
265 /// returns an error, the transaction is dropped without committing and no
266 /// database is created.
267 ///
268 /// # Arguments
269 /// * `instance` - Instance handle for storage and coordination
270 /// * `signing_key` - The private key for this database (becomes Admin(0))
271 /// * `initial_settings` - Initial settings document (must not contain auth)
272 /// * `init` - Callback run against the genesis transaction before commit
273 pub async fn create_with_init<F>(
274 instance: &Instance,
275 signing_key: PrivateKey,
276 initial_settings: Doc,
277 init: F,
278 ) -> Result<Self>
279 where
280 F: AsyncFnOnce(&Transaction) -> Result<()>,
281 {
282 let mut initial_settings = initial_settings;
283 let pubkey = signing_key.public_key();
284
285 // Reject preconfigured auth — Database::create owns auth bootstrapping entirely.
286 if initial_settings.get("auth").is_some() {
287 return Err(Error::Auth(Box::new(AuthError::InvalidAuthConfiguration {
288 reason: "initial_settings must not contain auth configuration; \
289 Database::create bootstraps auth with the signing key as Admin(0)"
290 .to_string(),
291 })));
292 }
293
294 // Bootstrap auth with the signing key as Admin(0)
295 let mut auth_settings = AuthSettings::new();
296 auth_settings.add_key(&pubkey, AuthKey::active(None, Permission::Admin(0)))?;
297 initial_settings.set("auth", auth_settings.as_doc().clone());
298
299 // Create the initial root entry using a temporary Database and Transaction.
300 // This placeholder ID should not exist in the backend, so snapshot will be empty.
301 let bootstrap_placeholder_id = format!(
302 "bootstrap_root_{}",
303 rand::thread_rng()
304 .sample_iter(&Alphanumeric)
305 .take(10)
306 .map(char::from)
307 .collect::<String>()
308 );
309
310 // Create temporary database for bootstrap with DatabaseKey.
311 // This allows the bootstrap transaction to use the provided key directly.
312 let temp_database_for_bootstrap = Database {
313 root: ID::from_bytes(bootstrap_placeholder_id.as_bytes()),
314 instance: instance.downgrade(),
315 ops: instance.backend().clone(),
316 key: Some(DatabaseKey::new(signing_key.clone())),
317 user_database: None,
318 allow_unverified: false,
319 };
320
321 // Create the transaction - it will use the provided key automatically
322 let txn = temp_database_for_bootstrap.new_transaction().await?;
323
324 // IMPORTANT: For the root entry, we need to set the database root to empty/default
325 // so that is_root() returns true and all_roots() can find it
326 txn.set_entry_root(ID::default())?;
327
328 // Populate the SETTINGS and ROOT subtrees for the very first entry
329 txn.update_subtree(SETTINGS, serde_json::to_vec(&initial_settings)?)
330 .await?;
331 txn.update_subtree(ROOT, serde_json::to_vec("")?).await?; // Standard practice for root entry's _root
332
333 // Add entropy to the entry metadata to ensure unique database IDs even with identical settings
334 txn.set_metadata_entropy(rand::thread_rng().next_u64())?;
335
336 // Lock system subtrees (`_settings`, `_root`, `_index`) for the
337 // duration of the init callback. The callback gets a `&Transaction`
338 // that rejects `_*` opens via `get_store` / `Store::open`, so callers
339 // can't accidentally clobber subtrees that `create_with_init` itself
340 // manages. Internal paths (Registry, Store::register's `_index`
341 // writes) bypass `get_store` and remain functional. The guard releases
342 // on drop, so the lock lifts on both early-return and panic-unwind.
343 {
344 let _lock = txn.lock_system_subtrees();
345 init(&txn).await?;
346 }
347
348 // Commit the initial entry
349 let new_root_id = txn.commit().await?;
350
351 // Construct the returned Database, wiring its `ops` to match the
352 // instance's flavour:
353 //
354 // - **Connected (remote) instance** — bind a `RemoteBackend` to the
355 // *new database's* identity (the signing-key's pubkey, self-signed
356 // as `Admin(0)` by the genesis). The connection's login pubkey is
357 // the *caller's* (e.g. the registering admin), which is **not** a
358 // member of the new tree's auth. If we cloned the instance's
359 // session backend here, every read `Transaction::commit` performs
360 // on this database would carry the connection's login identity, and
361 // the server's per-tree gate would deny it. A per-database identity
362 // makes all reads use the tree's own member key — but the server's
363 // gate also requires that key to be in the connection's *session
364 // keyset*, so we `register_session_key(signing_key)` first to do the
365 // proof-of-possession handshake that adds it.
366 // - **Local instance** — clone the instance's backend, unchanged.
367 #[cfg(all(unix, feature = "service"))]
368 if let Some(conn) = instance.remote_connection() {
369 let pubkey_for_identity = signing_key.public_key();
370 conn.register_session_key(&signing_key).await?;
371 return Ok(Self {
372 root: new_root_id.clone(),
373 instance: instance.downgrade(),
374 ops: Arc::new(RemoteBackend::new(
375 conn,
376 Some(SigKey::from_pubkey(&pubkey_for_identity)),
377 )),
378 key: Some(DatabaseKey::new(signing_key)),
379 user_database: None,
380 allow_unverified: false,
381 });
382 }
383
384 Ok(Self {
385 root: new_root_id,
386 instance: instance.downgrade(),
387 ops: instance.backend().clone(),
388 key: Some(DatabaseKey::new(signing_key)),
389 user_database: None,
390 allow_unverified: false,
391 })
392 }
393
394 /// Opens an existing database by its root ID.
395 ///
396 /// Verifies the root entry exists in the backend, then returns a handle
397 /// for read-only access. To perform authenticated writes, chain
398 /// `.with_key(key)` after opening.
399 ///
400 /// # Arguments
401 /// * `instance` - Instance handle for storage and coordination
402 /// * `root_id` - The root entry ID of the database to open
403 ///
404 /// # Errors
405 /// Returns an error if the root entry does not exist in the backend.
406 ///
407 /// # Example
408 /// ```rust,no_run
409 /// # use eidetica::*;
410 /// # use eidetica::backend::database::InMemory;
411 /// # use eidetica::auth::crypto::generate_keypair;
412 /// # #[tokio::main]
413 /// # async fn main() -> Result<()> {
414 /// # let instance = Instance::open_backend(Box::new(InMemory::new())).await?;
415 /// # let (signing_key, _verifying_key) = generate_keypair();
416 /// # let root_id = ID::from_bytes(b"existing_database_root_id");
417 /// // Open database for reading
418 /// let db = Database::open(&instance, &root_id).await?;
419 ///
420 /// // Open database with a signing key for writes
421 /// let db = Database::open(&instance, &root_id).await?.with_key(signing_key);
422 /// let tx = db.new_transaction().await?;
423 /// # Ok(())
424 /// # }
425 /// ```
426 pub async fn open(instance: &Instance, root_id: &ID) -> Result<Self> {
427 // Verify the root entry exists. Surfaces "database doesn't exist"
428 // at open time instead of at first read/write.
429 instance.backend().get(root_id).await?;
430
431 Ok(Self {
432 root: root_id.clone(),
433 instance: instance.downgrade(),
434 ops: instance.backend().clone(),
435 key: None,
436 user_database: None,
437 allow_unverified: false,
438 })
439 }
440
441 /// Open a database for remote access through a service connection.
442 ///
443 /// Constructs a [`Database`] handle whose backing
444 /// [`Backend`](crate::instance::backend::Backend) is a
445 /// [`RemoteBackend`](crate::instance::backend::RemoteBackend) bound to
446 /// `identity`, so every [`Transaction`]/[`Store`] read and write travels
447 /// over the connection as a `DatabaseOp` under that identity. The
448 /// `identity` must match the database's auth settings for the caller's
449 /// key. `Instance::connect` must be used to create the instance.
450 #[cfg(all(unix, feature = "service"))]
451 pub async fn open_remote(
452 instance: &Instance,
453 conn: RemoteConnection,
454 root_id: &ID,
455 identity: SigKey,
456 ) -> Result<Self> {
457 // Probe through the identity-bound backend, not `instance.backend()`.
458 // The instance-level backend carries no per-handle identity, so it acts
459 // as the connection's login pubkey — and a user's login key is not a
460 // member of every tree they hold a key for. A database created under a
461 // per-database key (`User::add_private_key` + `create_database`) grants
462 // only that key, so probing as the login identity denies the existence
463 // check on a database the caller is authorised to open. `identity` has
464 // already been proof-of-possession registered into the connection's
465 // session keyset by the caller, so this is the same identity every
466 // subsequent read and write on the handle travels as.
467 let ops = Arc::new(RemoteBackend::new(conn, Some(identity)));
468 ops.get(root_id).await?;
469 Ok(Self {
470 root: root_id.clone(),
471 instance: instance.downgrade(),
472 ops,
473 key: None,
474 user_database: None,
475 allow_unverified: false,
476 })
477 }
478
479 /// Attach a signing key to this database handle.
480 ///
481 /// The key is stored for use by future transactions. No validation is
482 /// performed; invalid keys will cause errors at commit time or when
483 /// calling [`current_permission`](Self::current_permission).
484 ///
485 /// Calling `with_key` again replaces any previously-attached key — the
486 /// most recent call wins.
487 ///
488 /// To discover which `SigKey` identity to use for a given public key,
489 /// use [`Database::find_sigkeys`].
490 pub fn with_key(self, key: impl Into<DatabaseKey>) -> Self {
491 Self {
492 key: Some(key.into()),
493 ..self
494 }
495 }
496
497 pub(crate) fn with_user_database(self, user_database: Database) -> Self {
498 Self {
499 user_database: Some(Box::new(user_database)),
500 ..self
501 }
502 }
503
504 /// Include `Unverified` entries in this handle's reads.
505 ///
506 /// By default a `Database` exposes only the **Verified frontier**: the
507 /// maximal prefix of the DAG (an ancestor-closed set, starting at the
508 /// root) in which every entry is `Verified`. Tips that are still
509 /// `Unverified` — and everything reachable only through them — are hidden,
510 /// so a default read never reflects state this node could not authenticate.
511 ///
512 /// Calling `allow_unverified` opts this handle into the looser view that
513 /// also includes `Unverified` entries (everything except `Failed`, which
514 /// is always dropped). Use it when you explicitly want to observe
515 /// not-yet-verified data — e.g. freshly synced entries whose pinned
516 /// `_settings` this node does not hold yet.
517 ///
518 /// This is a per-handle setting and composes with [`with_key`](Self::with_key):
519 ///
520 /// ```rust,no_run
521 /// # use eidetica::*;
522 /// # use eidetica::auth::crypto::generate_keypair;
523 /// # async fn example(instance: Instance, root_id: ID) -> Result<()> {
524 /// # let (signing_key, _) = generate_keypair();
525 /// let db = Database::open(&instance, &root_id)
526 /// .await?
527 /// .with_key(signing_key)
528 /// .allow_unverified();
529 /// # Ok(())
530 /// # }
531 /// ```
532 ///
533 /// # Note on CRDT coherence
534 ///
535 /// The Verified frontier is a *prefix* cut, not a per-value filter. An
536 /// interior `Unverified` entry hides all of its descendants from the
537 /// default view even if those descendants are themselves `Verified`,
538 /// because exposing them without their unverifiable ancestor would yield
539 /// an incoherent CRDT state. Run [`verify`](Self::verify) to promote the
540 /// blocking entry, or `allow_unverified` to read past it.
541 pub fn allow_unverified(self) -> Self {
542 Self {
543 allow_unverified: true,
544 ..self
545 }
546 }
547
548 /// Validate a `DatabaseKey` against this database's auth settings.
549 ///
550 /// Checks that:
551 /// 1. The signing key derives to the public key claimed by the identity
552 /// 2. The identity exists in the database's auth settings
553 ///
554 /// Returns the effective permission for the validated key. Callers wanting
555 /// to fail fast on an invalid key should call
556 /// [`current_permission`](Self::current_permission), which wraps this.
557 async fn validate_key(&self, key: &DatabaseKey) -> Result<Permission> {
558 let settings_store = self.get_settings().await?;
559 let auth_settings = settings_store.auth_snapshot().await?;
560 let actual_pubkey = key.public_key();
561 let instance = match key.identity() {
562 // Delegation resolution needs an Instance for cross-tree lookups;
563 // direct identities resolve from `auth_settings` alone.
564 SigKey::Delegation { .. } => Some(self.instance()?),
565 _ => None,
566 };
567 crate::auth::validation::permissions::resolve_identity_permission(
568 &actual_pubkey,
569 key.identity(),
570 &auth_settings,
571 instance.as_ref(),
572 )
573 .await
574 }
575
576 /// Find all SigKeys that a public key can use to access a database.
577 ///
578 /// This static helper method loads a database's authentication settings and returns
579 /// all possible SigKeys that can be used with the given public key. This is useful for
580 /// discovering authentication options before opening a database.
581 ///
582 /// Returns all matching SigKeys including:
583 /// - Specific key names where the pubkey matches
584 /// - Global permission if available
585 /// - Single-hop delegation paths (pubkey found in a directly delegated tree)
586 ///
587 /// The results are **sorted by permission level, highest first**, making it easy to
588 /// select the most privileged access available.
589 ///
590 /// # Arguments
591 /// * `instance` - Instance handle for storage and coordination
592 /// * `root_id` - Root entry ID of the database to check
593 /// * `pubkey` - Public key string (e.g., "Ed25519:abc123...") to look up
594 ///
595 /// # Returns
596 /// A vector of (SigKey, Permission) tuples, sorted by permission (highest first).
597 /// Returns empty vector if no valid access methods are found.
598 ///
599 /// # Errors
600 /// Returns an error if:
601 /// - Database cannot be loaded
602 /// - Auth settings cannot be parsed
603 ///
604 /// # Example
605 /// ```rust,no_run
606 /// # use eidetica::*;
607 /// # use eidetica::database::DatabaseKey;
608 /// # use eidetica::backend::database::InMemory;
609 /// # use eidetica::auth::crypto::generate_keypair;
610 /// # use eidetica::auth::types::SigKey;
611 /// # #[tokio::main]
612 /// # async fn main() -> Result<()> {
613 /// # let instance = Instance::open_backend(Box::new(InMemory::new())).await?;
614 /// # let (signing_key, pubkey) = generate_keypair();
615 /// # let root_id = ID::from_bytes(b"database_root_id");
616 /// // Find all SigKeys this pubkey can use (sorted highest permission first)
617 /// let sigkeys = Database::find_sigkeys(&instance, &root_id, &pubkey).await?;
618 ///
619 /// // Use the first available SigKey (highest permission)
620 /// if let Some((sigkey, _permission)) = sigkeys.first() {
621 /// let key = DatabaseKey::with_identity(signing_key, sigkey.clone());
622 /// let database = Database::open(&instance, &root_id).await?.with_key(key);
623 /// }
624 /// # Ok(())
625 /// # }
626 /// ```
627 pub async fn find_sigkeys(
628 instance: &Instance,
629 root_id: &ID,
630 pubkey: &PublicKey,
631 ) -> Result<Vec<(SigKey, Permission)>> {
632 use crate::auth::{
633 types::DelegationStep,
634 validation::{AuthValidator, permissions::select_effective_permission},
635 };
636
637 // Create temporary database to load settings (no key source needed for reading)
638 let temp_db = Self::open(instance, root_id).await?;
639
640 // Load auth settings
641 let settings_store = temp_db.get_settings().await?;
642 let auth_settings = settings_store.auth_snapshot().await?;
643
644 // Find direct SigKeys for this pubkey
645 let mut results = auth_settings.find_all_sigkeys_for_pubkey(pubkey);
646
647 // Scan single-hop delegation paths. Resolution — permission-bounds
648 // clamping, key status, the tip floor — is delegated to the validator,
649 // the single delegation walker, so this stays pure discovery: find which
650 // delegated trees list the pubkey, then resolve each through the shared
651 // path at the delegated tree's *current* tips.
652 // FIXME: deep nested delegations can't use this
653 if let Ok(delegated_trees) = auth_settings.get_all_delegated_trees() {
654 let mut validator = AuthValidator::new();
655 for delegated_root_id in delegated_trees.keys() {
656 // Load the delegated tree's auth to see which of the pubkey's
657 // hints it lists (enumeration only — the validator re-resolves).
658 let delegated_auth = match Self::open(instance, delegated_root_id).await {
659 Ok(db) => match db.get_settings().await {
660 Ok(s) => match s.auth_snapshot().await {
661 Ok(a) => a,
662 Err(_) => continue,
663 },
664 Err(_) => continue,
665 },
666 Err(_) => continue,
667 };
668 let delegated_sigkeys = delegated_auth.find_all_sigkeys_for_pubkey(pubkey);
669 if delegated_sigkeys.is_empty() {
670 continue;
671 }
672
673 // Current tips = the "live authority" question (caller-supplied),
674 // distinct from a signature's claimed tips. They are at or above
675 // the committed floor, so the validator's floor check passes.
676 let tips = match instance.backend().snapshot(delegated_root_id).await {
677 Ok(snap) => snap.tips().to_vec(),
678 Err(_) => continue,
679 };
680
681 for (delegated_sk, _) in delegated_sigkeys {
682 let delegation_sigkey = SigKey::Delegation {
683 path: vec![DelegationStep {
684 tree: delegated_root_id.clone(),
685 tips: tips.clone(),
686 }],
687 hint: delegated_sk.hint().clone(),
688 };
689 // Resolve through the single delegation walker; bounds,
690 // status, and the tip floor all live there, not here.
691 let Ok(resolved) = validator
692 .resolve_sig_key(&delegation_sigkey, &auth_settings, Some(instance))
693 .await
694 else {
695 continue;
696 };
697 if let Some(perm) = select_effective_permission(&resolved, pubkey) {
698 results.push((delegation_sigkey, perm));
699 }
700 }
701 }
702 }
703
704 // Sort by permission, highest first
705 results.sort_by_key(|b| std::cmp::Reverse(b.1));
706 Ok(results)
707 }
708
709 /// Whether `pubkey` may access the database at `root_id` with at least
710 /// `permission`.
711 ///
712 /// This is the **pubkey-only** access decision — the caller holds a key but
713 /// has not presented a signing identity (bootstrap deciding whether a key
714 /// already qualifies, for example). It resolves the same authority
715 /// [`find_sigkeys`](Self::find_sigkeys) does — direct grants, the global
716 /// `*` grant, and delegated authority — and never counts revoked keys.
717 ///
718 /// Delegated authority is discovered **one hop deep**: only databases this
719 /// one delegates to directly are searched, so a key reachable through a
720 /// chain of delegations answers `false` here. The limit belongs to
721 /// discovery, not to delegation — when an identity *is* presented (signing
722 /// an entry, a service operation), authorization goes through the resolver
723 /// behind [`validate_key`](Self::validate_key), which walks a delegation
724 /// path of any length because the signer names the path rather than making
725 /// the resolver search for it.
726 pub async fn can_access(
727 instance: &Instance,
728 root_id: &ID,
729 pubkey: &PublicKey,
730 permission: &Permission,
731 ) -> Result<bool> {
732 let sigkeys = Self::find_sigkeys(instance, root_id, pubkey).await?;
733 Ok(sigkeys
734 .first()
735 .is_some_and(|(_, granted)| *granted >= *permission))
736 }
737
738 /// Get the auth identity for this database's configured key.
739 pub fn auth_identity(&self) -> Option<&SigKey> {
740 self.key.as_ref().map(|k| &k.identity)
741 }
742
743 /// Register a callback to be invoked when entries are written to this database.
744 ///
745 /// The callback fires for **both** local writes (transaction commits) and remote
746 /// writes (sync). Branch on [`WriteEvent::source`](crate::WriteEvent::source) inside
747 /// the closure if you only care about one.
748 ///
749 /// Returns a [`WriteCallback`] handle. **Drop it to unregister.** Call
750 /// [`WriteCallback::detach`] to leave the callback registered for the life
751 /// of the [`Instance`] without holding the handle.
752 ///
753 /// **Important:** Callbacks are registered at the Instance level and fire for all
754 /// writes to the database tree (identified by root ID), regardless of which
755 /// `Database` handle performed the write or registered the callback.
756 ///
757 /// # Callback contract
758 ///
759 /// - **Local writes**: fires once per transaction commit. The
760 /// [`WriteEvent`] carries cursor brackets only — call
761 /// [`Database::ids_added`] with `event.previous_tips()` and
762 /// `event.post_tips()` to enumerate the single new entry.
763 /// - **Remote writes**: fires once per sync batch (not per entry).
764 /// `ids_added(prev, post)` yields the full set of new IDs in topo
765 /// order.
766 /// - All entries that fall inside the cursor advance are fully
767 /// persisted before the callback fires.
768 /// - Errors are logged but do not prevent other callbacks from running.
769 /// - The `db` argument is a **read-only** [`Database`] handle (no
770 /// [`DatabaseKey`] configured): you can read settings, entries, and
771 /// metadata, but cannot commit transactions through it. To write from
772 /// inside a callback, resolve through `db.instance()?` and open the
773 /// database with the appropriate key.
774 /// - **Reentrance**: writes are serialized per-tree via an async lock
775 /// that is held while callbacks run. A callback must not commit a
776 /// transaction on the same tree it was invoked for — that would
777 /// deadlock. Spawn a task or write to a different tree instead.
778 ///
779 /// # Callback timing
780 ///
781 /// On a **local** [`Instance`], callbacks complete before
782 /// [`Transaction::commit`](crate::Transaction::commit)`.await` returns —
783 /// the firing path is inline with the write and the per-tree lock is
784 /// held the whole way.
785 ///
786 /// On a **connected** [`Instance`] (one built via
787 /// [`Instance::connect`](crate::Instance::connect)), callbacks fire
788 /// when a `Notification::DatabaseWrite` arrives back from the daemon —
789 /// typically microseconds after `commit().await` returns, but
790 /// asynchronous to it. This is by design: the daemon is the **sole
791 /// orderer** of writes on a connected setup, so every subscriber —
792 /// including the committing client — observes callbacks in the same
793 /// daemon-canonical order. Two clients submitting concurrently see
794 /// each other's writes in the same sequence; a single client's
795 /// callbacks never race a remote-source event into a different
796 /// observable order than the daemon has it. The trade-off is a small
797 /// asynchrony at the commit boundary — if you need a synchronous
798 /// "callback fired before commit returned" guarantee, use a local
799 /// [`Instance`].
800 ///
801 /// Ordering is preserved end-to-end, per tree. On the daemon the
802 /// per-tree write lock is held across the callback fan-out, and each
803 /// subscription's notification is emitted by a *synchronous* channel
804 /// send inside that locked section (see the `SubscribeWrites` handler
805 /// in `service::server`), so two clients writing the same tree
806 /// concurrently cannot interleave their frames. On the client,
807 /// notifications are routed by `root_id` to a per-tree worker that
808 /// `await`s each user callback to completion before pulling the next —
809 /// so two callbacks for the *same* tree never overlap and never
810 /// reorder. Callbacks for *different* trees run on independent workers
811 /// and progress concurrently: a slow callback on one tree does not
812 /// stall another. Spawn from inside the callback if you need fanout
813 /// within a single tree.
814 ///
815 /// # Settled-state trigger — but not a settled-state bracket
816 ///
817 /// Callbacks fire **only for entries that have passed local
818 /// verification** — i.e. entries the system considers `Verified`.
819 /// Direct `put_entry(Verified, _)` (local commits, trusted
820 /// in-process writes) fires immediately. Off-node entries that
821 /// arrive `Unverified` (sync ingest, wire-submitted entries) are
822 /// stored silently; the subsequent verification pass is where
823 /// promotion to `Verified` happens, and a Verified fire follows
824 /// from there once that pass runs.
825 ///
826 /// Both ingest paths run that pass inline: `put_remote_entries`
827 /// (sync) and the service `SubmitSignedEntry` handler each call
828 /// `verify()` immediately after storing the batch, so a promoted
829 /// entry fires its `Verified` event without waiting for a reader to
830 /// trigger the access-time auto-verify hook. A listen-only
831 /// subscriber therefore sees sync-arrived content promptly.
832 ///
833 /// What fires is settled; what the event *brackets* is not. The
834 /// event's cursors are raw DAG frontiers, so an `Unverified` or
835 /// `Failed` entry that happens to be a tip lands inside the bracket
836 /// and [`Self::ids_added`] will enumerate it. Filter on
837 /// `Backend::get_verification_status` before acting on IDs derived
838 /// from a bracket if your callback ingests entry contents.
839 ///
840 /// On a connected instance the first `on_write` registration for a
841 /// given tree lazily sends a `SubscribeWrites` op to the daemon;
842 /// further registrations on the same tree reuse that subscription.
843 /// Dropping the last callback for a tree marks the subscription idle
844 /// rather than unsubscribing immediately, so a quick re-registration
845 /// costs no round-trip; a sweep sends `UnsubscribeWrites` once the
846 /// grace window elapses. Disconnecting the client unsubscribes
847 /// everything.
848 ///
849 /// # Example
850 /// ```rust,no_run
851 /// # use eidetica::*;
852 /// # use eidetica::crdt::Doc;
853 /// # use eidetica::backend::database::InMemory;
854 /// # use eidetica::auth::crypto::PrivateKey;
855 /// # #[tokio::main]
856 /// # async fn main() -> Result<()> {
857 /// let instance = Instance::open_backend(Box::new(InMemory::new())).await?;
858 /// # let signing_key = PrivateKey::generate();
859 /// # let database = Database::create(&instance, signing_key, Doc::new()).await?;
860 ///
861 /// let cb = database.on_write(|event, db| {
862 /// let source = event.source();
863 /// let prev = event.previous_tips().clone();
864 /// let post = event.post_tips().clone();
865 /// let db = db.clone();
866 /// async move {
867 /// let new_ids = db.ids_added(&prev, &post).await?;
868 /// println!("{} entries written to {} ({source:?})", new_ids.len(), db.root_id());
869 /// Ok(())
870 /// }
871 /// }).await?;
872 ///
873 /// // Drop `cb` to unregister, or:
874 /// cb.detach(); // keep registered for the life of the Instance
875 /// # Ok(())
876 /// # }
877 /// ```
878 ///
879 /// If a callback needs the [`Instance`], call [`Database::instance`] on
880 /// the `db` argument.
881 pub async fn on_write<F, Fut>(&self, callback: F) -> Result<WriteCallback>
882 where
883 F: for<'a> Fn(&'a WriteEvent, &'a Database) -> Fut + Send + std::marker::Sync + 'static,
884 Fut: std::future::Future<Output = Result<()>> + Send + 'static,
885 {
886 // Convenience: read current tips, then register at those tips.
887 //
888 // The "current tips" here come from this method's own
889 // `get_tips` call, **not** from any read the caller may have
890 // done before calling `on_write`. If the caller built initial
891 // state from a separate read at some `T0`, a write may have
892 // landed between that read and this `get_tips`, putting the
893 // cursor at `T1 ≥ T0`. The first callback fire's
894 // `previous_tips` will then be `T1`, not `T0`, and the caller
895 // observes `previous_tips != their_initial_tips` (a small
896 // race-window mismatch — typically one lock-acquisition apart
897 // on a local instance).
898 //
899 // Use [`Self::on_write_at_tips`] to close that window: pass
900 // exactly the tips your initial-state read used, and the first
901 // fire's `previous_tips` will match.
902 // Propagate rather than defaulting: `Snapshot::EMPTY` is the
903 // documented "I have no initial state; replay from the beginning"
904 // cursor, so masking a transient read error here would silently turn
905 // "subscribe from now" into a full-history replay on the first fire.
906 let tips = self.snapshot().await?;
907 self.on_write_at_tips(tips, callback).await
908 }
909
910 /// Register a callback with an explicit initial cursor.
911 ///
912 /// Primitive form of [`Self::on_write`]. The `tips` you pass become
913 /// this callback's initial cursor, and the first event the
914 /// callback receives will have `previous_tips = tips`. Subsequent
915 /// events' `previous_tips` are the previous event's post-write
916 /// tips — each callback owns its own continuous timeline.
917 ///
918 /// Use this when you need a hard guarantee that the cursor matches
919 /// some other tip set you've already used (typically the tips
920 /// returned by an initial-state read you did before subscribing):
921 ///
922 /// ```rust,no_run
923 /// # use eidetica::*;
924 /// # async fn doit(db: Database) -> Result<()> {
925 /// let tips = db.snapshot().await?;
926 /// // ... read initial state at `tips` ...
927 /// let _cb = db.on_write_at_tips(tips, |_event, _db| async move { Ok(()) }).await?;
928 /// // First callback fire's previous_tips will exactly equal the
929 /// // `tips` we read at — no race-window mismatch with subsequent
930 /// // commits.
931 /// # Ok(())
932 /// # }
933 /// ```
934 pub async fn on_write_at_tips<F, Fut>(
935 &self,
936 tips: Snapshot,
937 callback: F,
938 ) -> Result<WriteCallback>
939 where
940 F: for<'a> Fn(&'a WriteEvent, &'a Database) -> Fut + Send + std::marker::Sync + 'static,
941 Fut: std::future::Future<Output = Result<()>> + Send + 'static,
942 {
943 let instance = self.instance()?;
944 let tree_id = self.root_id().clone();
945 // Local registry uses `tips` as the per-callback cursor; the
946 // wire path (if any) also needs it so the daemon-side
947 // subscription cursor is pinned at the same value. Clone once
948 // here; a `Snapshot` is a tiny tip-set.
949 let id = instance.register_write_callback(tree_id.clone(), tips.clone(), callback);
950 let cb = WriteCallback::new_per_database(instance.downgrade(), tree_id.clone(), id);
951
952 // On a connected (daemon-backed) instance, ensure the daemon is
953 // pushing notifications for this tree before returning — otherwise
954 // an immediately-following commit could race the subscribe and
955 // lose its notification, since the daemon doesn't replay missed
956 // events. `subscribe_writes` is itself concurrency-safe and
957 // idempotent: it short-circuits when the tree is already
958 // subscribed and serialises racing registrations through a
959 // per-tree `Notify`, so every caller observes the subscription
960 // as live before returning. Database handles built with a key
961 // carry an explicit identity; keyless handles fall back to
962 // `SigKey::default()`, which the daemon resolves to the
963 // connection's login pubkey.
964 //
965 // The `tips` we hand to `subscribe_writes` become the daemon's
966 // subscription cursor, so the first notification's
967 // `previous_tips` exactly equals the tips this caller passed
968 // in — no race-window mismatch between the daemon's view and
969 // the client's local cursor.
970 #[cfg(all(unix, feature = "service"))]
971 if let Some(conn) = instance.remote_connection() {
972 let identity = self.auth_identity().cloned().unwrap_or_default();
973 conn.subscribe_writes(tree_id, identity, tips).await?;
974 } else {
975 let _ = tips;
976 }
977
978 Ok(cb)
979 }
980
981 /// Get the ID of the root entry
982 pub fn root_id(&self) -> &ID {
983 &self.root
984 }
985
986 /// Upgrade the weak instance reference to a strong reference.
987 ///
988 /// `Database` holds a [`WeakInstance`](crate::WeakInstance), so this can
989 /// fail if the owning [`Instance`] has already been dropped.
990 pub fn instance(&self) -> Result<Instance> {
991 self.instance
992 .upgrade()
993 .ok_or_else(|| Error::Instance(Box::new(InstanceError::InstanceDropped)))
994 }
995
996 /// Get a clone of the backend seam.
997 pub fn backend(&self) -> Result<Arc<dyn Backend>> {
998 Ok(self.instance()?.backend().clone())
999 }
1000
1001 /// The storage seam this handle's `Transaction`/`Store` reads/writes flow
1002 /// through (a [`LocalBackend`](crate::instance::backend::LocalBackend) clone
1003 /// of the instance's backend, or a per-handle
1004 /// [`RemoteBackend`](crate::instance::backend::RemoteBackend)).
1005 pub(crate) fn ops(&self) -> &dyn Backend {
1006 self.ops.as_ref()
1007 }
1008
1009 /// Retrieve the root entry from the backend
1010 pub async fn get_root(&self) -> Result<Entry> {
1011 let instance = self.instance()?;
1012 instance.get(&self.root).await
1013 }
1014
1015 /// Get a read-only settings store for the database.
1016 ///
1017 /// Returns a SettingsStore that provides access to the database's settings.
1018 /// Since this creates an internal transaction that is never committed, any
1019 /// modifications made through the returned store will not persist.
1020 ///
1021 /// For making persistent changes to settings, create a transaction and use
1022 /// `Transaction::get_settings()` instead.
1023 ///
1024 /// # Returns
1025 /// A `Result` containing the `SettingsStore` for settings or an error.
1026 ///
1027 /// # Example
1028 /// ```rust,no_run
1029 /// # use eidetica::Database;
1030 /// # async fn example(database: Database) -> eidetica::Result<()> {
1031 /// // Read-only access
1032 /// let settings = database.get_settings().await?;
1033 /// let name = settings.get_name().await?;
1034 ///
1035 /// // For modifications, use a transaction:
1036 /// let txn = database.new_transaction().await?;
1037 /// let settings = txn.get_settings()?;
1038 /// settings.set_name("new_name").await?;
1039 /// txn.commit().await?;
1040 /// # Ok(())
1041 /// # }
1042 /// ```
1043 pub async fn get_settings(&self) -> Result<SettingsStore> {
1044 let txn = self.new_transaction().await?;
1045 txn.get_settings()
1046 }
1047
1048 /// Read this handle owner's private synchronization preferences.
1049 ///
1050 /// Database settings belong to the replicated database. Synchronization
1051 /// preferences instead live in the owning user's private database, so only
1052 /// handles returned by [`User::create_database`](crate::user::User::create_database)
1053 /// or [`User::open_database`](crate::user::User::open_database) carry this
1054 /// capability.
1055 pub async fn sync_settings(&self) -> Result<SyncSettings> {
1056 self.authorize_sharing().await?;
1057 Ok(self.tracked_database().await?.sync_settings)
1058 }
1059
1060 /// Whether this handle owner has enabled sharing for this database.
1061 pub async fn is_shared(&self) -> Result<bool> {
1062 Ok(self.sync_settings().await?.sync_enabled)
1063 }
1064
1065 /// Set this handle owner's durable sharing preference.
1066 ///
1067 /// If the connection fails while committing, the write may have succeeded.
1068 /// Callers can safely retry this idempotent write or read
1069 /// [`Self::sync_settings`] again.
1070 pub async fn set_shared(&self, shared: bool) -> Result<()> {
1071 self.authorize_sharing().await?;
1072 let user_database = self.user_database()?;
1073 let tx = user_database.new_transaction().await?;
1074 let table = tx.get_store::<Table<TrackedDatabase>>("databases").await?;
1075 let key = self.root_id().to_string();
1076 let mut tracked = table
1077 .get(&key)
1078 .await
1079 .map_err(|_| UserError::DatabaseNotTracked {
1080 database_id: self.root_id().clone(),
1081 })?;
1082
1083 if tracked.sync_settings.sync_enabled == shared {
1084 return Ok(());
1085 }
1086
1087 tracked.sync_settings.sync_enabled = shared;
1088 table.set(&key, tracked).await?;
1089 tx.commit().await?;
1090 Ok(())
1091 }
1092
1093 /// Enable this handle owner's sharing preference.
1094 pub async fn share(&self) -> Result<()> {
1095 self.set_shared(true).await
1096 }
1097
1098 /// Disable this handle owner's sharing preference.
1099 pub async fn stop_sharing(&self) -> Result<()> {
1100 self.set_shared(false).await
1101 }
1102
1103 /// Return a point-in-time locator when this handle owner is sharing.
1104 pub async fn ticket(&self) -> Result<DatabaseTicket> {
1105 if !self.is_shared().await? {
1106 return Err(UserError::DatabaseNotShared {
1107 database_id: self.root_id().clone(),
1108 }
1109 .into());
1110 }
1111
1112 #[cfg(all(unix, feature = "service"))]
1113 if let Some(conn) = self.instance()?.remote_connection() {
1114 return conn
1115 .database_ticket(
1116 self.root_id(),
1117 self.auth_identity().cloned().unwrap_or_default(),
1118 )
1119 .await;
1120 }
1121
1122 crate::user::ticket_locator(&self.instance()?, self.root_id()).await
1123 }
1124
1125 fn user_database(&self) -> Result<&Database> {
1126 self.user_database.as_deref().ok_or_else(|| {
1127 UserError::MissingDatabaseCapability {
1128 database_id: self.root_id().clone(),
1129 }
1130 .into()
1131 })
1132 }
1133
1134 async fn tracked_database(&self) -> Result<TrackedDatabase> {
1135 let user_database = self.user_database()?;
1136 user_database
1137 .get_store_viewer::<Table<TrackedDatabase>>("databases")
1138 .await?
1139 .get(&self.root_id().to_string())
1140 .await
1141 .map_err(|_| {
1142 UserError::DatabaseNotTracked {
1143 database_id: self.root_id().clone(),
1144 }
1145 .into()
1146 })
1147 }
1148
1149 async fn authorize_sharing(&self) -> Result<()> {
1150 self.user_database()?;
1151 if self.current_permission().await? >= Permission::Read {
1152 Ok(())
1153 } else {
1154 Err(UserError::InsufficientPermissions.into())
1155 }
1156 }
1157
1158 /// Get the name of the database from its settings store
1159 pub async fn get_name(&self) -> Result<String> {
1160 let settings = self.get_settings().await?;
1161 settings.get_name().await
1162 }
1163
1164 /// Create a new atomic transaction on this database
1165 ///
1166 /// This creates a new atomic transaction containing a new Entry.
1167 /// The atomic transaction will be initialized with the current state of the database.
1168 /// If a default authentication key is set, the transaction will use it for signing.
1169 ///
1170 /// # Returns
1171 /// A `Result<Transaction>` containing the new atomic transaction
1172 pub async fn new_transaction(&self) -> Result<Transaction> {
1173 let snapshot = self.snapshot().await?;
1174 self.new_transaction_at(&snapshot).await
1175 }
1176
1177 /// Create a new atomic transaction on this database anchored at a specific snapshot.
1178 ///
1179 /// The transaction's parents are taken from the provided snapshot's tips instead of
1180 /// the database's current state. This allows creating complex DAG structures
1181 /// like diamond patterns for testing and advanced use cases.
1182 ///
1183 /// # Arguments
1184 /// * `snapshot` - The snapshot to anchor the transaction at
1185 ///
1186 /// # Returns
1187 /// A `Result<Transaction>` containing the new atomic transaction
1188 pub async fn new_transaction_at(&self, snapshot: &Snapshot) -> Result<Transaction> {
1189 let mut txn = Transaction::new_at(self, snapshot).await?;
1190
1191 // Set provided signing key from DatabaseKey
1192 if let Some(key) = &self.key {
1193 txn.set_provided_key(*key.signing_key.clone(), key.identity.clone());
1194 }
1195
1196 Ok(txn)
1197 }
1198
1199 /// Gather everything a client needs to build and sign a transaction
1200 /// locally for the given stores, with parents drawn from `scope`'s
1201 /// projection.
1202 ///
1203 /// This is **single-sourced**: both the server's `BeginTransaction`
1204 /// handler and the Phase-3 remote seam call it, so
1205 /// `Transaction::commit`'s build-sign path has one source of truth for
1206 /// context gathering.
1207 ///
1208 /// `scope=AllowUnverified` opens against the raw DAG (only `Failed`
1209 /// dropped); the default `Verified` scope uses the Verified frontier.
1210 /// The returned [`TransactionContext`] carries everything needed for
1211 /// one round-trip transaction build: main parents + heights, per-store
1212 /// subtree parents + heights, settings tips, and the merged `_settings`
1213 /// CRDT state this entry is authored against.
1214 #[cfg(all(unix, feature = "service"))]
1215 pub async fn transaction_context(
1216 &self,
1217 stores: &[String],
1218 scope: crate::service::protocol::ReadScope,
1219 ) -> Result<crate::service::protocol::TransactionContext> {
1220 use crate::service::protocol::{ReadScope, TransactionContext};
1221
1222 // -- scope-sensitive main tips --------------------------------
1223 let db_for_tips = Database {
1224 allow_unverified: matches!(scope, ReadScope::AllowUnverified),
1225 ..self.clone()
1226 };
1227 let main_snapshot = db_for_tips.snapshot().await?;
1228 let main_tips = main_snapshot.tips();
1229
1230 // -- main parents: (tip, height) ------------------------------
1231 let mut main_parents = Vec::with_capacity(main_tips.len());
1232 for tip in main_tips {
1233 let entry = self.ops().get(tip).await?;
1234 main_parents.push((tip.clone(), entry.height()));
1235 }
1236
1237 // -- per-store subtree parents: (tip, subtree_height) ---------
1238 let mut subtree_parents = std::collections::BTreeMap::new();
1239 for store in stores {
1240 let child_snap = self
1241 .ops()
1242 .store_snapshot_at(self.root_id(), store, &main_snapshot)
1243 .await?;
1244 let mut pairs = Vec::with_capacity(child_snap.len());
1245 for tip in child_snap.tips() {
1246 let entry = self.ops().get(tip).await?;
1247 let height = entry.subtree_height(store).unwrap_or(0);
1248 pairs.push((tip.clone(), height));
1249 }
1250 subtree_parents.insert(store.clone(), pairs);
1251 }
1252
1253 // -- settings tips (pinned in entry metadata) -----------------
1254 let settings_tips = self
1255 .ops()
1256 .store_snapshot_at(self.root_id(), SETTINGS, &main_snapshot)
1257 .await?
1258 .into_tips();
1259
1260 // -- merged _settings state as serde_json::Value --------------
1261 let txn = Transaction::new_at(self, &main_snapshot).await?;
1262 let settings_doc: Doc = txn.get_full_state(SETTINGS).await?;
1263 let settings_value = serde_json::to_value(&settings_doc)?;
1264
1265 Ok(TransactionContext {
1266 main_parents,
1267 subtree_parents,
1268 settings_tips,
1269 settings_value,
1270 })
1271 }
1272
1273 /// Server-materialized merged state of an **unencrypted** store, as a
1274 /// `serde_json::Value` against the database's Verified frontier.
1275 ///
1276 /// Creates an ephemeral transaction, deserializes every entry's
1277 /// store data as [`Doc`], and merges them via Doc's LWW merge —
1278 /// the same merge `Store<T>` would perform client-side. All current
1279 /// store types (DocStore, Table, Settings) serialize their data as
1280 /// JSON, so `Doc`-typed deserialization works universally.
1281 ///
1282 /// # Encrypted stores
1283 ///
1284 /// Encrypted stores cannot be materialized this way (the ephemeral
1285 /// transaction has no encryptor, so `serde_json::from_slice::<Doc>`
1286 /// would fail on ciphertext). The caller must use
1287 /// [`get_store_entries`](Self::get_store_entries) for encrypted
1288 /// stores and decrypt+merge client-side.
1289 pub async fn get_store_state(&self, store: &str) -> Result<serde_json::Value> {
1290 let txn = self.new_transaction().await?;
1291 let state: Doc = txn.get_full_state(store).await?;
1292 Ok(serde_json::to_value(&state)?)
1293 }
1294
1295 /// Ordered (by subtree height), verifiable, opaque store entries
1296 /// reachable from `tips` within `scope`.
1297 ///
1298 /// This is the **universal** primitive — works for encrypted and
1299 /// unencrypted stores alike because it returns raw [`Entry`] records
1300 /// with opaque [`RawData`](crate::entry::RawData); no deserialization
1301 /// or merge runs server-side. The per-subtree-height ordering
1302 /// (ascending, then by ID for tiebreaking) is exactly the canonical
1303 /// CRDT replay order produced by
1304 /// [`sort_entries_by_subtree_height`](crate::backend::database::in_memory::cache::sort_entries_by_subtree_height).
1305 ///
1306 /// When `scope` is [`ReadScope::Verified`] and `tips` are the
1307 /// Verified-frontier tips from [`snapshot`](Self::snapshot),
1308 /// every returned entry is guaranteed `Verified` (the frontier is
1309 /// ancestor-closed). For [`ReadScope::AllowUnverified`], entries
1310 /// reachable from unverified tips are included.
1311 #[cfg(all(unix, feature = "service"))]
1312 pub async fn get_store_entries(
1313 &self,
1314 store: &str,
1315 tips: &[ID],
1316 _scope: crate::service::protocol::ReadScope,
1317 ) -> Result<Vec<Entry>> {
1318 let snapshot = Snapshot::from(tips.to_vec());
1319 self.ops().store_at(self.root_id(), store, &snapshot).await
1320 }
1321
1322 /// Execute a closure within a transaction and commit the result.
1323 ///
1324 /// This is a convenience wrapper for the common pattern of creating a transaction,
1325 /// performing store operations, and committing. The transaction is committed after
1326 /// the closure returns `Ok`. If the closure returns `Err`, the transaction is
1327 /// dropped without committing.
1328 ///
1329 /// For read-only access, use [`get_store_viewer`](Self::get_store_viewer) instead.
1330 ///
1331 /// # Arguments
1332 /// * `f` - A closure that receives the [`Transaction`] and performs store operations.
1333 /// The closure should return `Ok(R)` on success.
1334 ///
1335 /// # Returns
1336 /// On success, returns the value produced by the closure after committing.
1337 /// The commit ID is not returned; use [`new_transaction`](Self::new_transaction)
1338 /// directly if you need it.
1339 ///
1340 /// # Errors
1341 /// Returns an error if transaction creation, the closure, or commit fails.
1342 /// If the closure fails, the transaction is not committed.
1343 ///
1344 /// # Example
1345 /// ```rust,no_run
1346 /// # use eidetica::*;
1347 /// # use eidetica::store::Table;
1348 /// # use serde::{Serialize, Deserialize};
1349 /// # #[derive(Clone, Serialize, Deserialize)]
1350 /// # struct Todo { title: String }
1351 /// # async fn example(db: Database) -> Result<()> {
1352 /// // Insert a record and get its generated key
1353 /// let key = db.with_transaction(|txn| async move {
1354 /// let store = txn.get_store::<Table<Todo>>("todos").await?;
1355 /// store.insert(Todo { title: "Buy milk".into() }).await
1356 /// }).await?;
1357 ///
1358 /// // Multiple operations in one atomic transaction
1359 /// db.with_transaction(|txn| async move {
1360 /// let store = txn.get_store::<Table<Todo>>("todos").await?;
1361 /// store.insert(Todo { title: "First".into() }).await?;
1362 /// store.insert(Todo { title: "Second".into() }).await?;
1363 /// Ok(())
1364 /// }).await?;
1365 /// # Ok(())
1366 /// # }
1367 /// ```
1368 pub async fn with_transaction<F, Fut, R>(&self, f: F) -> Result<R>
1369 where
1370 F: FnOnce(Transaction) -> Fut + Send,
1371 Fut: Future<Output = Result<R>> + Send,
1372 {
1373 let txn = self.new_transaction().await?;
1374 let commit_handle = txn.clone();
1375 let result = f(txn).await?;
1376 commit_handle.commit().await?;
1377 Ok(result)
1378 }
1379
1380 /// Insert an entry into the database without modifying or validating it.
1381 /// Primarily for testing / full control over raw entry storage.
1382 ///
1383 /// The entry is stored `Unverified`: this path runs no validation, so it
1384 /// cannot honestly claim the entry is verified, and the storage API no
1385 /// longer accepts a caller-asserted status. Only the local validation
1386 /// pass promotes entries to `Verified`.
1387 pub async fn insert_raw(&self, entry: Entry) -> Result<ID> {
1388 let instance = self.instance()?;
1389 let id = entry.id();
1390
1391 instance.put(entry).await?;
1392
1393 Ok(id)
1394 }
1395
1396 /// Get a Store type that will handle accesses to the Store
1397 /// This will return a Store initialized to point at the current state of the database.
1398 ///
1399 /// The returned store should NOT be used to modify the database, as it intentionally does not
1400 /// expose the Transaction. Since the Transaction is never committed, it does not have any
1401 /// effect on the database.
1402 pub async fn get_store_viewer<T>(&self, name: impl Into<String>) -> Result<T>
1403 where
1404 T: Store,
1405 {
1406 let txn = self.new_transaction().await?;
1407 T::load(&txn, name.into()).await
1408 }
1409
1410 /// Get the current tips (leaf entries) of the main database branch.
1411 ///
1412 /// Tips represent the latest entries in the database's main history, forming the heads of the DAG.
1413 ///
1414 /// If any raw tip is `Unverified`, an opportunistic [`Self::verify`] pass
1415 /// runs first (entries arrive `Unverified` from sync; this promotes the
1416 /// ones whose pinned `_settings` are now held).
1417 ///
1418 /// The returned tips then depend on the handle's view:
1419 ///
1420 /// - **default** — the **Verified frontier**: the tips of the maximal
1421 /// ancestor-closed, all-`Verified` prefix of the DAG. A still-`Unverified`
1422 /// tip is replaced by its nearest `Verified` ancestors; anything reachable
1423 /// only through an `Unverified` entry is excluded.
1424 /// - **[`allow_unverified`](Self::allow_unverified)** — the raw tips with
1425 /// only `Failed` entries dropped (`Unverified` tips are kept).
1426 ///
1427 /// `Failed` entries are dropped in both cases. While a [`verify`](Self::verify)
1428 /// pass is on the stack the frontier is bypassed (its own reads must see the
1429 /// raw DAG to reconstruct pinned `_settings`); a remote backend returns its
1430 /// raw tips unchanged (the server owns verification).
1431 ///
1432 /// # Returns
1433 /// A `Result` containing the [`Snapshot`] of tip entries or an error.
1434 pub async fn snapshot(&self) -> Result<Snapshot> {
1435 let instance = self.instance()?;
1436
1437 // On a remote instance the server owns verification: `snapshot`
1438 // already returns the server-side Verified frontier (or empty for a
1439 // not-yet-propagated tree, e.g. `Database::create`'s bootstrap
1440 // placeholder root — `EntryNotFound` is mapped to empty to match
1441 // `Backend::snapshot`'s contract). Return it directly: the local
1442 // verification machinery below (status probe, auto-verify,
1443 // `verified_frontier`) is local-only and would fail on a remote
1444 // backend anyway (e.g. `verified_frontier`'s `backend.get_tree(...)`).
1445 //
1446 // Delegate to `self.ops()` rather than calling the connection
1447 // directly: when this handle was built via `Database::create` or
1448 // `Database::open_remote` its `ops` is a `RemoteBackend` carrying the
1449 // *per-database* identity (the new tree's own member key, or the
1450 // caller's chosen identity), which the server's per-tree gate accepts.
1451 // Routing through `conn.session_identity()` here would instead use the
1452 // connection's (caller's) session pubkey, which is not a member of a
1453 // freshly-created tree and gets denied. A handle from `Database::open`
1454 // on a connected instance instead clones the instance's session
1455 // backend, keeping the session-identity semantics for that path.
1456 #[cfg(all(unix, feature = "service"))]
1457 if instance.remote_connection().is_some() {
1458 return match self.ops().snapshot(&self.root).await {
1459 Ok(snap) => Ok(snap),
1460 Err(e) if e.is_not_found() => Ok(Snapshot::EMPTY),
1461 Err(e) => Err(e),
1462 };
1463 }
1464
1465 // Local path: verification-status probing needs the concrete engine.
1466 let backend = instance.require_local_engine()?;
1467 let tips = self.ops().snapshot(&self.root).await?.into_tips();
1468
1469 // Verification status ops are local-only. On a remote backend the
1470 // server owns verification (and stores everything Unverified until
1471 // it verifies); the client returns verified tips unchanged.
1472 if let Some(first) = tips.first()
1473 && backend.get_verification_status(first).await.is_err()
1474 {
1475 return Ok(Snapshot::new(tips));
1476 }
1477
1478 // Access-time opportunistic verification: if any tip is still
1479 // Unverified, attempt to resolve it now. Best-effort — a failure or a
1480 // still-incomplete pin must not block the read. Suppressed while a
1481 // verify pass is already on the stack (its own reads land here).
1482 let tips = if auto_verify_suppressed() {
1483 tips
1484 } else {
1485 let mut any_unverified = false;
1486 for t in &tips {
1487 if backend
1488 .get_verification_status(t)
1489 .await
1490 .unwrap_or(VerificationStatus::Unverified)
1491 == VerificationStatus::Unverified
1492 {
1493 any_unverified = true;
1494 break;
1495 }
1496 }
1497 if any_unverified {
1498 // Boxed: this call closes a snapshot → verify →
1499 // validate_entry → delegation → get_settings → snapshot
1500 // async cycle; the box gives it a finite future size.
1501 let _ = Box::pin(self.verify()).await;
1502 self.ops().snapshot(&self.root).await?.into_tips()
1503 } else {
1504 tips
1505 }
1506 };
1507
1508 // Default view: cut to the Verified frontier. Suppressed while a
1509 // verify pass is on the stack — its reads must see the raw DAG to
1510 // reconstruct pinned `_settings` (the frontier filter itself depends
1511 // on verification status, which is exactly what verify is computing).
1512 if !self.allow_unverified && !auto_verify_suppressed() {
1513 return self.verified_frontier().await.map(Snapshot::new);
1514 }
1515
1516 // `allow_unverified` view: keep Unverified tips, drop only Failed.
1517 let mut visible = Vec::with_capacity(tips.len());
1518 for t in tips {
1519 if backend.get_verification_status(&t).await? != VerificationStatus::Failed {
1520 visible.push(t);
1521 }
1522 }
1523 Ok(Snapshot::new(visible))
1524 }
1525
1526 /// Compute the tips of the maximal all-`Verified` prefix of the DAG.
1527 ///
1528 /// An entry is in the prefix iff it is `Verified` **and** every one of its
1529 /// parents is in the prefix (the prefix is ancestor-closed). The frontier
1530 /// is the set of prefix entries that are not the parent of any other
1531 /// prefix entry — i.e. the tips of the verified subgraph.
1532 ///
1533 /// Returns an empty vector if the root itself is not `Verified` (nothing
1534 /// is observable in the default view until verification reaches the root).
1535 async fn verified_frontier(&self) -> Result<Vec<ID>> {
1536 let instance = self.instance()?;
1537 let backend = instance.require_local_engine()?;
1538
1539 // Topologically sorted (height then ID): every parent precedes its
1540 // children, so a single forward pass can decide prefix membership.
1541 let entries = backend.get_tree(self.root_id()).await?;
1542
1543 let mut in_prefix: std::collections::HashSet<ID> = std::collections::HashSet::new();
1544 let mut covered: std::collections::HashSet<ID> = std::collections::HashSet::new();
1545
1546 for e in &entries {
1547 let id = e.id();
1548 if backend.get_verification_status(&id).await? != VerificationStatus::Verified {
1549 continue;
1550 }
1551 let parents = e.parents().unwrap_or_default();
1552 if parents.iter().all(|p| in_prefix.contains(p)) {
1553 in_prefix.insert(id);
1554 // Every parent now has a verified child, so it is interior to
1555 // the prefix and cannot itself be a frontier tip.
1556 for p in parents {
1557 covered.insert(p);
1558 }
1559 }
1560 }
1561
1562 let frontier: Vec<ID> = entries
1563 .into_iter()
1564 .map(|e| e.id())
1565 .filter(|id| in_prefix.contains(id) && !covered.contains(id))
1566 .collect();
1567 Ok(frontier)
1568 }
1569
1570 /// Get the full `Entry` objects for the current tips of the main database branch.
1571 ///
1572 /// # Returns
1573 /// A `Result` containing a vector of the tip `Entry` objects or an error.
1574 pub async fn get_tip_entries(&self) -> Result<Vec<Entry>> {
1575 let instance = self.instance()?;
1576 let snapshot = self.snapshot().await?;
1577 let mut entries = Vec::new();
1578 for id in snapshot.tips() {
1579 entries.push(instance.get(id).await?);
1580 }
1581 Ok(entries)
1582 }
1583
1584 /// Get a single entry by ID from this database.
1585 ///
1586 /// This is the primary method for retrieving entries after commit operations.
1587 /// It provides safe, high-level access to entry data without exposing backend details.
1588 ///
1589 /// The method verifies that the entry belongs to this database by checking its root ID.
1590 /// If the entry exists but belongs to a different database, an error is returned.
1591 ///
1592 /// # Arguments
1593 /// * `entry_id` - The ID of the entry to retrieve (accepts anything that converts to ID/String)
1594 ///
1595 /// # Returns
1596 /// A `Result` containing the `Entry` or an error if not found or not part of this database
1597 ///
1598 /// # Example
1599 /// ```rust,no_run
1600 /// # use eidetica::*;
1601 /// # use eidetica::Instance;
1602 /// # use eidetica::backend::database::InMemory;
1603 /// # use eidetica::crdt::Doc;
1604 /// # #[tokio::main]
1605 /// # async fn main() -> Result<()> {
1606 /// # let (_instance, mut user) = Instance::create_backend(
1607 /// # Box::new(InMemory::new()),
1608 /// # NewUser::passwordless("test"),
1609 /// # ).await?;
1610 /// # let key_id = user.add_private_key(None).await?;
1611 /// # let tree = user.create_database(Doc::new(), &key_id).await?;
1612 /// # let txn = tree.new_transaction().await?;
1613 /// let entry_id = txn.commit().await?;
1614 /// let entry = tree.get_entry(&entry_id).await?; // Using &ID
1615 /// let entry = tree.get_entry(entry_id.clone()).await?; // Using ID
1616 /// println!("Entry signature: {:?}", entry.auth());
1617 /// # Ok(())
1618 /// # }
1619 /// ```
1620 pub async fn get_entry<I: Into<ID>>(&self, entry_id: I) -> Result<Entry> {
1621 let id = entry_id.into();
1622 // Route through `self.ops()` so handles built via `Database::create`
1623 // or `Database::open_remote` read with the per-DB identity from
1624 // their `RemoteBackend`. Going through `instance.get(id)` would
1625 // use the connection's login pubkey, which on a remote instance is
1626 // denied by the per-tree gate when the login key isn't a member of
1627 // this tree (e.g. user-tree key created via `User::add_private_key`
1628 // and used to author a database that doesn't grant the root key).
1629 let entry = self.ops().get(&id).await?;
1630
1631 // Check if the entry belongs to this database
1632 if !entry.in_tree(&self.root) {
1633 return Err(InstanceError::EntryNotInDatabase {
1634 entry_id: id,
1635 database_id: self.root.clone(),
1636 }
1637 .into());
1638 }
1639
1640 Ok(entry)
1641 }
1642
1643 /// Get multiple entries by ID efficiently.
1644 ///
1645 /// This method retrieves multiple entries more efficiently than multiple `get_entry()` calls
1646 /// by minimizing conversion overhead and pre-allocating the result vector.
1647 ///
1648 /// The method verifies that all entries belong to this database by checking their root IDs.
1649 /// If any entry exists but belongs to a different database, an error is returned.
1650 ///
1651 /// # Parameters
1652 /// * `entry_ids` - An iterable of entry IDs to retrieve
1653 ///
1654 /// # Returns
1655 /// A `Result` containing a vector of `Entry` objects or an error if any entry is not found or not part of this database
1656 ///
1657 /// # Example
1658 /// ```rust,no_run
1659 /// # use eidetica::*;
1660 /// # use eidetica::Instance;
1661 /// # use eidetica::backend::database::InMemory;
1662 /// # use eidetica::crdt::Doc;
1663 /// # #[tokio::main]
1664 /// # async fn main() -> Result<()> {
1665 /// # let (_instance, mut user) = Instance::create_backend(
1666 /// # Box::new(InMemory::new()),
1667 /// # NewUser::passwordless("test"),
1668 /// # ).await?;
1669 /// # let key_id = user.add_private_key(None).await?;
1670 /// # let tree = user.create_database(Doc::new(), &key_id).await?;
1671 /// let entry_ids = vec![ID::from_bytes("id1"), ID::from_bytes("id2")];
1672 /// let entries = tree.get_entries(entry_ids).await?;
1673 /// # Ok(())
1674 /// # }
1675 /// ```
1676 pub async fn get_entries<I, T>(&self, entry_ids: I) -> Result<Vec<Entry>>
1677 where
1678 I: IntoIterator<Item = T>,
1679 T: std::borrow::Borrow<ID>,
1680 {
1681 let ids: Vec<ID> = entry_ids.into_iter().map(|t| t.borrow().clone()).collect();
1682 let instance = self.instance()?;
1683 let mut entries = Vec::with_capacity(ids.len());
1684
1685 for id in ids {
1686 let entry = instance.get(&id).await?;
1687
1688 // Check if the entry belongs to this database
1689 if !entry.in_tree(&self.root) {
1690 return Err(InstanceError::EntryNotInDatabase {
1691 entry_id: id,
1692 database_id: self.root.clone(),
1693 }
1694 .into());
1695 }
1696
1697 entries.push(entry);
1698 }
1699
1700 Ok(entries)
1701 }
1702
1703 // === AUTHENTICATION HELPERS ===
1704
1705 /// Verify an entry's signature and authentication against the database's configuration that was valid at the time of entry creation.
1706 ///
1707 /// This method validates that:
1708 /// 1. The entry belongs to this database
1709 /// 2. The entry is properly signed with a key that was authorized in the database's authentication settings at the time the entry was created
1710 /// 3. The signature is cryptographically valid
1711 ///
1712 /// The method uses the entry's metadata to determine which authentication settings were active when the entry was signed,
1713 /// ensuring that entries remain valid even if keys are later revoked or settings change.
1714 ///
1715 /// # Arguments
1716 /// * `entry_id` - The ID of the entry to verify (accepts anything that converts to ID/String)
1717 ///
1718 /// # Returns
1719 /// A `Result` containing `true` if the entry is valid and properly authenticated, `false` if authentication fails
1720 ///
1721 /// # Errors
1722 /// Returns an error if:
1723 /// - The entry is not found
1724 /// - The entry does not belong to this database
1725 /// - The entry's metadata cannot be parsed
1726 /// - The historical authentication settings cannot be retrieved
1727 pub async fn verify_entry_signature<I: Into<ID>>(&self, entry_id: I) -> Result<bool> {
1728 let entry = self.get_entry(entry_id).await?;
1729
1730 // Validate against the `_settings` the entry pins, not current
1731 // settings — so a later key revocation cannot retroactively
1732 // invalidate (or validate) historical entries.
1733 match self.get_historical_settings_for_entry(&entry).await? {
1734 // We do not hold the pinned `_settings` set, so we cannot make a
1735 // verification decision: report not-verified rather than guess.
1736 PinnedSettings::Incomplete => Ok(false),
1737 PinnedSettings::Complete(auth_settings) => {
1738 let instance = self.instance()?;
1739 let mut validator = AuthValidator::new();
1740 validator
1741 .validate_entry(&entry, &auth_settings, Some(&instance))
1742 .await
1743 }
1744 }
1745 }
1746
1747 /// Get the permission level for this database's configured signing key.
1748 ///
1749 /// Returns the effective permission for the key that was configured when opening
1750 /// or creating this database. This uses the already-resolved identity stored in
1751 /// the database's `DatabaseKey`.
1752 ///
1753 /// # Returns
1754 /// The effective Permission for the configured signing key.
1755 ///
1756 /// # Errors
1757 /// Returns an error if:
1758 /// - No signing key is configured (database opened without authentication)
1759 /// - The database settings cannot be retrieved
1760 /// - The key is no longer valid in the current auth settings
1761 ///
1762 /// # Example
1763 /// ```rust,no_run
1764 /// # use eidetica::*;
1765 /// # use eidetica::crdt::Doc;
1766 /// # use eidetica::backend::database::InMemory;
1767 /// # use eidetica::auth::crypto::generate_keypair;
1768 /// # #[tokio::main]
1769 /// # async fn main() -> Result<()> {
1770 /// # let instance = Instance::open_backend(Box::new(InMemory::new())).await?;
1771 /// # let (signing_key, _public_key) = generate_keypair();
1772 /// # let database = Database::create(&instance, signing_key, Doc::new()).await?;
1773 /// // Check if the current key has Admin permission
1774 /// let permission = database.current_permission().await?;
1775 /// if permission.can_admin() {
1776 /// println!("Current key has Admin permission!");
1777 /// }
1778 /// # Ok(())
1779 /// # }
1780 /// ```
1781 pub async fn current_permission(&self) -> Result<Permission> {
1782 let key = self
1783 .key
1784 .as_ref()
1785 .ok_or(AuthError::InvalidAuthConfiguration {
1786 reason: "No signing key configured for this database".to_string(),
1787 })?;
1788 self.validate_key(key).await
1789 }
1790
1791 /// Reconstruct the `_settings` auth config an entry's signature is pinned
1792 /// to, from the `settings_snapshot` recorded in its signed metadata.
1793 ///
1794 /// Validation must run against the settings the entry pinned — not the
1795 /// current settings — so granting authority later cannot retroactively
1796 /// invalidate an entry that pinned less, and (once revocation lands)
1797 /// removals are handled on a separate, current-settings path.
1798 ///
1799 /// Returns [`PinnedSettings::Incomplete`] when this node does not hold the
1800 /// full pinned `_settings` ancestor set; the caller must then leave the
1801 /// entry `Unverified` rather than guess against whatever it does hold.
1802 async fn get_historical_settings_for_entry(&self, entry: &Entry) -> Result<PinnedSettings> {
1803 let instance = self.instance()?;
1804 let backend = instance.backend();
1805
1806 // The pin: `_settings` snapshot recorded in the entry's signed metadata.
1807 let settings_tips: Vec<ID> = match entry.metadata() {
1808 Some(raw) => match serde_json::from_slice::<crate::transaction::EntryMetadata>(raw) {
1809 Ok(md) => md.settings_snapshot.into_tips(),
1810 // Unparsable metadata ⇒ we cannot establish the pin.
1811 Err(_) => return Ok(PinnedSettings::Incomplete),
1812 },
1813 None => Vec::new(),
1814 };
1815
1816 // Resolve the effective `_settings` tips to validate against.
1817 let effective_tips: Vec<ID> = if settings_tips.is_empty() {
1818 if entry.in_subtree(SETTINGS) {
1819 // Genesis / bootstrap: no prior `_settings` exists, so the
1820 // entry is self-authorising — validate against the auth it
1821 // itself establishes (TOFU), mirroring how the transaction
1822 // validates initial database creation. Seeding the
1823 // reconstruction with the entry itself folds in its own
1824 // `_settings` contribution.
1825 vec![entry.id()]
1826 } else {
1827 // No auth context at all (no settings ever configured) —
1828 // mirrors the transaction path's "auth never configured" case.
1829 return Ok(PinnedSettings::Complete(AuthSettings::new()));
1830 }
1831 } else {
1832 settings_tips
1833 };
1834
1835 // Completeness: every pinned tip and its full `_settings` ancestor
1836 // closure must be present locally. `store_at` silently
1837 // skips absent ancestors, so an explicit walk is required — a missing
1838 // ancestor would otherwise yield a wrong (partial) auth config.
1839 let mut stack: Vec<ID> = effective_tips.clone();
1840 let mut seen: std::collections::HashSet<ID> = std::collections::HashSet::new();
1841 while let Some(id) = stack.pop() {
1842 if !seen.insert(id.clone()) {
1843 continue;
1844 }
1845 let Ok(e) = backend.get(&id).await else {
1846 return Ok(PinnedSettings::Incomplete);
1847 };
1848 // Walk both the `_settings` subtree DAG and the main parents that
1849 // carry it, so the closure can't be short-circuited.
1850 for p in e.subtree_parents(SETTINGS).unwrap_or_default() {
1851 stack.push(p);
1852 }
1853 for p in e.parents().unwrap_or_default() {
1854 stack.push(p);
1855 }
1856 }
1857
1858 // Reconstruct the merged `_settings` Doc as of the pinned tips.
1859 // Entries come back root-first; `_settings` is a system subtree and
1860 // is never encrypted, so deserialize directly.
1861 let effective_snapshot = Snapshot::from(effective_tips.clone());
1862 let entries = backend
1863 .store_at(self.root_id(), SETTINGS, &effective_snapshot)
1864 .await?;
1865 let mut settings_doc = Doc::default();
1866 for e in &entries {
1867 if let Ok(data) = e.data(SETTINGS) {
1868 let part: Doc = serde_json::from_slice(data)?;
1869 settings_doc = settings_doc.merge(&part)?;
1870 }
1871 }
1872
1873 let auth_settings = match settings_doc.get("auth") {
1874 Some(crate::crdt::doc::Value::Doc(auth_doc)) => auth_doc.clone().into(),
1875 _ => AuthSettings::new(),
1876 };
1877 Ok(PinnedSettings::Complete(auth_settings))
1878 }
1879
1880 /// Return the IDs of entries reachable from `post_tips` but not from
1881 /// `previous_tips`, in topological order (parents before children).
1882 ///
1883 /// This is the canonical way to enumerate the entries added between two
1884 /// cursors of this database — typically the cursors supplied by a
1885 /// [`WriteEvent`](crate::instance::WriteEvent)'s
1886 /// [`previous_tips()`](crate::instance::WriteEvent::previous_tips) and
1887 /// [`post_tips()`](crate::instance::WriteEvent::post_tips). Callbacks that
1888 /// only care that *something* changed can ignore this; callbacks that
1889 /// need to enumerate or fetch entry contents call this to expand the
1890 /// cursor advance into a concrete set of IDs.
1891 ///
1892 /// The walk is bounded by the cursor diff — cost is proportional to the
1893 /// number of entries *added* between the two cursors, not to the full
1894 /// DAG.
1895 ///
1896 /// # Behavior
1897 ///
1898 /// - If `previous_tips == post_tips`, returns an empty vector.
1899 /// - If `previous_tips` is empty, returns every ID reachable from
1900 /// `post_tips` (i.e. the full ancestor closure of those tips).
1901 /// - If `post_tips` references an entry that does not exist locally,
1902 /// returns an `EntryNotFound` error.
1903 /// - Verification status is **not** filtered, and event-driven callers
1904 /// are not exempt. `WriteEvent` is *triggered* only by Verified
1905 /// writes, but its cursors are raw DAG frontiers, so an `Unverified`
1906 /// or `Failed` entry sitting as a tip is inside the bracket and is
1907 /// enumerated here. Every caller that fetches or ingests these IDs
1908 /// should filter via `Backend::get_verification_status`.
1909 ///
1910 /// # Errors
1911 ///
1912 /// - `EntryNotFound` if any entry reachable from `post_tips` *above* the
1913 /// `previous_tips` frontier is missing locally. Entries at or below the
1914 /// frontier are never fetched, so a partial history there is tolerated —
1915 /// the conservative direction (we may over-report rather than
1916 /// under-report).
1917 pub async fn ids_added(
1918 &self,
1919 previous_tips: &Snapshot,
1920 post_tips: &Snapshot,
1921 ) -> Result<Vec<ID>> {
1922 use std::collections::{HashMap, HashSet, VecDeque};
1923
1924 // Set-equality on the canonical tip-sets: order- and
1925 // duplication-insensitive, so a cursor that hasn't advanced
1926 // short-circuits regardless of how its tips were ordered.
1927 if previous_tips == post_tips {
1928 return Ok(Vec::new());
1929 }
1930
1931 // `previous_tips` is the stop-frontier: everything reachable from it has
1932 // already been observed. Walk *backward from post_tips* and halt at the
1933 // first entry in that frontier, so cost is bounded by the cursor diff —
1934 // the entries added between the two cursors — not the full DAG. (Same
1935 // shape as `sync::utils::collect_ancestors_to_send`.)
1936 //
1937 // The callback cursor model always advances `previous_tips` as a
1938 // complete frontier (a cut), so halting at its members is exact. A
1939 // partial or stale `previous_tips` can only over-report — never
1940 // under-report, since any genuinely-new entry is reached before the walk
1941 // meets the frontier — which is the conservative direction.
1942 let boundary: HashSet<ID> = previous_tips.tips().iter().cloned().collect();
1943
1944 // Routes through `self.ops()` so handles built via
1945 // `Database::create` / `Database::open_remote` walk over the wire with
1946 // the per-DB identity. On a local instance this is a clone of the local
1947 // backend; on a connected instance each `get(id)` is a permission-checked
1948 // round-trip to the daemon, which is the security gate the cursor-only
1949 // push model relies on.
1950 let ops = self.ops();
1951
1952 // Guard against a **backward** bracket — `previous_tips` ahead of
1953 // `post_tips`. The boundary is then made of *descendants* of the
1954 // walk's starting point, so it is never reached and the walk descends
1955 // to the root, reporting the whole history as "added". On a connected
1956 // instance every step is a permission-checked round-trip, so the cost
1957 // is O(history) fetches rather than a merely-wrong answer.
1958 //
1959 // Detected via entry height, which increases monotonically from parent
1960 // to child: every ancestor of a `post_tips` entry has height at most
1961 // `max(post heights)`, so if every boundary tip sits strictly above
1962 // that, none of them can ever be reached by the walk. Bounded by the
1963 // tip counts (both small), and works identically on a local and a
1964 // connected instance — unlike the backend-level reachability
1965 // primitive, which a remote handle cannot call.
1966 //
1967 // Nothing was *added* moving backward, so the answer is empty.
1968 //
1969 // Best-effort: if any tip cannot be resolved we skip the guard and
1970 // fall through to the normal walk rather than failing the caller. This
1971 // only detects the strictly-backward case; forked/incomparable cursors
1972 // still over-report, which is the documented conservative direction.
1973 if !boundary.is_empty() {
1974 let mut max_post: Option<u64> = None;
1975 let mut min_prev: Option<u64> = None;
1976 let mut resolved = true;
1977 for id in post_tips.tips() {
1978 match ops.get(id).await {
1979 Ok(e) => max_post = Some(max_post.map_or(e.height(), |h| h.max(e.height()))),
1980 Err(_) => {
1981 resolved = false;
1982 break;
1983 }
1984 }
1985 }
1986 if resolved {
1987 for id in &boundary {
1988 match ops.get(id).await {
1989 Ok(e) => {
1990 min_prev = Some(min_prev.map_or(e.height(), |h| h.min(e.height())))
1991 }
1992 Err(_) => {
1993 resolved = false;
1994 break;
1995 }
1996 }
1997 }
1998 }
1999 if resolved
2000 && let (Some(post_h), Some(prev_h)) = (max_post, min_prev)
2001 && prev_h > post_h
2002 {
2003 return Ok(Vec::new());
2004 }
2005 }
2006
2007 // Walk parents from `post_tips`, collecting every ID until the boundary.
2008 // A non-boundary entry that's missing locally is a hard error — the
2009 // cursor references unknown state; boundary entries are never fetched.
2010 let mut added: HashMap<ID, Entry> = HashMap::new();
2011 let mut visited: HashSet<ID> = HashSet::new();
2012 let mut queue: VecDeque<ID> = post_tips.tips().iter().cloned().collect();
2013 while let Some(id) = queue.pop_front() {
2014 if boundary.contains(&id) || !visited.insert(id.clone()) {
2015 continue;
2016 }
2017 let entry = ops.get(&id).await?;
2018 for p in entry.parents().unwrap_or_default() {
2019 queue.push_back(p);
2020 }
2021 added.insert(id, entry);
2022 }
2023
2024 // 3. Topo sort the added set (Kahn's, scoped). Parents outside
2025 // the set are pre-observed boundary entries and don't count
2026 // toward in-degree.
2027 let mut in_degree: HashMap<ID, usize> = HashMap::with_capacity(added.len());
2028 let mut children: HashMap<ID, Vec<ID>> = HashMap::new();
2029 for (id, entry) in &added {
2030 let mut d = 0usize;
2031 for p in entry.parents().unwrap_or_default() {
2032 if added.contains_key(&p) {
2033 d += 1;
2034 children.entry(p).or_default().push(id.clone());
2035 }
2036 }
2037 in_degree.insert(id.clone(), d);
2038 }
2039 let mut topo_queue: VecDeque<ID> = in_degree
2040 .iter()
2041 .filter(|&(_, &d)| d == 0)
2042 .map(|(id, _)| id.clone())
2043 .collect();
2044 let mut order: Vec<ID> = Vec::with_capacity(added.len());
2045 while let Some(id) = topo_queue.pop_front() {
2046 if let Some(kids) = children.get(&id) {
2047 for kid in kids {
2048 let d = in_degree.get_mut(kid).expect("in_degree entry exists");
2049 *d -= 1;
2050 if *d == 0 {
2051 topo_queue.push_back(kid.clone());
2052 }
2053 }
2054 }
2055 order.push(id);
2056 }
2057
2058 Ok(order)
2059 }
2060
2061 /// Attempt to verify every `Unverified` entry in this database.
2062 ///
2063 /// For each `Unverified` entry, reconstruct the `_settings` it pins
2064 /// (see [`Self::get_historical_settings_for_entry`]) and validate its
2065 /// signature + permissions against that:
2066 ///
2067 /// - an ancestor is `Failed` → this entry is `Failed` too (quarantine
2068 /// propagates down the branch);
2069 /// - an ancestor is still `Unverified`, or not held locally yet (partial
2070 /// sync) → left `Unverified` (retried once the ancestor verifies / the
2071 /// missing entry arrives);
2072 /// - pinned `_settings` not fully held locally → left `Unverified`
2073 /// (a later pass retries once the set syncs in);
2074 /// - signature + permissions valid → promoted to `Verified`;
2075 /// - definitively invalid → marked `Failed` (dropped from reads).
2076 ///
2077 /// Verification is **prefix-closed**: an entry is `Verified` only if its
2078 /// entire ancestor history is `Verified`. It is therefore impossible for a
2079 /// tip to be `Verified` while one of its ancestors is not, which is what
2080 /// makes the Verified set ancestor-closed (see [`Self::allow_unverified`]).
2081 ///
2082 /// Already-`Verified` entries are never demoted by this pass. An explicit
2083 /// offline backend trust reset can clear all statuses before a rule upgrade.
2084 /// Local-only — verification is a per-node decision
2085 /// and is never delegated to a peer.
2086 pub async fn verify(&self) -> Result<VerifyReport> {
2087 self.verify_with_source(WriteSource::Remote, None).await
2088 }
2089
2090 pub(crate) async fn verify_with_source(
2091 &self,
2092 source: WriteSource,
2093 previous_tips: Option<Snapshot>,
2094 ) -> Result<VerifyReport> {
2095 use std::collections::{HashMap, HashSet, VecDeque};
2096
2097 let instance = self.instance()?;
2098
2099 // Hold the per-tree write lock for the whole pass + fire so the
2100 // promoted-batch event serialises against concurrent
2101 // `put_entry` fires on this tree. Callers must not hold this
2102 // lock when calling `verify`. (The two main callers,
2103 // `put_remote_entries` and the service `SubmitSignedEntry`
2104 // handler, release their own lock before calling verify.)
2105 let lock = instance.tree_lock(self.root_id());
2106 let _guard = lock.lock().await;
2107
2108 // Suppress the access-time auto-verify hook for the whole pass:
2109 // validation reads the database (delegation → settings → tips) and
2110 // must not recurse back into verification.
2111 let (report, any_promoted, fire_tips) = IN_VERIFY
2112 .scope(true, async move {
2113 let backend = instance.require_local_engine()?;
2114
2115 // 1. Outer boundary of the Unverified region: raw DAG
2116 // tips. (Failed/Verified tips both terminate the walk
2117 // below; no pre-filter needed.)
2118 let raw_tips = backend.snapshot(self.root_id()).await?;
2119
2120 // 2. Walk parents from those tips, collecting the
2121 // Unverified region. `Verified` and `Failed` entries
2122 // are the inner boundary — by prefix-closure their
2123 // ancestors are already settled. Demotions cascade
2124 // through `Instance::demote_to_unverified`, so a
2125 // `Verified` entry hiding an `Unverified` descendant
2126 // cannot occur and we never need to descend past
2127 // a Verified entry.
2128 let mut unverified: HashMap<ID, Entry> = HashMap::new();
2129 let mut visited: HashSet<ID> = HashSet::new();
2130 let mut queue: VecDeque<ID> = raw_tips.tips().iter().cloned().collect();
2131 while let Some(id) = queue.pop_front() {
2132 if !visited.insert(id.clone()) {
2133 continue;
2134 }
2135 let status = match backend.get_verification_status(&id).await {
2136 Ok(s) => s,
2137 // Tree-internal reference we don't hold yet
2138 // (partial sync): skip; descendants stay
2139 // Unverified and a later pass picks them up
2140 // once the parent arrives.
2141 Err(e) if e.is_not_found() => continue,
2142 Err(e) => return Err(e),
2143 };
2144 match status {
2145 VerificationStatus::Verified | VerificationStatus::Failed => continue,
2146 VerificationStatus::Unverified => {
2147 let entry = backend.get(&id).await?;
2148 for p in entry.parents().unwrap_or_default() {
2149 queue.push_back(p);
2150 }
2151 unverified.insert(id, entry);
2152 }
2153 }
2154 }
2155
2156 // 3. Topo-sort the Unverified subgraph (Kahn's,
2157 // scoped). Parents outside the set are pre-settled
2158 // boundary entries and don't contribute to
2159 // in-degree.
2160 let mut in_degree: HashMap<ID, usize> = HashMap::with_capacity(unverified.len());
2161 let mut children: HashMap<ID, Vec<ID>> = HashMap::new();
2162 for (id, entry) in &unverified {
2163 let mut d = 0usize;
2164 for p in entry.parents().unwrap_or_default() {
2165 if unverified.contains_key(&p) {
2166 d += 1;
2167 children.entry(p).or_default().push(id.clone());
2168 }
2169 }
2170 in_degree.insert(id.clone(), d);
2171 }
2172 let mut topo_queue: VecDeque<ID> = in_degree
2173 .iter()
2174 .filter(|&(_, &d)| d == 0)
2175 .map(|(id, _)| id.clone())
2176 .collect();
2177 let mut order: Vec<ID> = Vec::with_capacity(unverified.len());
2178 while let Some(id) = topo_queue.pop_front() {
2179 order.push(id.clone());
2180 if let Some(kids) = children.get(&id) {
2181 for kid in kids {
2182 let d = in_degree.get_mut(kid).expect("in_degree entry exists");
2183 *d -= 1;
2184 if *d == 0 {
2185 topo_queue.push_back(kid.clone());
2186 }
2187 }
2188 }
2189 }
2190
2191 // 4. Process parents-before-children. Same per-entry
2192 // logic as the legacy `get_tree`-walk verify, just
2193 // bounded to the Unverified region.
2194 let mut report = VerifyReport::default();
2195 let mut any_promoted = false;
2196 for id in &order {
2197 let entry = unverified.get(id).expect("topo id is in set");
2198 let parents = entry.parents().unwrap_or_default();
2199 let mut compromised = false;
2200 let mut blocked = false;
2201 for p in &parents {
2202 match backend.get_verification_status(p).await {
2203 Ok(VerificationStatus::Verified) => {}
2204 Ok(VerificationStatus::Failed) => compromised = true,
2205 Ok(VerificationStatus::Unverified) => blocked = true,
2206 Err(e) if e.is_not_found() => blocked = true,
2207 Err(e) => return Err(e),
2208 }
2209 }
2210 if compromised {
2211 backend
2212 .update_verification_status(id, VerificationStatus::Failed)
2213 .await?;
2214 report.failed += 1;
2215 continue;
2216 }
2217 if blocked {
2218 report.still_unverified += 1;
2219 continue;
2220 }
2221
2222 match self.get_historical_settings_for_entry(entry).await? {
2223 PinnedSettings::Incomplete => report.still_unverified += 1,
2224 PinnedSettings::Complete(auth_settings) => {
2225 let mut validator = AuthValidator::new();
2226 let valid = validator
2227 .validate_entry(entry, &auth_settings, Some(&instance))
2228 .await
2229 .unwrap_or(false);
2230 if valid {
2231 backend
2232 .update_verification_status(id, VerificationStatus::Verified)
2233 .await?;
2234 report.verified += 1;
2235 any_promoted = true;
2236 } else {
2237 backend
2238 .update_verification_status(id, VerificationStatus::Failed)
2239 .await?;
2240 report.failed += 1;
2241 }
2242 }
2243 }
2244 }
2245
2246 Ok::<_, crate::Error>((report, any_promoted, raw_tips))
2247 })
2248 .await?;
2249
2250 // Fire one batched `Verified` event for the promotion, if any.
2251 // Outside the IN_VERIFY scope so the callbacks' own reads aren't
2252 // suppressed. Still under the per-tree write lock — the lock is
2253 // dropped when this method returns.
2254 //
2255 // Pre- and post-pass tips are *the same* raw backend tips: verify
2256 // only mutates verification statuses, never the DAG structure.
2257 // For per-callback cursor advancement that means the cursor moves
2258 // to `pre_tips` (which includes the just-promoted entries), which
2259 // is the correct frontier — subscribers see the same
2260 // `previous_tips` they'd see for any subsequent fire until
2261 // something else writes to the tree.
2262 let joins = if any_promoted {
2263 let instance = self.instance()?;
2264 Some(
2265 instance
2266 .spawn_write_callbacks(
2267 self.root_id(),
2268 previous_tips.as_ref().unwrap_or(&fire_tips),
2269 &fire_tips,
2270 source,
2271 )
2272 .await,
2273 )
2274 } else {
2275 None
2276 };
2277
2278 // Release the per-tree lock before awaiting the user callbacks. The
2279 // cursor advances were already committed under the lock by
2280 // `spawn_write_callbacks` (which is what preserves event ordering);
2281 // the closures must run lock-free, because a callback that reads
2282 // tips can trip the access-time auto-verify hook → `verify()` →
2283 // `tree_lock` and would deadlock against a still-held `_guard`.
2284 drop(_guard);
2285
2286 if let Some(mut joins) = joins {
2287 while joins.join_next().await.is_some() {}
2288 }
2289
2290 Ok(report)
2291 }
2292
2293 // === DATABASE QUERIES ===
2294
2295 /// Get all entries in this database.
2296 ///
2297 /// ⚠️ **Warning**: This method loads all entries into memory. Use with caution on large databases.
2298 /// Consider using `snapshot()` or `get_tip_entries()` for more efficient access patterns.
2299 ///
2300 /// # Returns
2301 /// A `Result` containing a vector of all `Entry` objects in the database
2302 pub async fn get_all_entries(&self) -> Result<Vec<Entry>> {
2303 let instance = self.instance()?;
2304 instance.require_local_engine()?.get_tree(&self.root).await
2305 }
2306}