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

eidetica/service/
protocol.rs

1//! Wire protocol types for the Eidetica service.
2//!
3//! The protocol uses length-prefixed JSON frames over a Unix domain socket.
4//! Each frame is a 4-byte big-endian length followed by the JSON payload.
5//!
6//! ## Request shape
7//!
8//! `ServiceRequest` is a flat enum holding pre-authentication lifecycle messages
9//! (`TrustedLoginUser`, `TrustedLoginProve`), the pre-auth `GetInstanceMetadata`
10//! query, and an `AuthenticatedDb` wrapper that carries every storage operation
11//! (including any user-management writes against `_users`). The wrapper
12//! bundles the `(root_id, identity)` scope so the server can validate each
13//! op against the connection's session keyset and the target database's
14//! auth settings. Pre-auth verification of the session pubkey happens once at
15//! login via a challenge-response handshake.
16//!
17//! The login lifecycle is **trusted** in the sense that the daemon ships the
18//! user's encrypted credentials (salt + AES-GCM ciphertext) to anyone who can
19//! connect to the socket and asks for them. This is safe in the local-socket
20//! model — filesystem permissions on the socket already bound the caller set
21//! to processes that could read the underlying DB files directly. A network
22//! transport would need a different shape (PAKE: OPAQUE/SRP) so the server
23//! doesn't release the blob until the client proves password knowledge in a
24//! way that doesn't leak it. The `TrustedLogin*` naming is a load-bearing
25//! reminder of that assumption — see § Trusted login threat model in the
26//! Service Architecture doc.
27//!
28//! `AuthenticatedDb` requests carry the caller's `root_id`/`identity` and are
29//! gated per-tree by the daemon's permission check; clients populate these
30//! from the session established by the `TrustedLogin*` flow.
31
32use std::collections::BTreeMap;
33
34use serde::{Deserialize, Serialize};
35use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt};
36
37use crate::auth::crypto::PublicKey;
38use crate::auth::types::{Permission, SigKey};
39use crate::backend::{InstanceMetadata, RecordPage, RecordRange, StoreStateRequest};
40use crate::entry::{Entry, ID};
41use crate::instance::WriteSource;
42use crate::service::error::ServiceError;
43use crate::snapshot::Snapshot;
44use crate::user::UserInfo;
45
46/// Protocol version. Version 0 indicates an unstable protocol that may change
47/// without notice between releases.
48///
49/// This constant is the compatibility gate for serialized types in this
50/// protocol. `#[non_exhaustive]` does **not** protect wire compatibility: a
51/// peer on an older version fails to deserialize an unknown variant. Adding a
52/// variant to a serialized enum (e.g. [`WriteSource`](crate::instance::WriteSource)
53/// inside [`Notification::DatabaseWrite`]) is therefore a protocol version
54/// bump, not a backward-compatible addition.
55pub const PROTOCOL_VERSION: u32 = 0;
56
57/// Maximum frame size: 64 MiB.
58pub const MAX_FRAME_SIZE: u32 = 64 * 1024 * 1024;
59
60/// Record payload budget leaving room for the JSON envelope and length prefix.
61pub const MAX_RECORD_CHUNK_BYTES: u32 = MAX_FRAME_SIZE - 1024 * 1024;
62
63/// Default upper bound for an encoded record page.
64pub const MAX_RECORD_PAGE_BYTES: u32 = 4 * 1024 * 1024;
65
66/// Handshake message sent by the client on connection.
67#[derive(Debug, Clone, Serialize, Deserialize)]
68pub struct Handshake {
69    pub protocol_version: u32,
70}
71
72/// Handshake acknowledgment sent by the server.
73#[derive(Debug, Clone, Serialize, Deserialize)]
74pub struct HandshakeAck {
75    pub protocol_version: u32,
76}
77
78// ===========================================================================
79// Database-level wire API.
80//
81// Every storage operation rides this single op enum: the server runs the
82// `Database` layer on its local instance, so verify-on-read and the Verified
83// frontier are server-side by construction, and every op is intrinsically
84// (tree, store, identity)-scoped. Carried in `ServiceRequest::AuthenticatedDb`.
85// ===========================================================================
86
87/// Which snapshot of the DAG an op observes. Mirrors the `Database`
88/// read posture: a write's parent tips are the tips of the *same* snapshot
89/// the caller reads (see the Verification Model design doc).
90#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
91pub enum ReadScope {
92    /// Default-safe: only the maximal all-`Verified` ancestor-closed prefix.
93    #[default]
94    Verified,
95    /// Also include `Unverified` entries (`Failed` always dropped). The
96    /// caller explicitly opted in via `Database::allow_unverified()`.
97    AllowUnverified,
98}
99
100/// A CRDT store's materialized state on the wire. Concrete `Store<T>` typing
101/// stays client-side sugar over this; the cache path already ships
102/// `serde_json` bytes today, so this introduces no new representation.
103pub type WireCrdtValue = serde_json::Value;
104
105/// Everything a client needs to build **and sign** an entry locally without
106/// further round-trips. The client owns its keys, so signing stays
107/// client-side; only the inputs `Transaction::commit` reads from storage
108/// before signing travel here. Heights accompany each parent so the client
109/// computes entry height without a follow-up `GetEntry` per parent.
110#[derive(Debug, Clone, Serialize, Deserialize)]
111pub struct TransactionContext {
112    /// Main-tree parent tips with their heights, in the caller's `scope`.
113    pub main_parents: Vec<(ID, u64)>,
114    /// Per-store parent tips (with heights) reachable from `main_parents`.
115    pub subtree_parents: BTreeMap<String, Vec<(ID, u64)>>,
116    /// `_settings` tips this transaction pins in signed metadata.
117    pub settings_tips: Vec<ID>,
118    /// Merged `_settings` state the entry is authored against (used to build
119    /// the auth settings the signature is validated under).
120    pub settings_value: WireCrdtValue,
121}
122
123/// Response for ComputeMergeState: lowest common ancestor + path to tips.
124#[derive(Debug, Clone, Serialize, Deserialize)]
125pub struct MergeState {
126    /// `None` when the entries share no common ancestor: the state is folded
127    /// from a default over the tips' full ancestry, which the client fetches
128    /// as whole entries (`GetStoreEntries`) rather than as a path of IDs.
129    pub merge_base: Option<ID>,
130    /// Entries between the base and the tips; empty when `merge_base` is
131    /// `None`.
132    pub path: Vec<ID>,
133}
134
135/// JSON-safe wire representation of opaque byte-keyed record mutations.
136pub type WireRecordMutations = Vec<(Vec<u8>, Option<Vec<u8>>)>;
137
138/// Database-level operations the server runs on its local `Database`.
139///
140/// The target database (`root_id`) and identity claim travel in
141/// [`AuthenticatedDbRequest`]; the per-tree gate runs against `root_id`
142/// (Read for begin/get*, Write for submit, Admin-on-`_databases` for
143/// set-metadata) before dispatch.
144#[derive(Debug, Clone, Serialize, Deserialize)]
145pub enum DatabaseOp {
146    /// Resolve cached state through an opaque view onto one published record set.
147    ResolveStoreState { request: StoreStateRequest },
148    /// Begin a private build.
149    BeginStoreStateStaging { request: StoreStateRequest },
150    /// Upload one idempotent chunk into the private build.
151    StageStoreStateRecords {
152        token: String,
153        chunk_id: u64,
154        records: WireRecordMutations,
155    },
156    /// Publish the private build and return a view onto the published record set.
157    PublishStoreState { token: String },
158    /// Discard an unfinished private build.
159    AbortStoreState { token: String },
160    /// Fetch one record from a published record set.
161    StoreStateRecordGet { view: String, key: Vec<u8> },
162    /// Fetch one bounded page from a published record set.
163    StoreStateRecordScan {
164        view: String,
165        range: RecordRange,
166        after: Option<Vec<u8>>,
167        max_records: u32,
168        max_encoded_bytes: u32,
169    },
170    /// Acquire everything needed to build+sign a transaction locally for the
171    /// given stores, with parents drawn from `scope`'s snapshot. Gate Read.
172    BeginTransaction {
173        stores: Vec<String>,
174        scope: ReadScope,
175    },
176    /// Submit a finished, client-signed entry. The server stores it
177    /// `Unverified` and runs its **own** verification pass — it never trusts
178    /// a submitted entry's claimed validity. Submit is *verification-gated,
179    /// not session-gated*: it requires only an authenticated connection, and
180    /// the per-tree permission gate is **not** applied (the server's
181    /// verification pass against the tree's pinned auth is the boundary). The
182    /// `required_permission()` value below is advisory only for this variant.
183    SubmitSignedEntry { entry: Box<Entry> },
184    /// The database's Verified-frontier tips (server runs `Database::snapshot`
185    /// on its local instance). Gate Read.
186    GetVerifiedTips,
187    /// Server-materialized merged state of an **unencrypted** store, against
188    /// the server's own Verified frontier. Gate Read.
189    GetStoreState { store: String },
190    /// Ordered (by subtree height), verified, opaque store entries reachable
191    /// from `tips` in `scope` — the universal primitive, incl. encrypted
192    /// stores (client decrypts+merges locally). Gate Read.
193    GetStoreEntries {
194        store: String,
195        tips: Vec<ID>,
196        scope: ReadScope,
197    },
198    /// Subtree tips reachable from given main-tree entry IDs.
199    /// Used by Transaction internals to discover store entries.
200    GetStoreTipsUpToEntries { store: String, up_to: Vec<ID> },
201
202    /// Lowest common ancestor + path to tip entries in a store DAG.
203    /// Fused to one RPC so base and path resolve against a single server
204    /// view; answered from separate requests they can straddle a sync
205    /// ingest and disagree.
206    ComputeMergeState { store: String, entry_ids: Vec<ID> },
207
208    /// Fetch a single entry by id (gated post-fetch by its owning tree). Gate
209    /// Read.
210    GetEntry { id: ID },
211
212    /// Rewrite the daemon's instance metadata (system-DB pointers). Gated by
213    /// `Admin` on `_databases` (a daemon-global system tree, resolved
214    /// server-side — *not* the request's `root_id`), so the per-tree gate is
215    /// special-cased for this variant in the dispatcher. Boxed to keep the
216    /// enum's stack footprint small — `InstanceMetadata` dominates its size.
217    SetInstanceMetadata { metadata: Box<InstanceMetadata> },
218
219    /// Subscribe this connection to write notifications for the request's
220    /// `root_id`, with an explicit initial cursor (`tips`).
221    ///
222    /// After the server returns `Ok`, every write the daemon observes on
223    /// that tree (local commits via `SubmitSignedEntry`, sync ingest via
224    /// `put_remote_entries`, etc.) is pushed back to this connection as a
225    /// [`Notification::DatabaseWrite`] frame. The frame's `previous_tips`
226    /// is computed from the daemon-side subscription cursor — initially
227    /// the `tips` supplied here, and advanced to each event's `post_tips`
228    /// as the daemon fires.
229    ///
230    /// **Cursor semantics**: pass the tips you just read your initial
231    /// state at. The first notification's `previous_tips` will exactly
232    /// equal `tips`, so the client can diff `tips → notification.post_tips`
233    /// to discover anything that happened between the initial read and
234    /// the daemon recognising the subscription. An empty `tips` means
235    /// "I have no initial state; start from the daemon's current view"
236    /// (the first notification's `previous_tips` will be the daemon's
237    /// tips at subscribe-time, captured under the per-tree lock).
238    ///
239    /// Idempotent: re-subscribing a tree this connection already
240    /// subscribed to is a no-op (`tips` on the re-call is ignored;
241    /// the cursor stays at whatever it was). Gate Read on `root_id`.
242    /// Subscriptions are cleared automatically when the connection
243    /// drops.
244    SubscribeWrites { tips: Snapshot },
245
246    /// Stop pushing write notifications for the request's `root_id` to this
247    /// connection. Idempotent: unsubscribing a tree that wasn't subscribed
248    /// is a no-op. Gate Read on `root_id`.
249    UnsubscribeWrites,
250
251    /// Return a point-in-time locator for the request's `root_id`.
252    /// Gate Read.
253    CreateTicket,
254}
255
256impl DatabaseOp {
257    /// Minimum permission the caller needs against the target database.
258    ///
259    /// Only `SubmitSignedEntry` mutates; everything else is a read. Every
260    /// read variant is tree-scoped via the request's `root_id`, so the
261    /// per-tree gate always runs for reads — there is no tree-less
262    /// fall-through. `SubmitSignedEntry` is the exception: the server skips
263    /// the per-tree gate for submit and relies on its own verification pass,
264    /// so the `Write(0)` returned here is advisory only for that variant
265    /// (kept for completeness / non-submit callers that inspect it).
266    pub fn required_permission(&self) -> Permission {
267        match self {
268            DatabaseOp::SubmitSignedEntry { .. }
269            | DatabaseOp::BeginStoreStateStaging { .. }
270            | DatabaseOp::StageStoreStateRecords { .. }
271            | DatabaseOp::PublishStoreState { .. }
272            | DatabaseOp::AbortStoreState { .. } => Permission::Write(0),
273            // Gated against `_databases`, not the request's `root_id`; the
274            // dispatcher special-cases this so the value here is advisory.
275            DatabaseOp::SetInstanceMetadata { .. } => Permission::Admin(0),
276            DatabaseOp::BeginTransaction { .. }
277            | DatabaseOp::GetVerifiedTips
278            | DatabaseOp::GetStoreState { .. }
279            | DatabaseOp::GetStoreEntries { .. }
280            | DatabaseOp::GetStoreTipsUpToEntries { .. }
281            | DatabaseOp::ComputeMergeState { .. }
282            | DatabaseOp::GetEntry { .. }
283            | DatabaseOp::ResolveStoreState { .. }
284            | DatabaseOp::StoreStateRecordGet { .. }
285            | DatabaseOp::StoreStateRecordScan { .. }
286            | DatabaseOp::SubscribeWrites { .. }
287            | DatabaseOp::UnsubscribeWrites
288            | DatabaseOp::CreateTicket => Permission::Read,
289        }
290    }
291}
292
293/// Payload of an `AuthenticatedDb` service request.
294///
295/// Bundles the database scope (`root_id`) and identity claim (`identity`) with
296/// the [`DatabaseOp`] to run. Boxed inside `ServiceRequest::AuthenticatedDb` to
297/// keep the top-level enum's stack footprint flat — `SigKey` and
298/// `DatabaseOp::SubmitSignedEntry` are large.
299#[derive(Debug, Clone, Serialize, Deserialize)]
300pub struct AuthenticatedDbRequest {
301    /// Root entry of the database this op targets (auth-settings lookup +
302    /// the implicit tree scope every `DatabaseOp` carries by construction).
303    pub root_id: ID,
304    /// Identity claim; verified against the connection's session keyset
305    /// before dispatch.
306    pub identity: SigKey,
307    /// Database operation to execute.
308    pub op: DatabaseOp,
309}
310
311/// Top-level request from client to server.
312///
313/// The shape is intentionally flat: pre-auth lifecycle and queries sit beside
314/// the `AuthenticatedDb` wrapper rather than under a nested enum. This makes the
315/// pre-auth surface visible at a glance and keeps the server's dispatch
316/// branches symmetric.
317#[derive(Debug, Clone, Serialize, Deserialize)]
318pub enum ServiceRequest {
319    // === Pre-auth: trusted login handshake ===
320    /// Step 1 of the trusted login flow. Client names a user; server responds
321    /// with a `TrustedLoginChallenge` carrying random bytes the client must
322    /// sign. The "Trusted" qualifier is a load-bearing reminder that this flow
323    /// assumes the caller is already trusted by the socket's filesystem
324    /// permissions — over a network transport this would need PAKE instead.
325    TrustedLoginUser { username: String },
326    /// Step 2 of the trusted login flow. Client returns a signature over the
327    /// challenge from `TrustedLoginUser`, computed with the user's root key.
328    /// Server verifies against the stored pubkey and, on success, marks the
329    /// connection authenticated.
330    TrustedLoginProve { signature: Vec<u8> },
331
332    // === Pre-auth: queries safe before login ===
333    /// Fetch the server's instance metadata (including device id). Used by
334    /// `Instance::connect` during the handshake to establish server identity.
335    GetInstanceMetadata,
336
337    // === Post-auth: extend the connection's session keyset ===
338    /// Step 1 of registering an additional pubkey on an already-authenticated
339    /// connection. The client names a `pubkey`; the server issues a random
340    /// challenge bound to that pubkey. The pubkey is added to the keyset only
341    /// after the client returns a valid signature in `SessionKeyRegister`.
342    ///
343    /// Session-key registration extends the connection's identity from the
344    /// single `login_pubkey` (from `TrustedLogin*`) to a *set* of pubkeys the
345    /// client has proven possession of. Per-tree reads gate against this set,
346    /// so a user can drive operations on databases authored by any of their
347    /// per-DB keys without re-authenticating the whole connection.
348    SessionKeyChallenge { pubkey: PublicKey },
349    /// Step 2 of registering an additional pubkey. Carries a signature over
350    /// the challenge issued by the matching `SessionKeyChallenge`. Server
351    /// verifies the signature with the named `pubkey`; on success the pubkey
352    /// joins the connection's session keyset and the challenge is consumed.
353    SessionKeyRegister {
354        pubkey: PublicKey,
355        signature: Vec<u8>,
356    },
357
358    // === Authenticated wrapper for every storage operation ===
359    /// All storage ops travel inside this wrapper. The inner
360    /// `AuthenticatedDbRequest` carries `(root_id, identity, op)` and is boxed
361    /// to keep the enum's discriminated size compact.
362    AuthenticatedDb(Box<AuthenticatedDbRequest>),
363}
364
365/// Server-initiated push to the client, interleaved with normal responses
366/// at any point after a connection has authenticated.
367///
368/// Notifications are not solicited by a specific request; the client signs up
369/// for them with a [`DatabaseOp::SubscribeWrites`] and unsubscribes by
370/// dropping the connection or sending [`DatabaseOp::UnsubscribeWrites`].
371///
372/// Notifications are **triggered only by settled-state writes** — i.e.
373/// entries that have passed local verification on the daemon. An entry
374/// that arrives `Unverified` (via sync, or as a `SubmitSignedEntry` body)
375/// is ingested silently and only produces a notification once the daemon's
376/// verification pass promotes it to `Verified`.
377///
378/// **This does not mean the bracket contains only Verified entries.** The
379/// cursors are raw DAG frontiers (`Backend::snapshot`), not the Verified
380/// frontier, so an `Unverified` or `Failed` entry sitting as a raw tip
381/// falls inside every subsequent bracket and
382/// [`Database::ids_added`](crate::Database::ids_added) will enumerate it.
383/// A subscriber that expands a bracket and fetches the IDs can therefore
384/// observe entries that failed the daemon's auth settings. Consumers that
385/// care must filter on verification status themselves; locally that is
386/// `Backend::get_verification_status`, and over the wire there is
387/// currently no way to do it at all — treat bracket-derived IDs as
388/// untrusted until read back through a gated path.
389///
390/// A second consequence of the raw-frontier cursor: because it advances
391/// past a still-`Unverified` tip, a later `verify()` that promotes that
392/// entry fires with `previous_tips == post_tips` for the subscriber, so
393/// `ids_added` returns empty and the promotion is never signalled. An
394/// entry can be reported once while untrusted and never mentioned again.
395///
396/// The frame carries cursor brackets only — no entry payloads, no entry
397/// IDs. Subscribers that need to enumerate the new entries expand the
398/// brackets locally via
399/// [`Database::ids_added`](crate::Database::ids_added). This decision
400/// has two consequences:
401///
402/// 1. **Security**: per-tree Read is gated once at `SubscribeWrites`; the
403///    publisher fan-out does not re-check on each event. Shipping
404///    cursors only means a subscriber whose permission is revoked
405///    mid-session can at most learn that *some* write happened on the
406///    tree — never the contents of those writes, which would only reach
407///    the client through an explicit, currently-gated read.
408/// 2. **Efficiency**: sync ingest can batch many entries; a cursor pair
409///    is constant-size regardless of batch width.
410#[derive(Debug, Clone, Serialize, Deserialize)]
411pub enum Notification {
412    /// A settled-state write landed on the daemon for `root_id`.
413    ///
414    /// - `previous_tips` is the daemon-side subscription cursor at the
415    ///   moment of this fire — i.e. the `previous_tips` of the event
416    ///   the daemon dispatched to this subscription's callback. Useful
417    ///   for thin-forwarder topologies and trace/debug.
418    /// - `post_tips` is the daemon's tips *after* this write. The
419    ///   client uses it to advance every local per-callback cursor for
420    ///   this tree — each local callback's next event will have
421    ///   `previous_tips = post_tips` (the cursor moves forward by
422    ///   exactly one event).
423    /// - `source` distinguishes local-vs-sync for consumers that want
424    ///   to branch.
425    DatabaseWrite {
426        root_id: ID,
427        previous_tips: Snapshot,
428        post_tips: Snapshot,
429        source: WriteSource,
430    },
431}
432
433/// Envelope for every frame the server writes to a client.
434///
435/// Strict request/response responses ride `Response`; server-initiated
436/// pushes (subscribed write events) ride `Notification`. The reader task on
437/// the client demuxes by variant: `Response` frames go to the next pending
438/// oneshot in FIFO order, `Notification` frames go to the local callback
439/// dispatcher.
440///
441/// `Response` boxes its payload because `ServiceResponse` is large (the
442/// `TrustedLoginChallenge` variant carries a full `UserInfo`), and an
443/// unboxed enum would force every `Notification` frame to carry that
444/// stack footprint too.
445#[derive(Debug, Clone, Serialize, Deserialize)]
446pub enum ServerFrame {
447    /// A response to a previously-sent [`ServiceRequest`].
448    Response(Box<ServiceResponse>),
449    /// A server-initiated push (subscription-driven).
450    Notification(Notification),
451}
452
453/// Response from server to client.
454#[derive(Debug, Clone, Serialize, Deserialize)]
455pub enum ServiceResponse {
456    /// Single entry
457    Entry(Entry),
458    /// Multiple entries
459    Entries(Vec<Entry>),
460    /// Multiple IDs
461    Ids(Snapshot),
462    /// Success with no data
463    Ok,
464    /// One optional opaque record.
465    Record(Option<Vec<u8>>),
466    /// One bounded, ordered record page.
467    RecordPage(RecordPage),
468    /// View onto one published record set.
469    RecordView(Option<String>),
470    /// Capability for one private build.
471    Token(String),
472    /// Transaction-build context (response to `DatabaseOp::BeginTransaction`).
473    TransactionContext(TransactionContext),
474    /// Materialized CRDT store state (response to `DatabaseOp::GetStoreState`).
475    CrdtValue(WireCrdtValue),
476    /// Merge state: lowest common ancestor + path to tips (response to
477    /// `DatabaseOp::ComputeMergeState`).
478    MergeState(MergeState),
479    /// Optional instance metadata
480    InstanceMetadata(Option<InstanceMetadata>),
481    /// Error response
482    Error(ServiceError),
483    /// Point-in-time locator for one authorized database.
484    DatabaseTicket(crate::sync::DatabaseTicket),
485    /// Challenge bytes returned in response to `TrustedLoginUser`, plus the
486    /// user's full record so the client can derive the password→key, decrypt
487    /// the root signing key locally, sign the challenge in a single
488    /// round-trip, and then build the `User` session from data the daemon
489    /// already returned — no second wire read of `_users` is required.
490    ///
491    /// `user_info.credentials` carries the (encrypted) root private key, its
492    /// `KeyStorage` envelope (algorithm/ciphertext/nonce for password-protected
493    /// users, raw `PrivateKey` for passwordless users), and the Argon2id salt
494    /// when password-protected. The non-credential fields (user_database_id,
495    /// status, timestamps) are what `User::new` consumes after the proof
496    /// step succeeds. See § Trusted login threat model in the Service
497    /// Architecture doc for why this is safe to ship to anyone who can
498    /// reach the socket.
499    TrustedLoginChallenge {
500        challenge: Vec<u8>,
501        user_uuid: String,
502        user_info: UserInfo,
503    },
504    /// Trusted login succeeded; the connection is now authenticated.
505    TrustedLoginOk,
506    /// Challenge bytes returned in response to `SessionKeyChallenge`. The
507    /// client signs these with the named pubkey's private key and returns the
508    /// signature in `SessionKeyRegister`.
509    SessionKeyChallenge { challenge: Vec<u8> },
510}
511
512/// Write a length-prefixed JSON frame to an async writer.
513pub async fn write_frame<W: AsyncWrite + Unpin, T: Serialize>(
514    writer: &mut W,
515    value: &T,
516) -> crate::Result<()> {
517    let payload = serde_json::to_vec(value)?;
518    let len = payload.len() as u32;
519    if len > MAX_FRAME_SIZE {
520        return Err(crate::Error::Io(std::io::Error::new(
521            std::io::ErrorKind::InvalidData,
522            format!("frame too large: {len} bytes (max {MAX_FRAME_SIZE})"),
523        )));
524    }
525    writer.write_all(&len.to_be_bytes()).await?;
526    writer.write_all(&payload).await?;
527    writer.flush().await?;
528    Ok(())
529}
530
531/// Read a length-prefixed JSON frame from an async reader.
532///
533/// Returns `None` on clean EOF (connection closed).
534pub async fn read_frame<R: AsyncRead + Unpin, T: for<'de> Deserialize<'de>>(
535    reader: &mut R,
536) -> crate::Result<Option<T>> {
537    let mut len_buf = [0u8; 4];
538    match reader.read_exact(&mut len_buf).await {
539        Ok(_) => {}
540        Err(e) if e.kind() == std::io::ErrorKind::UnexpectedEof => return Ok(None),
541        Err(e) => return Err(e.into()),
542    }
543    let len = u32::from_be_bytes(len_buf);
544    if len > MAX_FRAME_SIZE {
545        return Err(crate::Error::Io(std::io::Error::new(
546            std::io::ErrorKind::InvalidData,
547            format!("frame too large: {len} bytes (max {MAX_FRAME_SIZE})"),
548        )));
549    }
550    let mut payload = vec![0u8; len as usize];
551    reader.read_exact(&mut payload).await?;
552    let value = serde_json::from_slice(&payload)?;
553    Ok(Some(value))
554}
555
556#[cfg(test)]
557mod tests {
558    use super::*;
559
560    // Helper to make a simple entry for testing
561    fn test_id() -> ID {
562        ID::from_bytes("test-entry-id")
563    }
564
565    fn wrap(op: DatabaseOp) -> ServiceRequest {
566        ServiceRequest::AuthenticatedDb(Box::new(AuthenticatedDbRequest {
567            root_id: ID::default(),
568            identity: SigKey::default(),
569            op,
570        }))
571    }
572
573    /// Extract the inner `DatabaseOp` from a deserialised request, panicking if
574    /// the variant isn't `AuthenticatedDb`.
575    fn unwrap_op(req: ServiceRequest) -> DatabaseOp {
576        match req {
577            ServiceRequest::AuthenticatedDb(inner) => inner.op,
578            other => panic!("expected AuthenticatedDb, got {other:?}"),
579        }
580    }
581
582    #[test]
583    fn test_handshake_serde() {
584        let h = Handshake {
585            protocol_version: PROTOCOL_VERSION,
586        };
587        let json = serde_json::to_string(&h).unwrap();
588        let h2: Handshake = serde_json::from_str(&json).unwrap();
589        assert_eq!(h2.protocol_version, PROTOCOL_VERSION);
590    }
591
592    #[test]
593    fn test_handshake_ack_serde() {
594        let h = HandshakeAck {
595            protocol_version: PROTOCOL_VERSION,
596        };
597        let json = serde_json::to_string(&h).unwrap();
598        let h2: HandshakeAck = serde_json::from_str(&json).unwrap();
599        assert_eq!(h2.protocol_version, PROTOCOL_VERSION);
600    }
601
602    #[test]
603    fn test_request_get_entry_serde() {
604        let req = wrap(DatabaseOp::GetEntry { id: test_id() });
605        let json = serde_json::to_string(&req).unwrap();
606        let req2: ServiceRequest = serde_json::from_str(&json).unwrap();
607        match unwrap_op(req2) {
608            DatabaseOp::GetEntry { id } => assert_eq!(id, test_id()),
609            _ => panic!("wrong variant"),
610        }
611    }
612
613    #[test]
614    fn test_request_get_instance_metadata_serde() {
615        let req = ServiceRequest::GetInstanceMetadata;
616        let json = serde_json::to_string(&req).unwrap();
617        let req2: ServiceRequest = serde_json::from_str(&json).unwrap();
618        assert!(matches!(req2, ServiceRequest::GetInstanceMetadata));
619    }
620
621    #[test]
622    fn test_request_trusted_login_user_serde() {
623        let req = ServiceRequest::TrustedLoginUser {
624            username: "alice".to_string(),
625        };
626        let json = serde_json::to_string(&req).unwrap();
627        let req2: ServiceRequest = serde_json::from_str(&json).unwrap();
628        match req2 {
629            ServiceRequest::TrustedLoginUser { username } => assert_eq!(username, "alice"),
630            _ => panic!("wrong variant"),
631        }
632    }
633
634    #[test]
635    fn test_request_trusted_login_prove_serde() {
636        let req = ServiceRequest::TrustedLoginProve {
637            signature: b"sig-bytes".to_vec(),
638        };
639        let json = serde_json::to_string(&req).unwrap();
640        let req2: ServiceRequest = serde_json::from_str(&json).unwrap();
641        match req2 {
642            ServiceRequest::TrustedLoginProve { signature } => assert_eq!(signature, b"sig-bytes"),
643            _ => panic!("wrong variant"),
644        }
645    }
646
647    #[test]
648    fn test_response_ok_serde() {
649        let resp = ServiceResponse::Ok;
650        let json = serde_json::to_string(&resp).unwrap();
651        let resp2: ServiceResponse = serde_json::from_str(&json).unwrap();
652        assert!(matches!(resp2, ServiceResponse::Ok));
653    }
654
655    #[test]
656    fn test_response_ids_serde() {
657        let resp = ServiceResponse::Ids(Snapshot::new(vec![test_id(), ID::from_bytes("other")]));
658        let json = serde_json::to_string(&resp).unwrap();
659        let resp2: ServiceResponse = serde_json::from_str(&json).unwrap();
660        match resp2 {
661            ServiceResponse::Ids(ids) => {
662                assert_eq!(ids.len(), 2);
663                assert!(ids.contains(&test_id()));
664                assert!(ids.contains(&ID::from_bytes("other")));
665            }
666            _ => panic!("wrong variant"),
667        }
668    }
669
670    #[test]
671    fn test_response_error_serde() {
672        let se = ServiceError {
673            module: "backend".to_string(),
674            kind: "EntryNotFound".to_string(),
675            message: "Entry not found: abc".to_string(),
676        };
677        let resp = ServiceResponse::Error(se);
678        let json = serde_json::to_string(&resp).unwrap();
679        let resp2: ServiceResponse = serde_json::from_str(&json).unwrap();
680        match resp2 {
681            ServiceResponse::Error(e) => {
682                assert_eq!(e.module, "backend");
683                assert_eq!(e.kind, "EntryNotFound");
684            }
685            _ => panic!("wrong variant"),
686        }
687    }
688
689    #[test]
690    fn test_response_instance_metadata_none_serde() {
691        let resp = ServiceResponse::InstanceMetadata(None);
692        let json = serde_json::to_string(&resp).unwrap();
693        let resp2: ServiceResponse = serde_json::from_str(&json).unwrap();
694        assert!(matches!(resp2, ServiceResponse::InstanceMetadata(None)));
695    }
696
697    #[test]
698    fn test_response_trusted_login_challenge_serde() {
699        use crate::auth::crypto::generate_keypair;
700        use crate::user::{KeyStorage, UserCredentials, UserInfo, UserStatus};
701
702        let (_signing, pubkey) = generate_keypair();
703        let user_info = UserInfo {
704            username: "alice".to_string(),
705            user_database_id: ID::from_bytes("alice-db"),
706            credentials: UserCredentials {
707                root_key_id: pubkey.clone(),
708                root_key: KeyStorage::Encrypted {
709                    algorithm: "aes-256-gcm".to_string(),
710                    ciphertext: b"ct".to_vec(),
711                    nonce: b"123456789012".to_vec(),
712                },
713                password_salt: Some("salt-string".to_string()),
714            },
715            created_at: 1_700_000_000,
716            status: UserStatus::Active,
717        };
718
719        let resp = ServiceResponse::TrustedLoginChallenge {
720            challenge: b"random-bytes".to_vec(),
721            user_uuid: "uuid-alice".to_string(),
722            user_info: user_info.clone(),
723        };
724        let json = serde_json::to_string(&resp).unwrap();
725        let resp2: ServiceResponse = serde_json::from_str(&json).unwrap();
726        match resp2 {
727            ServiceResponse::TrustedLoginChallenge {
728                challenge,
729                user_uuid,
730                user_info: ui2,
731            } => {
732                assert_eq!(challenge, b"random-bytes");
733                assert_eq!(user_uuid, "uuid-alice");
734                assert_eq!(ui2.username, user_info.username);
735                assert_eq!(ui2.user_database_id, user_info.user_database_id);
736                assert_eq!(ui2.credentials.root_key_id, pubkey);
737                assert_eq!(
738                    ui2.credentials.password_salt.as_deref(),
739                    Some("salt-string")
740                );
741            }
742            _ => panic!("wrong variant"),
743        }
744    }
745
746    #[test]
747    fn test_response_trusted_login_ok_serde() {
748        let resp = ServiceResponse::TrustedLoginOk;
749        let json = serde_json::to_string(&resp).unwrap();
750        let resp2: ServiceResponse = serde_json::from_str(&json).unwrap();
751        assert!(matches!(resp2, ServiceResponse::TrustedLoginOk));
752    }
753
754    #[test]
755    fn test_database_op_subscribe_writes_serde() {
756        let req = wrap(DatabaseOp::SubscribeWrites {
757            tips: Snapshot::new(vec![ID::from_bytes("t1"), ID::from_bytes("t2")]),
758        });
759        let json = serde_json::to_string(&req).unwrap();
760        let req2: ServiceRequest = serde_json::from_str(&json).unwrap();
761        match unwrap_op(req2) {
762            DatabaseOp::SubscribeWrites { tips } => assert_eq!(tips.len(), 2),
763            other => panic!("expected SubscribeWrites, got {other:?}"),
764        }
765    }
766
767    #[test]
768    fn test_database_op_subscribe_writes_empty_tips_serde() {
769        let req = wrap(DatabaseOp::SubscribeWrites {
770            tips: Snapshot::EMPTY,
771        });
772        let json = serde_json::to_string(&req).unwrap();
773        let req2: ServiceRequest = serde_json::from_str(&json).unwrap();
774        match unwrap_op(req2) {
775            DatabaseOp::SubscribeWrites { tips } => assert!(tips.is_empty()),
776            other => panic!("expected SubscribeWrites, got {other:?}"),
777        }
778    }
779
780    #[test]
781    fn test_database_op_unsubscribe_writes_serde() {
782        let req = wrap(DatabaseOp::UnsubscribeWrites);
783        let json = serde_json::to_string(&req).unwrap();
784        let req2: ServiceRequest = serde_json::from_str(&json).unwrap();
785        assert!(matches!(unwrap_op(req2), DatabaseOp::UnsubscribeWrites));
786    }
787
788    #[test]
789    fn test_subscribe_ops_gate_read() {
790        assert_eq!(
791            DatabaseOp::SubscribeWrites {
792                tips: Snapshot::EMPTY
793            }
794            .required_permission(),
795            Permission::Read
796        );
797        assert_eq!(
798            DatabaseOp::UnsubscribeWrites.required_permission(),
799            Permission::Read
800        );
801    }
802
803    #[test]
804    fn test_server_frame_response_serde() {
805        let frame = ServerFrame::Response(Box::new(ServiceResponse::Ok));
806        let json = serde_json::to_string(&frame).unwrap();
807        let frame2: ServerFrame = serde_json::from_str(&json).unwrap();
808        match frame2 {
809            ServerFrame::Response(resp) => match *resp {
810                ServiceResponse::Ok => {}
811                other => panic!("expected ServiceResponse::Ok, got {other:?}"),
812            },
813            other => panic!("expected ServerFrame::Response(Ok), got {other:?}"),
814        }
815    }
816
817    #[test]
818    fn test_server_frame_notification_serde() {
819        let notif = Notification::DatabaseWrite {
820            root_id: test_id(),
821            previous_tips: Snapshot::new(vec![ID::from_bytes("tip-1"), ID::from_bytes("tip-2")]),
822            post_tips: Snapshot::new(vec![ID::from_bytes("post-1")]),
823            source: WriteSource::Remote,
824        };
825        let frame = ServerFrame::Notification(notif);
826        let json = serde_json::to_string(&frame).unwrap();
827        let frame2: ServerFrame = serde_json::from_str(&json).unwrap();
828        match frame2 {
829            ServerFrame::Notification(Notification::DatabaseWrite {
830                root_id,
831                previous_tips,
832                post_tips,
833                source,
834            }) => {
835                assert_eq!(root_id, test_id());
836                assert_eq!(previous_tips.len(), 2);
837                assert_eq!(post_tips, Snapshot::new(vec![ID::from_bytes("post-1")]));
838                assert_eq!(source, WriteSource::Remote);
839            }
840            other => panic!("expected ServerFrame::Notification(DatabaseWrite), got {other:?}"),
841        }
842    }
843
844    #[test]
845    fn ticket_wire_remains_tree_scoped() {
846        let tree = test_id();
847        let request = ServiceRequest::AuthenticatedDb(Box::new(AuthenticatedDbRequest {
848            root_id: tree.clone(),
849            identity: SigKey::default(),
850            op: DatabaseOp::CreateTicket,
851        }));
852        let encoded = serde_json::to_string(&request).unwrap();
853        let decoded: ServiceRequest = serde_json::from_str(&encoded).unwrap();
854        match decoded {
855            ServiceRequest::AuthenticatedDb(request) => assert_eq!(request.root_id, tree),
856            other => panic!("expected database request, got {other:?}"),
857        }
858    }
859
860    #[tokio::test]
861    async fn test_frame_eof_returns_none() {
862        // Use a real Unix socket pair for proper EOF semantics
863        let dir = tempfile::tempdir().unwrap();
864        let sock_path = dir.path().join("eof-test.sock");
865        let listener = tokio::net::UnixListener::bind(&sock_path).unwrap();
866
867        let client = tokio::net::UnixStream::connect(&sock_path).await.unwrap();
868        let (server_stream, _) = listener.accept().await.unwrap();
869
870        // Drop the server stream to close the connection
871        drop(server_stream);
872
873        let (mut reader, _writer) = tokio::io::split(client);
874        let result: crate::Result<Option<ServiceRequest>> = read_frame(&mut reader).await;
875        assert!(result.unwrap().is_none());
876    }
877
878    #[tokio::test]
879    async fn test_frame_max_size_rejection_on_write() {
880        let (client, _server) = tokio::io::duplex(1024);
881        let (_read, mut write) = tokio::io::split(client);
882
883        // Create a payload that's too large
884        let huge_string = "x".repeat(MAX_FRAME_SIZE as usize + 1);
885        let result = write_frame(&mut write, &huge_string).await;
886        assert!(result.is_err());
887    }
888
889    #[tokio::test]
890    async fn test_frame_max_size_rejection_on_read() {
891        let (client, server) = tokio::io::duplex(1024 * 1024);
892        let (mut client_read, _client_write) = tokio::io::split(client);
893        let (_server_read, mut server_write) = tokio::io::split(server);
894
895        // Write a fake frame header with size > MAX_FRAME_SIZE from the server end
896        let fake_len = MAX_FRAME_SIZE + 1;
897        tokio::spawn(async move {
898            server_write
899                .write_all(&fake_len.to_be_bytes())
900                .await
901                .unwrap();
902        });
903
904        let result: crate::Result<Option<ServiceRequest>> = read_frame(&mut client_read).await;
905        assert!(result.is_err());
906    }
907}