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}