Development Documentation (main branch) - For stable release docs, see docs.rs/eidetica
Skip to main content

eidetica/service/
client.rs

1//! Remote connection client for the Eidetica service.
2//!
3//! `RemoteConnection` connects to an Eidetica service server and forwards
4//! storage operations as RPC calls. It backs the `RemoteBackend` implementation
5//! of the `Backend` seam and is not itself a `BackendImpl`.
6//!
7//! Authentication uses the client-side-signing flow described in the Service
8//! Architecture doc § Security Model: `RemoteConnection::trusted_login` drives
9//! the daemon's `TrustedLoginUser` / `TrustedLoginProve` challenge-response,
10//! decrypts the user's root signing key in-process, and signs the challenge
11//! locally. The daemon never sees the password or the plaintext signing key.
12//! After login, subsequent backend operations travel inside the `Authenticated`
13//! envelope and are dispatched against the user's identity; the daemon gates
14//! each one per-tree against the target database's auth settings.
15
16use std::collections::{HashMap, HashSet, VecDeque};
17use std::path::Path;
18use std::sync::atomic::{AtomicBool, Ordering};
19use std::sync::{Arc, RwLock, Weak};
20
21use tokio::io::{ReadHalf, WriteHalf};
22use tokio::net::UnixStream;
23use tokio::sync::{Mutex, Notify, mpsc, oneshot};
24
25use crate::auth::crypto::PrivateKey;
26use crate::auth::crypto::{PublicKey, create_challenge_response};
27use crate::auth::types::SigKey;
28use crate::backend::{
29    InstanceMetadata, RecordMutations, RecordPage, RecordRange, StoreStateRequest,
30};
31use crate::entry::{Entry, ID};
32use crate::instance::WeakInstance;
33use crate::service::error::service_error_to_eidetica_error;
34use crate::service::protocol::{
35    AuthenticatedDbRequest, DatabaseOp, Handshake, HandshakeAck, MergeState, Notification,
36    PROTOCOL_VERSION, ReadScope, ServerFrame, ServiceRequest, ServiceResponse, TransactionContext,
37    WireCrdtValue, read_frame, write_frame,
38};
39use crate::snapshot::Snapshot;
40use crate::user::UserError;
41use crate::user::crypto::{decrypt_private_key, derive_encryption_key};
42use crate::user::types::{KeyStorage, UserInfo};
43
44/// How long an `Idle` per-tree subscription is kept warm before the
45/// sweep sends `UnsubscribeWrites`. A re-registration arriving inside
46/// this window transitions back to `Subscribed` without a wire call.
47///
48/// Sized for a "user is briefly between two react renders" case rather
49/// than "user closes the app, comes back tomorrow." 60s is plenty for
50/// the churn case and small enough that abandoned subscriptions don't
51/// linger.
52const IDLE_GRACE_WINDOW: std::time::Duration = std::time::Duration::from_secs(60);
53
54/// How often the sweep task wakes up to check for expired `Idle`
55/// entries. Half the grace window so an entry that becomes Idle right
56/// after a sweep tick still gets unsubscribed within roughly one grace
57/// window's worth of clock time.
58const SWEEP_INTERVAL: std::time::Duration = std::time::Duration::from_secs(30);
59
60/// Returns the idle grace window.
61///
62/// Production callers get [`IDLE_GRACE_WINDOW`]. Under the `testing`
63/// feature, [`set_idle_grace_window_for_test`] can install an override
64/// so the lazy-unsubscribe sweep path is exercisable without waiting
65/// the full production grace window.
66fn idle_grace_window() -> std::time::Duration {
67    #[cfg(feature = "testing")]
68    {
69        if let Some(d) = TEST_IDLE_GRACE_WINDOW.get() {
70            return *d;
71        }
72    }
73    IDLE_GRACE_WINDOW
74}
75
76/// Returns the sweep interval.
77///
78/// Production callers get [`SWEEP_INTERVAL`]. Under the `testing`
79/// feature, [`set_sweep_interval_for_test`] can install an override
80/// so the sweep fires promptly in test suites.
81fn sweep_interval() -> std::time::Duration {
82    #[cfg(feature = "testing")]
83    {
84        if let Some(d) = TEST_SWEEP_INTERVAL.get() {
85            return *d;
86        }
87    }
88    SWEEP_INTERVAL
89}
90
91#[cfg(feature = "testing")]
92static TEST_IDLE_GRACE_WINDOW: std::sync::OnceLock<std::time::Duration> =
93    std::sync::OnceLock::new();
94
95#[cfg(feature = "testing")]
96static TEST_SWEEP_INTERVAL: std::sync::OnceLock<std::time::Duration> = std::sync::OnceLock::new();
97
98/// Override the idle grace window for tests.
99///
100/// Call once, before any connection that spawns the sweep task. The
101/// value is process-global (a `OnceLock`); subsequent calls are
102/// no-ops.
103#[cfg(feature = "testing")]
104pub fn set_idle_grace_window_for_test(d: std::time::Duration) {
105    let _ = TEST_IDLE_GRACE_WINDOW.set(d);
106}
107
108/// Override the sweep interval for tests.
109///
110/// Call once, before any connection that spawns the sweep task. The
111/// value is process-global (a `OnceLock`); subsequent calls are
112/// no-ops.
113#[cfg(feature = "testing")]
114pub fn set_sweep_interval_for_test(d: std::time::Duration) {
115    let _ = TEST_SWEEP_INTERVAL.set(d);
116}
117
118/// Per-connection session state, populated by `trusted_login` on success.
119///
120/// Holds only the public key the daemon verified during challenge-response.
121/// The plaintext signing key is intentionally **not** stored here — it lives
122/// in the `User::key_manager` session that owns this connection. The
123/// daemon-side `ConnectionState::Authenticated` is what carries the
124/// `user_uuid` (for chunk 6's cache scoping); the client doesn't need it.
125#[derive(Clone, Debug)]
126struct SessionState {
127    session_pubkey: PublicKey,
128}
129
130/// Per-tree subscription state for [`RemoteConnectionInner::subscribed_trees`].
131///
132/// State machine:
133/// ```text
134///     (absent)
135///        │
136///        │ first `on_write` for tree
137///        ▼
138///   InFlight(notify) ──leader wire failure──> (absent)
139///        │
140///        │ leader wire success
141///        ▼
142///    Subscribed ───drop last cb───> Idle { since: Instant }
143///        ▲                              │
144///        │ new `on_write` arrives        │ sweep determines past
145///        │ (no wire call)                │ grace window
146///        └──────────────────────────────┘
147///                                       │
148///                                       ▼
149///                              UnsubscribeWrites on wire → (absent)
150/// ```
151///
152/// `InFlight` carries a `Notify` whose waiters are released exactly
153/// once when the leader finishes the wire round-trip (success or
154/// failure). Followers re-check the map after waking — success
155/// transitions to `Subscribed` (they return `Ok`), failure removes
156/// the entry (one of them becomes the next leader on retry).
157/// Defensive against a future shape change: under the current
158/// [`RemoteConnectionInner::subscription_locks`] fence,
159/// [`RemoteConnection::subscribe_writes`] is per-tree-serialized and
160/// no caller can observe `InFlight` for a tree it's about to subscribe
161/// to. Kept so a future relaxation of the fence (e.g. per-connection
162/// rather than per-tree) doesn't silently re-introduce the
163/// concurrent-leader race the `Notify` originally guarded.
164///
165/// `Idle` records the moment the last local callback for this tree
166/// was dropped. The daemon-side subscription is still alive — we
167/// haven't sent `UnsubscribeWrites` — so a re-registration before the
168/// grace window expires can transition straight back to `Subscribed`
169/// without a wire round-trip. A periodic sweep task removes Idle
170/// entries that have been quiet long enough, sending
171/// `UnsubscribeWrites` to the daemon at that point under the same
172/// per-tree fence the subscribe path uses, so a sweep's Unsubscribe
173/// is fully acked before any racing re-subscribe can send its
174/// Subscribe — no daemon-side `Sub → Unsub` inversion possible.
175enum SubState {
176    InFlight(Arc<Notify>),
177    Subscribed {
178        identity: SigKey,
179    },
180    Idle {
181        since: std::time::Instant,
182        identity: SigKey,
183    },
184}
185
186/// Role assignment for one entry into [`RemoteConnection::subscribe_writes`].
187/// Decided under the `subscribed_trees` mutex and consumed outside it so the
188/// std::Mutex is never held across an `await`.
189enum SubRole {
190    /// This task owns the wire round-trip and must transition the state +
191    /// `notify_waiters` when it finishes (success or failure).
192    Leader(Arc<Notify>),
193    /// Another task is already subscribing; await its notify and re-check.
194    Follower(Arc<Notify>),
195}
196
197/// Internal state for a remote connection, wrapped in Arc for Clone.
198struct RemoteConnectionInner {
199    /// Owns the write half of the socket. `tokio::sync::Mutex` because
200    /// `write_frame` is async (held across awaits). Only ever held for the
201    /// duration of one frame's write plus the FIFO push into [`Self::pending`];
202    /// the await on the response itself happens *after* the lock is released
203    /// so concurrent callers don't serialise on read-side latency.
204    writer: Mutex<WriteHalf<UnixStream>>,
205    /// FIFO of awaiting response slots. `request()` pushes one before
206    /// releasing the writer lock; the reader task pops the front on every
207    /// `ServerFrame::Response` so request and response order line up. The
208    /// VecDeque is guarded by a plain `std::sync::Mutex` — never held
209    /// across an await — so it can't deadlock with the writer lock.
210    pending: std::sync::Mutex<VecDeque<oneshot::Sender<ServiceResponse>>>,
211    /// Set once by [`RemoteConnection::attach_instance`] right after
212    /// `Instance::connect` has finished building the Instance. The reader
213    /// task reads (cheap clone of the inner `Weak`) on each
214    /// [`Notification::DatabaseWrite`] to dispatch into the instance's
215    /// callback registry. Stays `None` until attach; notifications can't
216    /// arrive in that window because the client subscribes lazily on the
217    /// first `Database::on_write` registration, which itself can't run
218    /// until the Instance exists.
219    weak_instance: std::sync::Mutex<Option<WeakInstance>>,
220    /// Set on successful `trusted_login`; read by `backend_request` to populate
221    /// the `Authenticated` envelope's identity field. `RwLock` because reads
222    /// are far more frequent than the one-shot login write.
223    ///
224    /// Accessed poison-tolerantly via [`RemoteConnectionInner::session_read`]
225    /// and [`RemoteConnectionInner::session_write`]: a panic in one task
226    /// while holding the guard must not promote itself to a permanent connection
227    /// outage. The worst observable case is a half-written session field, which
228    /// the caller already treats as "unauthenticated" (`session_identity`
229    /// returns `None` and the per-tree gate rejects the op).
230    session: RwLock<Option<SessionState>>,
231    /// Pubkeys this client has already proven possession of on this
232    /// connection (via `SessionKeyChallenge`/`SessionKeyRegister`), plus the
233    /// login pubkey added in `trusted_login`. Lets `register_session_key`
234    /// short-circuit when the key has already been registered, avoiding
235    /// per-request wire chatter for the common case where a single per-DB
236    /// key is reused across many ops.
237    registered_keys: Mutex<HashSet<PublicKey>>,
238    /// Per-tree subscription state. Entries are inserted on first call to
239    /// [`RemoteConnection::subscribe_writes`] and never removed on success
240    /// (subscriptions live for the connection's lifetime; the daemon scrubs
241    /// them on disconnect).
242    ///
243    /// The two states coordinate concurrent registrations against the same
244    /// tree: exactly one task is the "leader" that drives the wire round-trip;
245    /// other tasks observe `InFlight(notify)`, await the notify, and re-check
246    /// state. On leader success the state transitions to `Subscribed` and
247    /// followers return `Ok`; on leader failure the entry is removed so the
248    /// next waker can take leadership and retry.
249    subscribed_trees: std::sync::Mutex<HashMap<ID, SubState>>,
250    /// Per-tree async mutexes that fence wire-subscription state
251    /// transitions on this connection: held across the full
252    /// `SubscribeWrites` / `UnsubscribeWrites` request-response by the
253    /// leader path of [`RemoteConnection::subscribe_writes`] and by the
254    /// lazy-unsubscribe sweep ([`run_sweep_task`]).
255    ///
256    /// Closes the latent sweep-vs-resubscribe race in the gap between
257    /// the sweep removing an `Idle` entry from `subscribed_trees` and
258    /// its `UnsubscribeWrites` reaching the daemon: a racing
259    /// `subscribe_writes` could observe `None`, send `SubscribeWrites`,
260    /// and — if the daemon processed Subscribe before the in-flight
261    /// Unsubscribe — end Sub → Unsub (silent broken delivery).
262    ///
263    /// **Why the race is currently latent**: the daemon's per-connection
264    /// request loop is serial today (`server.rs` SubscribeWrites
265    /// handler comment), so Subscribe queues behind in-flight
266    /// Unsubscribe and gets processed after — end Subscribed. The
267    /// fence is structural future-proofing against a daemon shape
268    /// change to per-connection parallel dispatch. Cheap to maintain
269    /// (one async mutex per active tree, held only across sweep and
270    /// subscribe wire RTTs) and removes the dependency on the daemon
271    /// invariant entirely. Holding this lock across the daemon's ack
272    /// means a `subscribe_writes` arriving while a sweep is in flight
273    /// on the same tree blocks until the daemon has fully processed
274    /// the unsubscribe, regardless of dispatch shape.
275    ///
276    /// **Correctness contract this fence depends on**: the daemon
277    /// must serialize `SubscribeWrites` / `UnsubscribeWrites` *per
278    /// tree* within a single connection. Today this holds trivially
279    /// via per-connection serial dispatch. A future shape change to
280    /// per-connection+tree-parallel dispatch (the natural next step,
281    /// mirroring the client's `tree_workers`) also satisfies the
282    /// contract: within tree X the daemon would still order
283    /// Unsubscribe → Subscribe, while unrelated work on tree Y
284    /// proceeds in parallel. The fence stays correct under that
285    /// shape with no further work.
286    ///
287    /// What would break the fence: a daemon that *parallelizes
288    /// requests within a single tree* on one connection, freely
289    /// reordering Subscribe/Unsubscribe processing for the same
290    /// `root_id`. That shape would also break verify, settled-state
291    /// cursor advancement, and other invariants — it's not a
292    /// realistic future direction. If it ever becomes one, this
293    /// fence is insufficient and the design needs to revisit
294    /// ack-then-Subscribe vs. an `Unsubscribing { notify }` sub-state.
295    ///
296    /// Deliberately a separate lock from [`crate::instance::Instance`]'s
297    /// `tree_lock`: that lock serializes local `put_entry`/`verify`
298    /// against callback-dispatch coherence; reusing it here would
299    /// stall local writes on the same tree for an Unsubscribe RTT
300    /// for no correctness benefit.
301    ///
302    /// Shape mirrors `Instance::tree_lock`: std mutex around a hashmap
303    /// of `Arc<tokio::sync::Mutex<()>>` so the per-tree guard can be
304    /// cloned out and held across awaits.
305    subscription_locks: std::sync::Mutex<HashMap<ID, Arc<Mutex<()>>>>,
306    /// Set to `true` when the reader task exits (clean EOF, socket error,
307    /// or deserialization failure). Once set, [`RemoteConnection::request`]
308    /// short-circuits with `ConnectionAborted` instead of pushing a fresh
309    /// oneshot that would never be matched.
310    ///
311    /// Required because dropping the user-visible `RemoteConnection` does
312    /// not tear down the inner Arc (the reader task holds its own clone);
313    /// post-reader-exit calls would otherwise queue a sender into
314    /// [`Self::pending`] and `await` indefinitely on a `recv()` that no
315    /// one can fulfil. `pending` is cleared on reader exit, but a fresh
316    /// request landing *after* the clear would push a new sender into
317    /// the now-orphan queue.
318    ///
319    /// Ordering: the reader task sets this with `Release` ordering before
320    /// clearing `pending`, so any `Acquire` load that observes `true` is
321    /// guaranteed to also observe the empty queue.
322    closed: AtomicBool,
323    /// Per-tree dispatch lanes. The reader routes each incoming
324    /// `Notification::DatabaseWrite` by `root_id` into the matching
325    /// tree's `mpsc<Notification>` (lazily creating one + spawning a
326    /// per-tree worker on first notification for the tree). Each
327    /// worker pulls from its own channel and `await`s
328    /// `Instance::fire_write_callbacks` sequentially — sequential
329    /// within a tree (cursor advancement is well-defined), concurrent
330    /// across trees (a slow callback on one tree doesn't stall any
331    /// other tree's dispatches on this connection).
332    ///
333    /// **Why per-tree, not per-connection.** User-callback work is
334    /// per-tree; cursor advancement is per-tree; the only ordering
335    /// constraint we actually need is per-tree. The previous
336    /// single-drain-task model serialised across trees and could
337    /// stall an entire connection on one slow callback.
338    ///
339    /// **Why the reader doesn't await inline.** User callbacks may
340    /// issue wire calls (e.g. `Database::open` on a connected
341    /// instance) whose responses land through *this same reader*.
342    /// Awaiting a callback inline would deadlock the reader against
343    /// the response it needs to deliver. Routing to a separate worker
344    /// task keeps the reader free.
345    ///
346    /// Each worker holds `Weak<RemoteConnectionInner>` so it doesn't
347    /// keep `inner` alive. When `inner` drops (every user-facing
348    /// `RemoteConnection` released *and* the reader has exited),
349    /// every sender in this map drops, each worker's `recv()` returns
350    /// `None`, and workers exit cleanly without prolonging
351    /// `inner`'s lifetime.
352    tree_workers: std::sync::Mutex<HashMap<ID, mpsc::UnboundedSender<Notification>>>,
353    /// Abort handle for the background reader task, so [`Self::mark_dead`] can
354    /// force the reader out when the connection is torn down against a *wedged*
355    /// daemon — one that neither answers nor closes the socket. Without it the
356    /// reader blocks in `read_frame` forever, holding its own `Arc<inner>` clone
357    /// (leaking the task + socket, since the field docs on [`Self::closed`] note
358    /// dropping the user-facing `RemoteConnection` does not tear down `inner`)
359    /// and re-spawning per-tree workers on the next notification after
360    /// `mark_dead` cleared them. Set once, immediately after the reader is
361    /// spawned in [`RemoteConnection::connect`]; `None` only in the brief window
362    /// before that assignment, and after `mark_dead` has taken it.
363    reader_abort: std::sync::Mutex<Option<tokio::task::AbortHandle>>,
364}
365
366impl RemoteConnectionInner {
367    /// Acquire a read guard on `session`, tolerating poisoning.
368    ///
369    /// See the field-level doc on [`Self::session`] for the recovery rationale.
370    fn session_read(&self) -> std::sync::RwLockReadGuard<'_, Option<SessionState>> {
371        self.session
372            .read()
373            .unwrap_or_else(|poisoned| poisoned.into_inner())
374    }
375
376    /// Acquire a write guard on `session`, tolerating poisoning.
377    fn session_write(&self) -> std::sync::RwLockWriteGuard<'_, Option<SessionState>> {
378        self.session
379            .write()
380            .unwrap_or_else(|poisoned| poisoned.into_inner())
381    }
382
383    /// Acquire the pending-queue lock, tolerating poisoning.
384    fn pending_lock(
385        &self,
386    ) -> std::sync::MutexGuard<'_, VecDeque<oneshot::Sender<ServiceResponse>>> {
387        self.pending
388            .lock()
389            .unwrap_or_else(|poisoned| poisoned.into_inner())
390    }
391
392    /// Get-or-insert the per-tree subscription mutex. The returned
393    /// `Arc` is cheap to clone; the caller takes `lock().await` on it
394    /// outside the std mutex guard.
395    ///
396    /// See the field-level doc on [`Self::subscription_locks`] for the
397    /// race this fence closes.
398    fn subscription_lock(&self, tree_id: &ID) -> Arc<Mutex<()>> {
399        let mut locks = self
400            .subscription_locks
401            .lock()
402            .unwrap_or_else(|p| p.into_inner());
403        Arc::clone(
404            locks
405                .entry(tree_id.clone())
406                .or_insert_with(|| Arc::new(Mutex::new(()))),
407        )
408    }
409
410    /// Mark this connection dead and drop all dependent state. Used by
411    /// the reader task on its own exit path and by the sweep when an
412    /// `UnsubscribeWrites` times out (a daemon that can't ack a trivial
413    /// hashmap removal in seconds is broken; tear down and let the
414    /// caller reconnect).
415    ///
416    /// 1. Mark `closed` with `Release` ordering so future `request()`
417    ///    calls short-circuit before pushing senders into an orphan
418    ///    queue.
419    /// 2. Drain `pending` — every awaiting caller sees `RecvError`
420    ///    and surfaces `ConnectionAborted`.
421    /// 3. Drop every per-tree worker sender so workers exit cleanly.
422    /// 4. Abort the reader task. On the reader's *own* exit path this is a
423    ///    no-op (the task is already returning). When the sweep calls this on a
424    ///    wedged daemon it forces the reader out of its blocking `read_frame`,
425    ///    dropping the reader's `Arc<inner>` clone so `inner` can finally drop,
426    ///    and stopping it from re-spawning the workers step 3 just cleared.
427    ///    The handle is `take`n so a later `mark_dead` is a no-op.
428    fn mark_dead(&self) {
429        self.closed.store(true, Ordering::Release);
430        self.pending_lock().clear();
431        self.tree_workers
432            .lock()
433            .unwrap_or_else(|p| p.into_inner())
434            .clear();
435        if let Some(handle) = self
436            .reader_abort
437            .lock()
438            .unwrap_or_else(|p| p.into_inner())
439            .take()
440        {
441            handle.abort();
442        }
443    }
444}
445
446/// A connection to a remote Eidetica service server over a Unix domain socket.
447///
448/// `RemoteConnection` backs the `RemoteBackend` implementation of the `Backend`
449/// seam. It provides the storage operations as inherent methods, plus additional
450/// coordination methods like `notify_entry_written`.
451///
452/// Cloning is cheap (Arc-backed).
453///
454/// **Teardown.** The reader task holds a strong `Arc<RemoteConnectionInner>`
455/// and `inner` owns the socket's `WriteHalf`, so neither half of the split
456/// stream can drop on its own: the reader parks in `read_frame` until the
457/// daemon EOFs, and the daemon only EOFs once our socket closes. The
458/// `live` token breaks that cycle — when the last user-facing handle drops
459/// it calls `RemoteConnectionInner::mark_dead`, aborting the reader so
460/// its `Arc` is released, `inner` drops, and the `WriteHalf` closes. The
461/// daemon then sees EOF and scrubs the connection's subscriptions via its
462/// `ConnectionGuard`.
463#[derive(Clone)]
464pub struct RemoteConnection {
465    inner: Arc<RemoteConnectionInner>,
466    /// `Some` on every user-facing handle; `None` on internal handles
467    /// minted from a task that must not keep the connection alive (the
468    /// sweep). Cloning a user handle clones the token, so teardown fires
469    /// only when the last one goes.
470    _live: Option<Arc<ConnLiveness>>,
471}
472
473/// Drop token wired to `RemoteConnection::_live`. Holds a strong `inner`
474/// so `mark_dead` is always callable; that `Arc` is released immediately
475/// after, as part of this drop.
476struct ConnLiveness(Arc<RemoteConnectionInner>);
477
478impl Drop for ConnLiveness {
479    fn drop(&mut self) {
480        self.0.mark_dead();
481    }
482}
483
484impl std::fmt::Debug for RemoteConnection {
485    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
486        f.debug_struct("RemoteConnection").finish_non_exhaustive()
487    }
488}
489
490impl RemoteConnection {
491    /// Connect to a service server at the given socket path.
492    ///
493    /// Performs the protocol handshake, then spawns a background reader
494    /// task that demuxes [`ServerFrame`]s: `Response` frames pop the next
495    /// pending oneshot in FIFO order, `Notification` frames dispatch into
496    /// the attached `Instance`'s callback registry (after
497    /// [`Self::attach_instance`] has been called).
498    pub async fn connect(path: impl AsRef<Path>) -> crate::Result<Self> {
499        let stream = UnixStream::connect(path.as_ref()).await?;
500        let (mut reader, mut writer) = tokio::io::split(stream);
501
502        // Send handshake
503        let handshake = Handshake {
504            protocol_version: PROTOCOL_VERSION,
505        };
506        write_frame(&mut writer, &handshake).await?;
507
508        // Read ack
509        let ack: HandshakeAck = read_frame(&mut reader).await?.ok_or_else(|| {
510            crate::Error::Io(std::io::Error::new(
511                std::io::ErrorKind::ConnectionAborted,
512                "Server closed connection during handshake",
513            ))
514        })?;
515
516        if ack.protocol_version != PROTOCOL_VERSION {
517            return Err(crate::Error::Io(std::io::Error::new(
518                std::io::ErrorKind::InvalidData,
519                format!(
520                    "Protocol version mismatch: client={}, server={}",
521                    PROTOCOL_VERSION, ack.protocol_version
522                ),
523            )));
524        }
525
526        let inner = Arc::new(RemoteConnectionInner {
527            writer: Mutex::new(writer),
528            pending: std::sync::Mutex::new(VecDeque::new()),
529            weak_instance: std::sync::Mutex::new(None),
530            session: RwLock::new(None),
531            registered_keys: Mutex::new(HashSet::new()),
532            subscribed_trees: std::sync::Mutex::new(HashMap::new()),
533            subscription_locks: std::sync::Mutex::new(HashMap::new()),
534            closed: AtomicBool::new(false),
535            tree_workers: std::sync::Mutex::new(HashMap::new()),
536            reader_abort: std::sync::Mutex::new(None),
537        });
538
539        // Spawn the reader task. It holds an Arc clone of `inner` so the
540        // connection (and its pending queue) stay live as long as any
541        // request is in flight, and exits cleanly on EOF / read error /
542        // failure-to-deserialize, or on the abort the `live` token fires
543        // when the last user handle drops. On exit it drops the remaining oneshot
544        // senders (surfaces as `RecvError` on awaiting `request()`s) and
545        // also drops every per-tree worker channel, which causes those
546        // workers to exit. No separate dispatch task — per-tree workers
547        // are spawned lazily by the reader on first notification per
548        // tree.
549        let inner_for_reader = inner.clone();
550        let reader_handle = tokio::spawn(run_reader_task(reader, inner_for_reader));
551        *inner.reader_abort.lock().unwrap_or_else(|p| p.into_inner()) =
552            Some(reader_handle.abort_handle());
553
554        // Spawn the lazy-unsubscribe sweep. It holds a `Weak<inner>` so
555        // it doesn't extend `inner`'s lifetime; exits when `weak.upgrade()`
556        // returns `None` (the connection is being torn down).
557        let weak_for_sweep = Arc::downgrade(&inner);
558        tokio::spawn(run_sweep_task(weak_for_sweep));
559
560        let _live = Some(Arc::new(ConnLiveness(inner.clone())));
561        Ok(Self { inner, _live })
562    }
563
564    /// Attach an `Instance` to this connection so the reader task can
565    /// dispatch incoming [`Notification::DatabaseWrite`]s into the
566    /// instance's callback registry. Called exactly once by
567    /// `Instance::connect` after the Instance has been constructed.
568    /// Subsequent calls overwrite the previous reference, but no caller
569    /// does that today.
570    pub(crate) fn attach_instance(&self, weak: WeakInstance) {
571        *self
572            .inner
573            .weak_instance
574            .lock()
575            .unwrap_or_else(|poisoned| poisoned.into_inner()) = Some(weak);
576    }
577
578    /// Send a request and await its response.
579    ///
580    /// Take the writer lock; push a oneshot into the pending FIFO; write
581    /// the frame; release the lock; await the oneshot. Pushing the
582    /// oneshot *while still holding the writer lock* guarantees the FIFO
583    /// order in `pending` lines up with the on-wire order so the reader
584    /// task pairs each `ServerFrame::Response` with the right caller.
585    /// Concurrent `request()` calls do not serialise on the response
586    /// wait — only on the (cheap) frame write.
587    ///
588    /// Two `closed` checks gate the path: a cheap `Acquire` load before
589    /// acquiring the writer lock (the common-case fast path) and a second
590    /// re-check inside the lock *after* pushing the oneshot, in case the
591    /// reader exited concurrently between the first check and the push.
592    /// The reader sets `closed` with `Release` ordering *before* clearing
593    /// `pending`, so a load that sees `true` is guaranteed to see the
594    /// empty (or about-to-be-empty) queue. Without the post-push check,
595    /// a fresh request landing right after the reader clears could push
596    /// a sender into the orphan queue and `rx.await` forever.
597    async fn request(&self, req: ServiceRequest) -> crate::Result<ServiceResponse> {
598        if self.inner.closed.load(Ordering::Acquire) {
599            return Err(connection_aborted());
600        }
601        let (tx, rx) = oneshot::channel::<ServiceResponse>();
602        {
603            let mut writer = self.inner.writer.lock().await;
604            // Re-check under the writer lock to close the race where the
605            // reader exits and clears `pending` between our pre-check and
606            // this point. If we observe `closed` now, drop our oneshot on
607            // the floor without pushing — no reader means no response.
608            if self.inner.closed.load(Ordering::Acquire) {
609                return Err(connection_aborted());
610            }
611            // Push *before* writing the frame so the FIFO is consistent
612            // with the order frames hit the wire. If the write fails we
613            // pop the just-pushed sender so a future caller doesn't get
614            // matched to a response that never comes.
615            self.inner.pending_lock().push_back(tx);
616            if let Err(e) = write_frame(&mut *writer, &req).await {
617                let _ = self.inner.pending_lock().pop_back();
618                return Err(e);
619            }
620        }
621        rx.await.map_err(|_| connection_aborted())
622    }
623
624    /// Send a request and convert error responses to `crate::Error`.
625    pub(crate) async fn request_ok(&self, req: ServiceRequest) -> crate::Result<ServiceResponse> {
626        let resp = self.request(req).await?;
627        match resp {
628            ServiceResponse::Error(e) => Err(service_error_to_eidetica_error(e)),
629            other => Ok(other),
630        }
631    }
632
633    /// Wrap a `DatabaseOp` in the `AuthenticatedDb` envelope and send it.
634    ///
635    /// `(root_id, identity)` scope and `request_ok` error conversion, carrying
636    /// a `DatabaseOp` in an `AuthenticatedDbRequest`.
637    async fn db_request(
638        &self,
639        root_id: ID,
640        identity: SigKey,
641        op: DatabaseOp,
642    ) -> crate::Result<ServiceResponse> {
643        self.request_ok(ServiceRequest::AuthenticatedDb(Box::new(
644            AuthenticatedDbRequest {
645                root_id,
646                identity,
647                op,
648            },
649        )))
650        .await
651    }
652
653    /// Authenticate this connection as `username` by completing the
654    /// `TrustedLogin*` handshake against the daemon.
655    ///
656    /// Flow: send `TrustedLoginUser` → receive challenge + the user's full
657    /// `UserInfo` (encrypted credentials, user-database id, status) → derive
658    /// the password-encryption key locally (Argon2id) and decrypt the root
659    /// signing key in-process (or take it raw for passwordless users) → sign
660    /// the challenge → send `TrustedLoginProve` → expect `TrustedLoginOk`.
661    ///
662    /// The daemon never sees the password or the plaintext signing key; the
663    /// trust model for shipping the encrypted blob over the socket is captured
664    /// in the Service Architecture doc § Trusted login threat model.
665    ///
666    /// On success the connection's server-side state is `Authenticated`. The
667    /// caller receives the user's record and the decrypted root key so it can
668    /// build the `User` session without a second wire read of `_users` —
669    /// reads through the wire always travel as the authenticated user, which
670    /// with the per-tree gate means a fresh user without permissions on
671    /// `_users` would not be able to re-fetch it.
672    pub(crate) async fn trusted_login(
673        &self,
674        username: &str,
675        password: Option<&str>,
676    ) -> crate::Result<(String, UserInfo, PrivateKey)> {
677        // Step 1: name the user, receive challenge + user record.
678        let resp = self
679            .request_ok(ServiceRequest::TrustedLoginUser {
680                username: username.to_string(),
681            })
682            .await?;
683        let (challenge, user_uuid, user_info) = match resp {
684            ServiceResponse::TrustedLoginChallenge {
685                challenge,
686                user_uuid,
687                user_info,
688            } => (challenge, user_uuid, user_info),
689            other => return Err(unexpected_response("TrustedLoginChallenge", &other)),
690        };
691
692        // Step 2: decrypt the root signing key locally. Cross-check that the
693        // caller's password/no-password matches the credential's salt/no-salt;
694        // a mismatch is the same UX-level error as a wrong password.
695        let credentials = &user_info.credentials;
696        let is_passwordless = credentials.password_salt.is_none();
697        let signing_key = match (&credentials.root_key, password, is_passwordless) {
698            (KeyStorage::Unencrypted { key }, None, true) => key.clone(),
699            (
700                KeyStorage::Encrypted {
701                    ciphertext, nonce, ..
702                },
703                Some(pwd),
704                false,
705            ) => {
706                let salt = credentials.password_salt.as_deref().ok_or_else(|| {
707                    UserError::PasswordRequired {
708                        operation: "decrypt root key for remote login".to_string(),
709                    }
710                })?;
711                let kek = derive_encryption_key(pwd, salt)?;
712                decrypt_private_key(ciphertext, nonce, &kek)?
713            }
714            _ => return Err(UserError::InvalidPassword.into()),
715        };
716
717        // Step 3: sign the challenge and send the proof.
718        let signature = create_challenge_response(&challenge, &signing_key);
719        let resp = self
720            .request_ok(ServiceRequest::TrustedLoginProve { signature })
721            .await?;
722        match resp {
723            ServiceResponse::TrustedLoginOk => {
724                // Stash the verified session pubkey so subsequent
725                // `backend_request` calls can populate the `Authenticated`
726                // envelope's identity field.
727                *self.inner.session_write() = Some(SessionState {
728                    session_pubkey: credentials.root_key_id.clone(),
729                });
730                // The login pubkey is in the server-side session keyset by
731                // construction (the server seeds it there in
732                // `handle_trusted_login_prove`). Mirror that here so
733                // `register_session_key` short-circuits without a wire
734                // round-trip when called for the login key.
735                self.inner
736                    .registered_keys
737                    .lock()
738                    .await
739                    .insert(credentials.root_key_id.clone());
740                Ok((user_uuid, user_info, signing_key))
741            }
742            other => Err(unexpected_response("TrustedLoginOk", &other)),
743        }
744    }
745
746    /// Prove possession of `signing_key` and add its public key to the
747    /// connection's session keyset.
748    ///
749    /// Used by every `Database` handle whose `RemoteBackend` carries a
750    /// per-database identity (e.g. `Database::create` on a connected
751    /// instance, or `user.open_database_with_key` over the wire): the daemon
752    /// gates reads against the *acting* pubkey from the identity hint, and
753    /// the acting pubkey must be in the keyset, so we register the per-DB
754    /// key before the first read.
755    ///
756    /// Idempotent and cheap on repeated calls: a successful registration
757    /// caches the pubkey in `registered_keys`, and a follow-up call with the
758    /// same key returns `Ok(())` without touching the wire. The login pubkey
759    /// is seeded into the cache by `trusted_login`.
760    ///
761    /// Cryptographically a two-step proof of possession:
762    /// 1. `SessionKeyChallenge { pubkey }` → server returns a single-use,
763    ///    pubkey-bound random challenge.
764    /// 2. Client signs the challenge with `signing_key`; `SessionKeyRegister
765    ///    { pubkey, signature }` → server verifies and inserts the pubkey
766    ///    into its `session_keyset`.
767    pub(crate) async fn register_session_key(&self, signing_key: &PrivateKey) -> crate::Result<()> {
768        let pubkey = signing_key.public_key();
769        {
770            let cache = self.inner.registered_keys.lock().await;
771            if cache.contains(&pubkey) {
772                return Ok(());
773            }
774        }
775        // Step 1: ask for a challenge bound to this pubkey.
776        let resp = self
777            .request_ok(ServiceRequest::SessionKeyChallenge {
778                pubkey: pubkey.clone(),
779            })
780            .await?;
781        let challenge = match resp {
782            ServiceResponse::SessionKeyChallenge { challenge } => challenge,
783            other => return Err(unexpected_response("SessionKeyChallenge", &other)),
784        };
785        // Step 2: sign and submit. The daemon verifies and joins the pubkey
786        // into the connection's keyset on Ok.
787        let signature = create_challenge_response(&challenge, signing_key);
788        let resp = self
789            .request_ok(ServiceRequest::SessionKeyRegister {
790                pubkey: pubkey.clone(),
791                signature,
792            })
793            .await?;
794        Self::expect_ok(resp)?;
795        self.inner.registered_keys.lock().await.insert(pubkey);
796        Ok(())
797    }
798
799    pub(crate) async fn database_ticket(
800        &self,
801        tree_id: &ID,
802        identity: SigKey,
803    ) -> crate::Result<crate::sync::DatabaseTicket> {
804        let response = self
805            .db_request(tree_id.clone(), identity, DatabaseOp::CreateTicket)
806            .await?;
807        match response {
808            ServiceResponse::DatabaseTicket(ticket) => Ok(ticket),
809            other => Err(unexpected_response("DatabaseTicket", &other)),
810        }
811    }
812
813    // === Response extraction helpers ===
814
815    fn expect_ok(resp: ServiceResponse) -> crate::Result<()> {
816        match resp {
817            ServiceResponse::Ok => Ok(()),
818            other => Err(unexpected_response("Ok", &other)),
819        }
820    }
821
822    // === Instance-level operations ===
823
824    /// Build a `SigKey` from the session pubkey, when logged in.
825    pub fn session_identity(&self) -> Option<SigKey> {
826        self.inner
827            .session_read()
828            .as_ref()
829            .map(|s| SigKey::from_pubkey(&s.session_pubkey))
830    }
831
832    pub async fn get_instance_metadata(&self) -> crate::Result<Option<InstanceMetadata>> {
833        let resp = self.request_ok(ServiceRequest::GetInstanceMetadata).await?;
834        match resp {
835            ServiceResponse::InstanceMetadata(meta) => Ok(meta),
836            other => Err(unexpected_response("InstanceMetadata", &other)),
837        }
838    }
839
840    pub async fn set_instance_metadata(&self, metadata: &InstanceMetadata) -> crate::Result<()> {
841        // Gated server-side as Admin on `_databases`, not on `root_id`, so the
842        // scope's `root_id` is unused for this op (default is fine).
843        let identity = self.session_identity().unwrap_or_default();
844        let resp = self
845            .db_request(
846                ID::default(),
847                identity,
848                DatabaseOp::SetInstanceMetadata {
849                    metadata: Box::new(metadata.clone()),
850                },
851            )
852            .await?;
853        Self::expect_ok(resp)
854    }
855
856    /// Subscribe this connection to write notifications for `tree_id`.
857    ///
858    /// Safe to call concurrently from multiple tasks for the same
859    /// `tree_id`: serialized per-tree via
860    /// [`RemoteConnectionInner::subscription_locks`] so only one task
861    /// at a time decides state + runs the wire round-trip. The fence
862    /// also serializes against the lazy-unsubscribe sweep — if a
863    /// sweep is mid-`UnsubscribeWrites` for this tree, this call
864    /// blocks until the daemon has fully acked the unsubscribe, then
865    /// observes `None` and sends a fresh `SubscribeWrites`. Wire
866    /// order Unsubscribe → Subscribe is preserved end-to-end,
867    /// independent of the daemon's request-dispatch shape.
868    ///
869    /// Returns `Ok` only *after* the daemon has registered the
870    /// subscription, so an immediately-following commit on this
871    /// connection cannot race the subscribe and lose its
872    /// notification. On failure the local state is rolled back so a
873    /// subsequent call can retry.
874    ///
875    /// Idempotent across calls: a tree that is already `Subscribed`
876    /// returns `Ok` without any wire activity. The server's
877    /// `ConnectionRegistry` is also idempotent. Identity is gated
878    /// server-side as Read on `tree_id`.
879    pub(crate) async fn subscribe_writes(
880        &self,
881        tree_id: ID,
882        identity: SigKey,
883        tips: Snapshot,
884    ) -> crate::Result<()> {
885        // Per-tree subscription fence. Held across the entire state
886        // decision + (leader path's) wire round-trip so a concurrent
887        // sweep that's mid-`UnsubscribeWrites` for this tree can't
888        // interleave its frame between our state read and our Subscribe.
889        // See `RemoteConnectionInner::subscription_locks` for the race
890        // this closes.
891        let sub_lock = self.inner.subscription_lock(&tree_id);
892        let _sub_guard = sub_lock.lock().await;
893        loop {
894            // Decide our role under the std::Mutex without holding it across
895            // any await: leader inserts `InFlight(notify)` and proceeds to the
896            // wire call; followers clone the notify and `await` it below.
897            // `Idle` re-entries transition straight to `Subscribed` without
898            // a wire round-trip — the daemon-side subscription is still
899            // alive (the sweep hasn't unsubscribed it yet).
900            let role = {
901                let mut subs = self.subscribed_trees_lock();
902                match subs.get(&tree_id) {
903                    // Already subscribed: a no-op. A new `identity` on this
904                    // re-call is intentionally ignored — the daemon's live
905                    // subscription is keyed off the pubkey the original
906                    // subscribe succeeded with (same rationale as the `Idle`
907                    // arm below).
908                    Some(SubState::Subscribed { .. }) => return Ok(()),
909                    Some(SubState::Idle {
910                        identity: existing_identity,
911                        ..
912                    }) => {
913                        // Daemon-side subscription is still live; re-mark
914                        // as Subscribed locally without a wire call. Reuse
915                        // the identity the original subscribe succeeded
916                        // with — the daemon's subscription is keyed off
917                        // that pubkey, not whatever this caller now holds.
918                        let existing_identity = existing_identity.clone();
919                        subs.insert(
920                            tree_id.clone(),
921                            SubState::Subscribed {
922                                identity: existing_identity,
923                            },
924                        );
925                        return Ok(());
926                    }
927                    Some(SubState::InFlight(n)) => SubRole::Follower(n.clone()),
928                    None => {
929                        let n = Arc::new(Notify::new());
930                        subs.insert(tree_id.clone(), SubState::InFlight(n.clone()));
931                        SubRole::Leader(n)
932                    }
933                }
934            };
935
936            match role {
937                SubRole::Follower(notify) => {
938                    notify.notified().await;
939                    // Re-check: success → return Ok; failure (entry removed)
940                    // → loop, where this task may become the next leader.
941                    continue;
942                }
943                SubRole::Leader(notify) => {
944                    let result = self
945                        .db_request(
946                            tree_id.clone(),
947                            identity.clone(),
948                            DatabaseOp::SubscribeWrites { tips: tips.clone() },
949                        )
950                        .await
951                        .and_then(Self::expect_ok);
952                    {
953                        let mut subs = self.subscribed_trees_lock();
954                        match &result {
955                            Ok(()) => {
956                                subs.insert(
957                                    tree_id.clone(),
958                                    SubState::Subscribed {
959                                        identity: identity.clone(),
960                                    },
961                                );
962                            }
963                            Err(_) => {
964                                subs.remove(&tree_id);
965                            }
966                        }
967                    }
968                    // Wake everyone exactly once; new waiters that arrive
969                    // after this point land in the post-transition state.
970                    notify.notify_waiters();
971                    return result;
972                }
973            }
974        }
975    }
976
977    fn subscribed_trees_lock(&self) -> std::sync::MutexGuard<'_, HashMap<ID, SubState>> {
978        self.inner
979            .subscribed_trees
980            .lock()
981            .unwrap_or_else(|p| p.into_inner())
982    }
983
984    /// Whether the attached `Instance` still holds a per-database write
985    /// callback for `tree_id`.
986    ///
987    /// `false` when no instance is attached or it has been dropped —
988    /// nothing can be observing writes in either case, so the caller is
989    /// free to release the subscription.
990    ///
991    /// Lock order: callers hold `subscribed_trees` across this, which
992    /// takes the instance's `write_callbacks`. Nothing acquires those in
993    /// the opposite order — the registry guard in `register_write_callback`,
994    /// `remove_write_callback` and `spawn_write_callbacks` is released
995    /// before any subscription-state access.
996    fn tree_has_local_callbacks(&self, tree_id: &ID) -> bool {
997        let weak = {
998            let guard = self
999                .inner
1000                .weak_instance
1001                .lock()
1002                .unwrap_or_else(|p| p.into_inner());
1003            guard.clone()
1004        };
1005        weak.and_then(|w| w.upgrade())
1006            .is_some_and(|instance| instance.has_write_callbacks(tree_id))
1007    }
1008
1009    /// Transition the subscription state for `tree_id` from `Subscribed`
1010    /// to `Idle`. Called from `WriteCallback::drop` when the last local
1011    /// callback for a tree on this connection is released.
1012    ///
1013    /// No wire call: the daemon-side subscription stays alive through
1014    /// the Idle grace window. If a new `on_write` registration arrives
1015    /// before the sweep, [`Self::subscribe_writes`] transitions back to
1016    /// `Subscribed` without touching the wire.
1017    ///
1018    /// If the state isn't `Subscribed` at the moment of the call —
1019    /// e.g. a concurrent re-registration already raced us, or the
1020    /// sweep already unsubscribed — this is a no-op.
1021    ///
1022    /// **Re-check under the state lock.** `WriteCallback::drop` removes
1023    /// the callback from the instance registry and calls this as two
1024    /// separate steps, so a registration can land in between: it inserts
1025    /// into the registry, then `subscribe_writes` observes `Subscribed`
1026    /// and returns without wire traffic. Flipping to `Idle` unconditionally
1027    /// would strand that live callback on an `Idle` subscription, and the
1028    /// sweep would silently unsubscribe it a grace window later. Probing
1029    /// the registry while holding `subscribed_trees` closes the window in
1030    /// both directions: a registration whose insert is already visible
1031    /// keeps us `Subscribed`, and one that isn't visible yet must take
1032    /// this same lock in `subscribe_writes`, where it observes `Idle` and
1033    /// transitions back.
1034    pub(crate) fn transition_to_idle(&self, tree_id: &ID) {
1035        let mut subs = self.subscribed_trees_lock();
1036        if self.tree_has_local_callbacks(tree_id) {
1037            return;
1038        }
1039        if let Some(SubState::Subscribed { identity }) = subs.get(tree_id) {
1040            let identity = identity.clone();
1041            subs.insert(
1042                tree_id.clone(),
1043                SubState::Idle {
1044                    since: std::time::Instant::now(),
1045                    identity,
1046                },
1047            );
1048        }
1049    }
1050
1051    /// Send `UnsubscribeWrites` to the daemon for `tree_id`. Called by
1052    /// the sweep task when an `Idle` entry's grace window has expired.
1053    pub(crate) async fn unsubscribe_writes(
1054        &self,
1055        tree_id: ID,
1056        identity: SigKey,
1057    ) -> crate::Result<()> {
1058        self.db_request(tree_id, identity, DatabaseOp::UnsubscribeWrites)
1059            .await
1060            .and_then(Self::expect_ok)
1061    }
1062
1063    // === Database operations (DatabaseOp via AuthenticatedDb envelope) ===
1064
1065    /// Acquire a [`TransactionContext`] for the given stores and scope.
1066    ///
1067    /// The returned context includes main-tree parents with heights,
1068    /// per-store subtree parents, `_settings` tips, and the merged
1069    /// `_settings` value — everything needed to build and sign an entry
1070    /// locally without further round-trips.
1071    pub async fn begin_transaction(
1072        &self,
1073        root_id: ID,
1074        identity: SigKey,
1075        stores: Vec<String>,
1076        scope: ReadScope,
1077    ) -> crate::Result<TransactionContext> {
1078        let resp = self
1079            .db_request(
1080                root_id,
1081                identity,
1082                DatabaseOp::BeginTransaction { stores, scope },
1083            )
1084            .await?;
1085        match resp {
1086            ServiceResponse::TransactionContext(ctx) => Ok(ctx),
1087            other => Err(unexpected_response("TransactionContext", &other)),
1088        }
1089    }
1090
1091    /// Fetch the server-materialized merged state of an unencrypted store.
1092    pub async fn get_store_state(
1093        &self,
1094        root_id: ID,
1095        identity: SigKey,
1096        store: String,
1097    ) -> crate::Result<WireCrdtValue> {
1098        let resp = self
1099            .db_request(root_id, identity, DatabaseOp::GetStoreState { store })
1100            .await?;
1101        match resp {
1102            ServiceResponse::CrdtValue(v) => Ok(v),
1103            other => Err(unexpected_response("CrdtValue", &other)),
1104        }
1105    }
1106
1107    /// Fetch ordered, verified, opaque store entries reachable from `tips`.
1108    ///
1109    /// Universal primitive — works for encrypted stores (client decrypts
1110    /// locally) as well as unencrypted ones.
1111    pub async fn get_store_entries(
1112        &self,
1113        root_id: ID,
1114        identity: SigKey,
1115        store: String,
1116        tips: Vec<ID>,
1117        scope: ReadScope,
1118    ) -> crate::Result<Vec<Entry>> {
1119        let resp = self
1120            .db_request(
1121                root_id,
1122                identity,
1123                DatabaseOp::GetStoreEntries { store, tips, scope },
1124            )
1125            .await?;
1126        match resp {
1127            ServiceResponse::Entries(entries) => Ok(entries),
1128            other => Err(unexpected_response("Entries", &other)),
1129        }
1130    }
1131
1132    /// Fetch the database's Verified-frontier tips.
1133    pub async fn get_verified_tips(
1134        &self,
1135        root_id: ID,
1136        identity: SigKey,
1137    ) -> crate::Result<Snapshot> {
1138        let resp = self
1139            .db_request(root_id, identity, DatabaseOp::GetVerifiedTips)
1140            .await?;
1141        match resp {
1142            ServiceResponse::Ids(ids) => Ok(ids),
1143            other => Err(unexpected_response("Ids", &other)),
1144        }
1145    }
1146
1147    /// Submit a client-signed entry to the server.
1148    ///
1149    /// The server stores the entry as `Unverified` and runs its own
1150    /// verification pass — it never trusts a submitted entry's claimed
1151    /// validity.
1152    pub async fn submit_signed_entry(
1153        &self,
1154        root_id: ID,
1155        identity: SigKey,
1156        entry: Entry,
1157    ) -> crate::Result<()> {
1158        let resp = self
1159            .db_request(
1160                root_id,
1161                identity,
1162                DatabaseOp::SubmitSignedEntry {
1163                    entry: Box::new(entry),
1164                },
1165            )
1166            .await?;
1167        match resp {
1168            ServiceResponse::Ok => Ok(()),
1169            other => Err(unexpected_response("Ok", &other)),
1170        }
1171    }
1172
1173    /// Fetch a single database entry by id.
1174    ///
1175    /// Gated post-fetch by the entry's owning tree, so the caller must hold
1176    /// at least `Read` on the database the entry belongs to.
1177    pub async fn db_get_entry(
1178        &self,
1179        root_id: ID,
1180        identity: SigKey,
1181        id: ID,
1182    ) -> crate::Result<Entry> {
1183        let resp = self
1184            .db_request(root_id, identity, DatabaseOp::GetEntry { id })
1185            .await?;
1186        match resp {
1187            ServiceResponse::Entry(entry) => Ok(entry),
1188            other => Err(unexpected_response("Entry", &other)),
1189        }
1190    }
1191
1192    /// Subtree tips reachable from given main-tree entries.
1193    pub async fn store_snapshot_at(
1194        &self,
1195        root_id: ID,
1196        identity: SigKey,
1197        store: String,
1198        up_to: Vec<ID>,
1199    ) -> crate::Result<Snapshot> {
1200        let resp = self
1201            .db_request(
1202                root_id,
1203                identity,
1204                DatabaseOp::GetStoreTipsUpToEntries { store, up_to },
1205            )
1206            .await?;
1207        match resp {
1208            ServiceResponse::Ids(ids) => Ok(ids),
1209            other => Err(unexpected_response("Ids", &other)),
1210        }
1211    }
1212
1213    /// Compute merge state: lowest common ancestor + path to tip entries.
1214    pub async fn compute_merge_state(
1215        &self,
1216        root_id: ID,
1217        identity: SigKey,
1218        store: String,
1219        entry_ids: Vec<ID>,
1220    ) -> crate::Result<MergeState> {
1221        let resp = self
1222            .db_request(
1223                root_id,
1224                identity,
1225                DatabaseOp::ComputeMergeState { store, entry_ids },
1226            )
1227            .await?;
1228        match resp {
1229            ServiceResponse::MergeState(state) => Ok(state),
1230            other => Err(unexpected_response("MergeState", &other)),
1231        }
1232    }
1233
1234    /// Resolve cached state, returning its opaque server-issued view onto the published record set.
1235    pub async fn resolve_store_state(
1236        &self,
1237        identity: SigKey,
1238        request: StoreStateRequest,
1239    ) -> crate::Result<Option<String>> {
1240        match self
1241            .db_request(
1242                request.database.clone(),
1243                identity,
1244                DatabaseOp::ResolveStoreState { request },
1245            )
1246            .await?
1247        {
1248            ServiceResponse::RecordView(view) => Ok(view),
1249            other => Err(unexpected_response("RecordView", &other)),
1250        }
1251    }
1252
1253    /// Begin a private build, returning its opaque capability.
1254    pub async fn begin_store_state_staging(
1255        &self,
1256        identity: SigKey,
1257        request: StoreStateRequest,
1258    ) -> crate::Result<String> {
1259        match self
1260            .db_request(
1261                request.database.clone(),
1262                identity,
1263                DatabaseOp::BeginStoreStateStaging { request },
1264            )
1265            .await?
1266        {
1267            ServiceResponse::Token(token) => Ok(token),
1268            other => Err(unexpected_response("Token", &other)),
1269        }
1270    }
1271
1272    /// Upload one idempotent chunk of records into the private build.
1273    pub async fn stage_store_state_records(
1274        &self,
1275        root: ID,
1276        identity: SigKey,
1277        token: String,
1278        chunk_id: u64,
1279        records: RecordMutations,
1280    ) -> crate::Result<()> {
1281        self.db_request(
1282            root,
1283            identity,
1284            DatabaseOp::StageStoreStateRecords {
1285                token,
1286                chunk_id,
1287                records: records.into_iter().collect(),
1288            },
1289        )
1290        .await
1291        .and_then(Self::expect_ok)
1292    }
1293
1294    /// Publish the private build and return a view onto the published record set.
1295    pub async fn publish_store_state(
1296        &self,
1297        root: ID,
1298        identity: SigKey,
1299        token: String,
1300    ) -> crate::Result<String> {
1301        match self
1302            .db_request(root, identity, DatabaseOp::PublishStoreState { token })
1303            .await?
1304        {
1305            ServiceResponse::RecordView(Some(view)) => Ok(view),
1306            other => Err(unexpected_response("RecordView", &other)),
1307        }
1308    }
1309
1310    /// Discard an unfinished private build and its bytes.
1311    pub async fn abort_store_state(
1312        &self,
1313        root: ID,
1314        identity: SigKey,
1315        token: String,
1316    ) -> crate::Result<()> {
1317        self.db_request(root, identity, DatabaseOp::AbortStoreState { token })
1318            .await
1319            .and_then(Self::expect_ok)
1320    }
1321
1322    /// Read one record through a resolved view.
1323    pub async fn store_state_record_get(
1324        &self,
1325        root: ID,
1326        identity: SigKey,
1327        view: String,
1328        key: Vec<u8>,
1329    ) -> crate::Result<Option<Vec<u8>>> {
1330        match self
1331            .db_request(
1332                root,
1333                identity,
1334                DatabaseOp::StoreStateRecordGet { view, key },
1335            )
1336            .await?
1337        {
1338            ServiceResponse::Record(record) => Ok(record),
1339            other => Err(unexpected_response("Record", &other)),
1340        }
1341    }
1342
1343    /// Read one bounded ordered page through a resolved view.
1344    #[allow(clippy::too_many_arguments)]
1345    pub async fn store_state_record_scan(
1346        &self,
1347        root: ID,
1348        identity: SigKey,
1349        view: String,
1350        range: RecordRange,
1351        after: Option<Vec<u8>>,
1352        max_records: u32,
1353        max_encoded_bytes: u32,
1354    ) -> crate::Result<RecordPage> {
1355        match self
1356            .db_request(
1357                root,
1358                identity,
1359                DatabaseOp::StoreStateRecordScan {
1360                    view,
1361                    range,
1362                    after,
1363                    max_records,
1364                    max_encoded_bytes,
1365                },
1366            )
1367            .await?
1368        {
1369            ServiceResponse::RecordPage(page) => Ok(page),
1370            other => Err(unexpected_response("RecordPage", &other)),
1371        }
1372    }
1373}
1374
1375fn unexpected_response(expected: &str, actual: &ServiceResponse) -> crate::Error {
1376    crate::Error::Io(std::io::Error::new(
1377        std::io::ErrorKind::InvalidData,
1378        format!("Expected {expected} response, got {actual:?}"),
1379    ))
1380}
1381
1382/// Canonical "connection torn down" error returned to any caller whose
1383/// request couldn't reach (or be answered by) the daemon.
1384fn connection_aborted() -> crate::Error {
1385    crate::Error::Io(std::io::Error::new(
1386        std::io::ErrorKind::ConnectionAborted,
1387        "Server closed connection unexpectedly",
1388    ))
1389}
1390
1391/// Background task driving the read half of the socket.
1392///
1393/// Loops on [`read_frame`] and demuxes by [`ServerFrame`] variant:
1394///
1395/// - `Response(r)`: pop the front of the pending FIFO and resolve its
1396///   oneshot. If the queue is empty something has gone badly wrong
1397///   (server sent more responses than the client issued requests) — log
1398///   and continue.
1399/// - `Notification(Notification::DatabaseWrite { … })`: upgrade the
1400///   attached `WeakInstance` and route the event into its callback
1401///   registry via [`crate::Instance::fire_write_callbacks`]. If no
1402///   instance is attached yet or the instance has been dropped, the
1403///   notification is silently dropped — both are expected end-states,
1404///   not errors.
1405///
1406/// Exit conditions: clean EOF (server closed), any read error, or any
1407/// deserialisation error. On exit the task drops its `Arc<inner>`, which
1408/// in turn drops every remaining oneshot sender in `pending`, surfacing
1409/// as a `RecvError` on each awaiting `request()` (translated to a
1410/// connection-closed `io::Error` there).
1411async fn run_reader_task(mut reader: ReadHalf<UnixStream>, inner: Arc<RemoteConnectionInner>) {
1412    loop {
1413        let frame_result: crate::Result<Option<ServerFrame>> = read_frame(&mut reader).await;
1414        let frame = match frame_result {
1415            Ok(Some(f)) => f,
1416            Ok(None) => break, // Clean EOF
1417            Err(e) => {
1418                tracing::debug!("RemoteConnection reader error: {e}");
1419                break;
1420            }
1421        };
1422
1423        match frame {
1424            ServerFrame::Response(resp) => {
1425                let next = inner.pending_lock().pop_front();
1426                match next {
1427                    Some(tx) => {
1428                        // Receiver dropped → the caller has already given
1429                        // up. Not an error worth logging.
1430                        let _ = tx.send(*resp);
1431                    }
1432                    None => {
1433                        tracing::warn!(
1434                            "RemoteConnection reader: response with no pending request; dropping"
1435                        );
1436                    }
1437                }
1438            }
1439            ServerFrame::Notification(notif) => {
1440                route_notification(&inner, notif);
1441            }
1442        }
1443    }
1444
1445    // Mark the connection dead with `Release` ordering paired against
1446    // `request()`'s `Acquire` load: any post-exit caller that observes
1447    // `true` is guaranteed to also see the drained `pending` queue and
1448    // bail with `ConnectionAborted` instead of pushing a sender no one
1449    // will ever pop. The helper also drops every per-tree worker
1450    // sender so workers exit cleanly without prolonging `inner`'s
1451    // lifetime.
1452    inner.mark_dead();
1453}
1454
1455/// Route a notification to its per-tree worker, spawning one if this is
1456/// the first notification for the tree on this connection.
1457///
1458/// Worker spawn is lazy: we don't create a worker for a tree until the
1459/// daemon actually pushes a notification for it. The map of per-tree
1460/// senders lives in `inner.tree_workers` (std mutex; the map is touched
1461/// for at most a single insert + clone per notification).
1462///
1463/// Sends are best-effort: if the worker has already exited (e.g. the
1464/// connection is winding down and `inner` is mid-drop), the send fails
1465/// and we silently drop. Same posture as the previous single-dispatch
1466/// shape.
1467///
1468/// TODO(dispatch-bound): per-tree channels are `unbounded`. Under
1469/// sustained write load on one tree, a slow user callback lets that
1470/// worker's queue grow without limit, holding all queued notifications
1471/// in client memory.
1472///
1473/// Under the cursor-only `Notification::DatabaseWrite` shape, drops are
1474/// *recoverable*: a worker that drops event N still receives event N+1
1475/// whose `post_tips` reflects the daemon's latest frontier, and the user
1476/// callback's next fire's `previous_tips = post_tips_of_N+1` lets
1477/// `ids_added` pick up any skipped IDs. So drop-oldest via
1478/// `Mutex<VecDeque<Notification>> + Notify` is the right v2 shape —
1479/// roughly a 30-line primitive isolated to this file. Bound size is the
1480/// tunable; `~256` is a reasonable starting point.
1481///
1482/// Deferred for its own PR (alongside server-side
1483/// [`TODO(backpressure)`]) so the drop-semantics tests get focused
1484/// review.
1485fn route_notification(inner: &Arc<RemoteConnectionInner>, notif: Notification) {
1486    let tree_id = match &notif {
1487        Notification::DatabaseWrite { root_id, .. } => root_id.clone(),
1488    };
1489    let tx = {
1490        let mut workers = inner.tree_workers.lock().unwrap_or_else(|p| p.into_inner());
1491        workers
1492            .entry(tree_id)
1493            .or_insert_with(|| {
1494                let (tx, rx) = mpsc::unbounded_channel::<Notification>();
1495                let weak = Arc::downgrade(inner);
1496                tokio::spawn(run_tree_worker(rx, weak));
1497                tx
1498            })
1499            .clone()
1500    };
1501    let _ = tx.send(notif);
1502}
1503
1504/// Drain one tree's notification queue, dispatching to the attached
1505/// `Instance`'s callback registry in arrival order.
1506///
1507/// **Ordering guarantee within the tree.** Notifications are processed
1508/// strictly one at a time — the next `recv()` doesn't run until the
1509/// previous callback's `fire_write_callbacks().await` has returned.
1510/// The reader pushes in the order frames hit the socket, so user
1511/// callbacks for this tree observe events in the daemon's canonical
1512/// order.
1513///
1514/// **No ordering guarantee across trees.** Different trees have their
1515/// own worker tasks; a slow callback on tree A doesn't stall tree B's
1516/// dispatches on the same connection. This is the load-bearing
1517/// difference from the previous single-drain-task shape.
1518///
1519/// **Why this can't be inline in the reader.** User callbacks may
1520/// issue wire ops (e.g. `Database::open` over the connected instance)
1521/// whose responses land through the same reader. Awaiting a callback
1522/// inline would deadlock the reader against the response it is
1523/// supposed to deliver. Per-tree workers keep the reader free.
1524///
1525/// **Lifecycle.** Holds `Weak<RemoteConnectionInner>` so it does not
1526/// extend `inner`'s lifetime. When the reader exits it clears
1527/// `tree_workers`, dropping every sender; `recv()` returns `None`;
1528/// this worker exits.
1529async fn run_tree_worker(
1530    mut rx: mpsc::UnboundedReceiver<Notification>,
1531    weak_inner: Weak<RemoteConnectionInner>,
1532) {
1533    while let Some(notif) = rx.recv().await {
1534        // Snapshot the attached `WeakInstance` per-notification under
1535        // the std mutex — never held across an await. Worst case is
1536        // `None`, which we treat as "instance not attached yet" (the
1537        // attach-vs-first-notification race is impossible in practice
1538        // because attach happens before `Instance::connect` returns,
1539        // and subscriptions only start after that).
1540        let weak_instance = {
1541            let Some(inner) = weak_inner.upgrade() else {
1542                tracing::debug!("RemoteConnection tree worker: inner gone; exiting");
1543                return;
1544            };
1545            let guard = inner
1546                .weak_instance
1547                .lock()
1548                .unwrap_or_else(|poisoned| poisoned.into_inner());
1549            guard.clone()
1550        };
1551        let Some(weak) = weak_instance else {
1552            tracing::debug!(
1553                "RemoteConnection tree worker: notification before attach_instance; dropping"
1554            );
1555            continue;
1556        };
1557        let Some(instance) = weak.upgrade() else {
1558            tracing::debug!(
1559                "RemoteConnection tree worker: instance dropped; ignoring notification"
1560            );
1561            continue;
1562        };
1563
1564        match notif {
1565            Notification::DatabaseWrite {
1566                root_id,
1567                previous_tips,
1568                post_tips,
1569                source,
1570            } => {
1571                instance
1572                    .fire_write_callbacks(&root_id, &previous_tips, &post_tips, source)
1573                    .await;
1574            }
1575        }
1576    }
1577}
1578
1579/// Bound on how long the sweep waits for the daemon's
1580/// `UnsubscribeWrites` ack before declaring the connection broken.
1581///
1582/// The daemon's handler is a single hashmap removal — well under a
1583/// millisecond on local transport, single-digit ms over loopback TCP.
1584/// Five seconds is a generous "is the daemon alive at all" bound; on
1585/// expiry we mark the connection dead via
1586/// [`RemoteConnectionInner::mark_dead`] and let callers reconnect.
1587const UNSUBSCRIBE_RTT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
1588
1589/// Periodically tear down `Idle` per-tree subscriptions whose grace
1590/// window has elapsed.
1591///
1592/// Wakes every [`SWEEP_INTERVAL`]; for each entry in
1593/// `subscribed_trees` that *appears* to be `Idle` past
1594/// [`IDLE_GRACE_WINDOW`], the sweep processes that tree under its
1595/// per-tree subscription fence (see
1596/// `RemoteConnectionInner::subscription_locks`):
1597///
1598/// 1. Acquire `subscription_lock(tree_id)`. This serializes with the
1599///    leader path in [`RemoteConnection::subscribe_writes`], which
1600///    must take the same lock before any mutation of
1601///    `subscribed_trees`. A racing re-registration that arrived after
1602///    the candidate snapshot blocks here until our unsubscribe is
1603///    fully acked.
1604/// 2. Re-check state under the `subscribed_trees` std mutex. If still
1605///    `Idle` past grace, remove the entry and capture its identity
1606///    for the wire call. If the state changed (re-registered, swept
1607///    by another path) skip.
1608/// 3. Send `UnsubscribeWrites` via [`RemoteConnection::unsubscribe_writes`]
1609///    wrapped in [`tokio::time::timeout`] of [`UNSUBSCRIBE_RTT_TIMEOUT`].
1610///    On timeout the daemon is wedged on a trivial op — mark the
1611///    connection dead and exit.
1612/// 4. Release `subscription_lock`. A pending `subscribe_writes` on
1613///    this tree can now proceed: it observes `None` in
1614///    `subscribed_trees`, takes the leader path, and sends a fresh
1615///    `SubscribeWrites` — guaranteed to land on the daemon *after*
1616///    our `UnsubscribeWrites` because we held the fence across the
1617///    ack, independent of the daemon's request-dispatch shape.
1618///
1619/// Holds `Weak<RemoteConnectionInner>` so it doesn't extend `inner`'s
1620/// lifetime. Exits when the upgrade fails (last `Arc<inner>` dropped →
1621/// connection torn down) or when the timeout path marks the connection
1622/// dead.
1623async fn run_sweep_task(weak_inner: Weak<RemoteConnectionInner>) {
1624    let grace = idle_grace_window();
1625    let mut ticker = tokio::time::interval(sweep_interval());
1626    // Skip the initial immediate tick so we don't sweep before any
1627    // subscription has had a chance to register.
1628    ticker.tick().await;
1629    loop {
1630        ticker.tick().await;
1631        let Some(inner) = weak_inner.upgrade() else {
1632            tracing::debug!("RemoteConnection sweep: inner gone; exiting");
1633            return;
1634        };
1635
1636        // Cheap snapshot: collect tree IDs that *appear* to be Idle
1637        // past grace. The per-tree fence below re-checks under the
1638        // std mutex so a re-registration that lands between snapshot
1639        // and fence acquisition is honored.
1640        let candidates: Vec<ID> = {
1641            let subs = inner
1642                .subscribed_trees
1643                .lock()
1644                .unwrap_or_else(|p| p.into_inner());
1645            let now = std::time::Instant::now();
1646            subs.iter()
1647                .filter_map(|(id, state)| match state {
1648                    SubState::Idle { since, .. } if now.duration_since(*since) >= grace => {
1649                        Some(id.clone())
1650                    }
1651                    _ => None,
1652                })
1653                .collect()
1654        };
1655
1656        // Resolve the connection handle once for the batch. `live: None`
1657        // is load-bearing: this is an internal handle, and minting a
1658        // token here would `mark_dead` the connection every time the
1659        // batch handle drops at the end of a tick. The strong `inner`
1660        // moved in here is released when `conn` drops at the end of the
1661        // tick, before the next `upgrade()`.
1662        let conn = RemoteConnection { inner, _live: None };
1663
1664        for tree_id in candidates {
1665            // Per-tree fence: hold across the full Unsubscribe RTT so
1666            // a concurrent `subscribe_writes` on the same tree must
1667            // wait for the daemon's ack before sending its Subscribe.
1668            // See `RemoteConnectionInner::subscription_locks`.
1669            let sub_lock = conn.inner.subscription_lock(&tree_id);
1670            let _sub_guard = sub_lock.lock().await;
1671
1672            // Re-check under the std mutex now that we hold the
1673            // fence. If a re-registration won the race before we got
1674            // the fence, state will be `Subscribed` (or `InFlight` /
1675            // absent on an unrelated race) and we skip — the next
1676            // sweep tick will pick it up if it goes Idle again.
1677            let identity = {
1678                let mut subs = conn
1679                    .inner
1680                    .subscribed_trees
1681                    .lock()
1682                    .unwrap_or_else(|p| p.into_inner());
1683                let now = std::time::Instant::now();
1684                match subs.get(&tree_id) {
1685                    Some(SubState::Idle { since, .. }) if now.duration_since(*since) >= grace => {
1686                        // Still due — pull the entry out so a racing
1687                        // `subscribe_writes` (queued behind our
1688                        // subscription_lock) will observe `None` when
1689                        // it finally proceeds and take the leader path.
1690                        match subs.remove(&tree_id) {
1691                            Some(SubState::Idle { identity, .. }) => Some(identity),
1692                            _ => unreachable!("just observed Idle under the same lock"),
1693                        }
1694                    }
1695                    _ => None,
1696                }
1697            };
1698
1699            let Some(identity) = identity else {
1700                continue;
1701            };
1702
1703            // Wire round-trip under the fence. Timeout-then-teardown
1704            // on hang: a daemon that can't ack a hashmap removal in
1705            // five seconds is broken; mark dead and let the next
1706            // caller reconnect. Don't release the fence and continue
1707            // — that re-opens the race we're fencing against.
1708            match tokio::time::timeout(
1709                UNSUBSCRIBE_RTT_TIMEOUT,
1710                conn.unsubscribe_writes(tree_id.clone(), identity),
1711            )
1712            .await
1713            {
1714                Ok(Ok(())) => tracing::debug!(?tree_id, "lazy unsubscribe complete"),
1715                Ok(Err(e)) => tracing::debug!(
1716                    ?tree_id,
1717                    "lazy unsubscribe failed (connection likely closing): {e}"
1718                ),
1719                Err(_elapsed) => {
1720                    tracing::error!(
1721                        ?tree_id,
1722                        "lazy unsubscribe timed out after {:?}; tearing down connection",
1723                        UNSUBSCRIBE_RTT_TIMEOUT,
1724                    );
1725                    conn.inner.mark_dead();
1726                    return;
1727                }
1728            }
1729        }
1730    }
1731}
1732
1733#[cfg(test)]
1734mod tests {
1735    use super::*;
1736
1737    /// Build a `RemoteConnection` over a socketpair — enough to exercise the
1738    /// subscription state machine without a daemon. No reader task is
1739    /// spawned and no wire traffic is sent; the peer end is returned so the
1740    /// caller keeps the socket open for the duration of the test.
1741    fn test_conn() -> (RemoteConnection, tokio::net::UnixStream) {
1742        let (client_side, peer) = tokio::net::UnixStream::pair().unwrap();
1743        let (_reader, writer) = tokio::io::split(client_side);
1744        let inner = Arc::new(RemoteConnectionInner {
1745            writer: Mutex::new(writer),
1746            pending: std::sync::Mutex::new(VecDeque::new()),
1747            weak_instance: std::sync::Mutex::new(None),
1748            session: RwLock::new(None),
1749            registered_keys: Mutex::new(HashSet::new()),
1750            subscribed_trees: std::sync::Mutex::new(HashMap::new()),
1751            subscription_locks: std::sync::Mutex::new(HashMap::new()),
1752            closed: AtomicBool::new(false),
1753            tree_workers: std::sync::Mutex::new(HashMap::new()),
1754            reader_abort: std::sync::Mutex::new(None),
1755        });
1756        (RemoteConnection { inner, _live: None }, peer)
1757    }
1758
1759    /// Regression: `transition_to_idle` must not flip a tree to `Idle` while
1760    /// the attached instance still holds a per-database callback for it.
1761    ///
1762    /// `WriteCallback::drop` calls `remove_write_callback` and
1763    /// `transition_to_idle` as two separate synchronous steps, so a
1764    /// registration on another thread can land between them: it inserts into
1765    /// the registry, then `subscribe_writes` observes `Subscribed` and returns
1766    /// with no wire traffic. An unconditional flip strands that live callback
1767    /// on an `Idle` subscription, and the sweep sends `UnsubscribeWrites` a
1768    /// grace window later — after which the callback never fires again, with
1769    /// no error surfaced anywhere.
1770    #[tokio::test]
1771    async fn transition_to_idle_skips_while_a_callback_is_registered() {
1772        use crate::auth::crypto::generate_keypair;
1773        use crate::backend::database::InMemory;
1774        use crate::crdt::Doc;
1775
1776        let (conn, _peer) = test_conn();
1777        // A local instance suffices — the guard only reads the callback
1778        // registry, which is the same registry on connected instances.
1779        let (instance, _admin) = crate::Instance::create_backend(
1780            Box::new(InMemory::new()),
1781            crate::NewUser::passwordless("admin"),
1782        )
1783        .await
1784        .unwrap();
1785        conn.attach_instance(instance.downgrade());
1786
1787        let (signing_key, _) = generate_keypair();
1788        let db = crate::Database::create(&instance, signing_key, Doc::new())
1789            .await
1790            .unwrap();
1791        let tree_id = db.root_id().clone();
1792
1793        conn.subscribed_trees_lock().insert(
1794            tree_id.clone(),
1795            SubState::Subscribed {
1796                identity: SigKey::default(),
1797            },
1798        );
1799
1800        let cb = db.on_write(|_event, _db| async { Ok(()) }).await.unwrap();
1801
1802        conn.transition_to_idle(&tree_id);
1803        assert!(
1804            matches!(
1805                conn.subscribed_trees_lock().get(&tree_id),
1806                Some(SubState::Subscribed { .. })
1807            ),
1808            "a live callback must keep the wire subscription Subscribed"
1809        );
1810
1811        // Releasing the last callback makes the transition legitimate. The
1812        // handle's own Drop is a no-op for the wire state here (a local
1813        // instance has no `remote_connection`), so drive it explicitly.
1814        drop(cb);
1815        conn.transition_to_idle(&tree_id);
1816        assert!(
1817            matches!(
1818                conn.subscribed_trees_lock().get(&tree_id),
1819                Some(SubState::Idle { .. })
1820            ),
1821            "with no callbacks left the subscription must go Idle"
1822        );
1823    }
1824}