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 ¬if {
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}