feat: http server for account and keypackages storage

This commit is contained in:
kaichaosun 2026-06-11 15:43:10 +08:00
parent cf09036f42
commit b2fd4f6ac1
No known key found for this signature in database
GPG Key ID: 223E0F992F4F03BF
9 changed files with 2120 additions and 1 deletions

7
.dockerignore Normal file
View File

@ -0,0 +1,7 @@
target
*.db
*.db-shm
*.db-wal
.git
.gitignore
.DS_Store

10
.gitignore vendored Normal file
View File

@ -0,0 +1,10 @@
# Rust build artifacts
/target
# Local SQLite databases created at runtime
*.db
*.db-shm
*.db-wal
# Editor / OS noise
.DS_Store

1185
Cargo.lock generated Normal file

File diff suppressed because it is too large Load Diff

27
Cargo.toml Normal file
View File

@ -0,0 +1,27 @@
# Standalone single-crate workspace: keeps this repo independent of any parent
# Cargo workspace it might be checked out beneath.
[workspace]
[package]
name = "keypackage-registry"
version = "0.1.0"
edition = "2024"
[[bin]]
name = "keypackage-registry"
path = "src/main.rs"
[dependencies]
anyhow = "1.0"
axum = "0.7"
base64 = "0.22"
clap = { version = "4", features = ["derive"] }
ed25519-dalek = "2.2.0"
hex = "0.4"
rusqlite = { version = "0.35", features = ["bundled-sqlcipher-vendored-openssl"] }
serde = { version = "1.0", features = ["derive"] }
serde_json = "1.0"
thiserror = "2"
tokio = { version = "1", features = ["rt-multi-thread", "macros", "signal", "sync", "time"] }
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }

48
Dockerfile Normal file
View File

@ -0,0 +1,48 @@
# syntax=docker/dockerfile:1
########################################
# Build stage
########################################
FROM rust:1-bookworm AS builder
# rusqlite's `bundled-sqlcipher-vendored-openssl` feature compiles SQLCipher and
# a vendored OpenSSL from source: the C toolchain ships in the base image, but
# the OpenSSL build also needs perl + make.
RUN apt-get update \
&& apt-get install -y --no-install-recommends perl make \
&& rm -rf /var/lib/apt/lists/*
WORKDIR /app
# Build dependencies first against a stub binary so the (slow) SQLCipher/OpenSSL
# compilation is cached and only re-runs when Cargo.toml/Cargo.lock change.
COPY Cargo.toml Cargo.lock ./
RUN mkdir src \
&& echo "fn main() {}" > src/main.rs \
&& cargo build --release --locked \
&& rm -rf src target/release/deps/keypackage_registry* target/release/keypackage-registry
# Now build the real binary; dependency artifacts above are reused.
COPY src ./src
RUN cargo build --release --locked --bin keypackage-registry
########################################
# Runtime stage
########################################
FROM debian:bookworm-slim AS runtime
RUN apt-get update \
&& apt-get install -y --no-install-recommends ca-certificates \
&& rm -rf /var/lib/apt/lists/*
COPY --from=builder /app/target/release/keypackage-registry /usr/local/bin/keypackage-registry
# Matches the default --bind 0.0.0.0:8080.
EXPOSE 8080
# Persist the SQLite database on a volume rather than the container layer.
VOLUME ["/data"]
ENV RUST_LOG=info
ENTRYPOINT ["keypackage-registry"]
CMD ["--bind", "0.0.0.0:8080", "--db", "/data/keypackage-registry.db"]

196
README.md
View File

@ -1,3 +1,197 @@
# Chat Store
For persistence of group chat users' key package.
Persistence for group-chat users' key packages — the **keypackage-registry**
HTTP service, extracted from [libchat](https://github.com/logos-messaging/libchat)
so it can be deployed on its own.
Standalone HTTP service that caches MLS KeyPackages keyed by **`device_id`**, so a
client can fetch a contact's keypackage without an out-of-band exchange.
Throwaway by design: scheduled to be replaced by a λLEZ-based service in v0.3, so
it intentionally has no overlap with the rest of libchat (axum + rusqlite only).
`device_id` is the hex-encoded 32-byte Ed25519 verifying key of a device.
It also runs a minimal **account service**: one signed blob per **`account_id`**
mapping an Account to its set of device (LocalIdentity) public keys, so clients
can invite every LocalIdentity of an account. `account_id` is the hex-encoded
32-byte Ed25519 AccountAddress verifying key. See
[Account device-list endpoints](#account-device-list-endpoints).
## Trust model
A bundle is an opaque **payload** plus its **signature**, published under a
**`device_id`** (the hex of the device's 32-byte Ed25519 verifying key).
The signed bytes and the wire bytes are identical, so a verifier checks the
signature over exactly what it received, no reconstruction.
The **server treats `payload` as a black box**: it never decodes it. It only
verifies that `signature` over the payload bytes is valid under `device_id`'s
key, then stores it. A valid signature is proof-of-possession — only the holder
of `device_id`'s key can publish under it — so an adversary can't publish under
a `device_id` it doesn't control, and junk is dropped before storage. The server
is not a trusted authority, so **consumers MUST also verify on retrieve**, and a
valid signature does not prove the device is authorized for any account (that
binding arrives with λLEZ in v0.3).
Consumers define the payload layout. Today it is:
```text
payload = timestamp_ms_le[8] || key_package[..]
```
Fixed-width field first with the variable `key_package` last makes it parse
exactly one way — no delimiter, even though `key_package` is arbitrary bytes.
## Building & running
```bash
cargo build --release
./target/release/keypackage-registry # binds 0.0.0.0:8080, db ./keypackage-registry.db
```
| Flag | Default | Description |
|------|---------|-------------|
| `--bind <addr>` | `0.0.0.0:8080` | HTTP bind address |
| `--db <path>` | `keypackage-registry.db` | SQLite database path |
| `--max-per-identity <n>` | `100` | Bundles retained per `device_id` |
| `--retention-days <n>` | `30` | Drop bundles older than this |
| `--prune-interval-secs <n>` | `3600` | How often the prune task runs |
Logs via `RUST_LOG` (default `info`).
## Docker
```bash
# Build the image
docker build -t chat-store .
# Run it, persisting the SQLite db on a named volume and exposing port 8080
docker run --rm -p 8080:8080 -v chat-store-data:/data chat-store
```
The image runs the binary with `--bind 0.0.0.0:8080 --db /data/keypackage-registry.db`
by default; override the `CMD` to change flags, e.g.:
```bash
docker run --rm -p 9000:9000 -v chat-store-data:/data chat-store \
--bind 0.0.0.0:9000 --db /data/registry.db --retention-days 14
```
## API
### `POST /v0/keypackage`
```json
{
"device_id": "hex(32-byte ed25519 verifying key)",
"payload": "base64(opaque signed bytes)",
"signature": "base64(64-byte ed25519 signature over payload)"
}
```
The server verifies `signature` over the (opaque) `payload` bytes under
`device_id`'s key before storing, keyed by `device_id`. It does not decode
`payload`. Returns `204` on success, `400` on malformed input or a signature
that fails to verify.
### `GET /v0/keypackage/{device_id}`
Returns the most recently submitted bundle for that `device_id`, or `404`:
```json
{
"payload": "base64(...)",
"signature": "base64(64-byte ed25519 signature)"
}
```
Consumers verify `signature` over the `payload` bytes using the key recovered
from `device_id`, then read `key_package` out of the payload. A bundle that
fails verification must be treated as not found.
## Account device-list endpoints
The account service stores **exactly one blob per `account_id`** mapping an
Account to its LocalIdentity device keys. Same trust model as keypackages: the
server verifies `signature` over `payload` under `account_id`'s key
(proof-of-possession), and consumers MUST re-verify on retrieve. Clients encode
a lamport-timestamped list of device public keys in `payload`; the rest of the
payload stays opaque to the server.
> Anti-replay: the server reads the lamport from the (signature-verified)
> `payload` and replaces the stored bundle only when the incoming lamport is
> strictly higher, returning `409` otherwise. Because the lamport is covered by
> the account signature it cannot be forged, so a replayed older-but-still-valid
> bundle cannot downgrade the device list, nor refresh the retention clock.
> Consumers should still compare lamports themselves as defence in depth.
### `POST /v0/account`
Upsert the device-list bundle for an account; replaces any previous value.
```json
{
"account_id": "hex(32-byte ed25519 AccountAddress verifying key)",
"payload": "base64(opaque signed bytes: lamport-ts + device pubkeys)",
"signature": "base64(64-byte ed25519 signature over payload by the account key)"
}
```
Returns `204` on success, `400` on malformed input or a signature that fails to
verify, and `409` when the bundle's lamport is not newer than the stored one
(replay / stale publish).
### `GET /v0/account/{account_id}`
Returns the stored bundle for that account, or `404`:
```json
{
"payload": "base64(...)",
"signature": "base64(64-byte ed25519 signature)",
"updated_at": 1700000000000
}
```
`updated_at` is the server's last-upsert time in Unix ms. Consumers verify
`signature` over `payload` under `account_id`'s key, then decode the device list.
## Storage & retention
Two SQLite tables: `keypackages` keyed by `device_id`, and `account_bundles`
(one row per `account_id`). A background task runs every `--prune-interval-secs`,
dropping keypackage bundles older than `--retention-days` (keeping at most
`--max-per-identity` per `device_id`) and dropping account bundles not refreshed
within `--retention-days`. The schema is an internal detail and may change.
## Smoke test
End-to-end check with the real `chat-cli` (which lives in the
[libchat](https://github.com/logos-messaging/libchat) repo) against a running
server:
```bash
# In this repo: start the server on a test port with a fresh db
cargo run -- --bind 127.0.0.1:18080 --db tmp/registry.db
# In a libchat checkout: register two identities (--smoketest exits after registering)
cargo build -p chat-cli
./target/debug/chat-cli --name alice --transport file --data tmp/alice \
--registry-url http://127.0.0.1:18080 --smoketest # exits 0 on success
./target/debug/chat-cli --name bob --transport file --data tmp/bob \
--registry-url http://127.0.0.1:18080 --smoketest
# Confirm both bundles landed
sqlite3 tmp/registry.db "SELECT substr(device_id,1,12), length(payload) FROM keypackages;"
```
A non-zero exit from `chat-cli` means the server rejected the submission — e.g.
the signature failed verification. `GET /v0/keypackage/{device_id}` returns `200`
for a registered device and `404` otherwise.
## Lifecycle
Exists to unblock contact-by-id flows on testnet; removed once λLEZ-based
discovery lands in v0.3. The seam is the `RegistrationService` trait in libchat
(`core/conversations/src/service_traits.rs`) — swapping implementations does not
touch the chat protocol.

245
src/handlers.rs Normal file
View File

@ -0,0 +1,245 @@
use std::sync::Arc;
use axum::extract::{Path, State};
use axum::http::StatusCode;
use axum::response::{IntoResponse, Response};
use axum::routing::{get, post};
use axum::{Json, Router};
use base64::Engine;
use base64::engine::general_purpose::STANDARD as BASE64;
use ed25519_dalek::{Signature, VerifyingKey};
use serde::{Deserialize, Serialize};
use crate::store::{Store, StoredAccountBundle, StoredKeyPackageBundle};
#[derive(Debug, Deserialize)]
pub struct SubmitRequest {
/// Hex of the 32-byte Ed25519 device verifying key. Used to verify the
/// signature and as the storage/lookup key. `payload` stays opaque.
pub device_id: String,
/// base64 of the signed payload. Opaque to the server — it never decodes it.
pub payload: String,
/// base64 of the 64-byte Ed25519 signature over `payload`. Verifying it
/// under `device_id`'s key is proof-of-possession: only the holder of that
/// key can publish under this `device_id`.
pub signature: String,
}
#[derive(Debug, Serialize)]
pub struct FetchResponse {
/// base64 of the stored payload; consumers verify `signature` over it.
pub payload: String,
pub signature: String,
}
#[derive(Debug, Serialize)]
struct ErrorBody {
error: String,
}
pub fn router(store: Arc<Store>) -> Router {
Router::new()
.route("/v0/keypackage", post(submit))
.route("/v0/keypackage/:device_id", get(fetch))
.route("/v0/account", post(submit_account))
.route("/v0/account/:account_id", get(fetch_account))
.with_state(store)
}
async fn submit(
State(store): State<Arc<Store>>,
Json(req): Json<SubmitRequest>,
) -> Result<StatusCode, ApiError> {
// Verify proof-of-possession before persisting. `payload` is opaque — the
// server only checks that `signature` over the received payload bytes is
// valid under `device_id`'s key. A valid signature means the submitter holds
// that key. This rejects junk early (DoS mitigation); consumers still verify
// on retrieve, the server is not a trusted authority.
let device_pubkey: [u8; 32] = hex::decode(&req.device_id)
.ok()
.and_then(|b| b.try_into().ok())
.ok_or_else(|| ApiError::bad("device_id: must be hex of a 32-byte key"))?;
let payload = BASE64
.decode(&req.payload)
.map_err(|_| ApiError::bad("payload: not valid base64"))?;
let signature: [u8; 64] = BASE64
.decode(&req.signature)
.ok()
.and_then(|b| b.try_into().ok())
.ok_or_else(|| ApiError::bad("signature: must be base64 of 64 bytes"))?;
let verifying_key = VerifyingKey::from_bytes(&device_pubkey)
.map_err(|_| ApiError::bad("device_id: not a valid ed25519 key"))?;
verifying_key
.verify_strict(&payload, &Signature::from_bytes(&signature))
.map_err(|_| ApiError::bad("signature: verification failed"))?;
store
.insert(
&req.device_id,
&StoredKeyPackageBundle {
payload,
signature: signature.to_vec(),
},
)
.map_err(ApiError::internal)?;
Ok(StatusCode::NO_CONTENT)
}
async fn fetch(
State(store): State<Arc<Store>>,
Path(device_id): Path<String>,
) -> Result<Json<FetchResponse>, ApiError> {
let Some(bundle) = store.latest(&device_id).map_err(ApiError::internal)? else {
return Err(ApiError::not_found("no keypackage for device"));
};
Ok(Json(FetchResponse {
payload: BASE64.encode(&bundle.payload),
signature: BASE64.encode(&bundle.signature),
}))
}
/// Request body for publishing a signed device-list bundle under an account.
///
/// The `payload` is intentionally opaque to the server. Clients are expected
/// to encode a lamport-timestamped list of device (LocalIdentity) Ed25519
/// public keys inside it so that consumers can detect stale bundles. The server
/// only verifies that `signature` is a valid Ed25519 signature over `payload`
/// made by the key identified by `account_id`.
#[derive(Debug, Deserialize)]
pub struct SubmitAccountRequest {
/// Hex of the 32-byte Ed25519 account (AccountAddress) verifying key.
/// Acts as both the storage key and the verification key.
pub account_id: String,
/// base64 of the opaque signed payload (lamport-ts + device pubkeys, etc.).
pub payload: String,
/// base64 of the 64-byte Ed25519 signature over `payload` made by the
/// account key. Proof-of-possession: only the account holder can publish.
pub signature: String,
}
#[derive(Debug, Serialize)]
pub struct FetchAccountResponse {
/// base64 of the stored payload.
pub payload: String,
/// base64 of the 64-byte Ed25519 signature.
pub signature: String,
/// Unix timestamp (ms) of the last successful upsert.
pub updated_at: i64,
}
/// `POST /v0/account` — upsert a signed device-list bundle for an account.
///
/// The server verifies the Ed25519 signature and then stores exactly one blob
/// per `account_id`, replacing any previous value. Clients should re-publish
/// whenever they add or rotate LocalIdentities.
async fn submit_account(
State(store): State<Arc<Store>>,
Json(req): Json<SubmitAccountRequest>,
) -> Result<StatusCode, ApiError> {
let account_pubkey: [u8; 32] = hex::decode(&req.account_id)
.ok()
.and_then(|b| b.try_into().ok())
.ok_or_else(|| ApiError::bad("account_id: must be hex of a 32-byte key"))?;
let payload = BASE64
.decode(&req.payload)
.map_err(|_| ApiError::bad("payload: not valid base64"))?;
let signature: [u8; 64] = BASE64
.decode(&req.signature)
.ok()
.and_then(|b| b.try_into().ok())
.ok_or_else(|| ApiError::bad("signature: must be base64 of 64 bytes"))?;
let verifying_key = VerifyingKey::from_bytes(&account_pubkey)
.map_err(|_| ApiError::bad("account_id: not a valid ed25519 key"))?;
verifying_key
.verify_strict(&payload, &Signature::from_bytes(&signature))
.map_err(|_| ApiError::bad("signature: verification failed"))?;
// Read the bundle's lamport so the store can reject replays. Safe to trust:
// the signature over `payload` was just verified, so the lamport can't be
// forged without the account key.
let lamport = crate::store::payload_lamport(&payload)
.ok_or_else(|| ApiError::bad("payload: too short to contain a lamport header"))?;
let applied = store
.upsert_account(
&req.account_id,
lamport,
&StoredAccountBundle {
payload,
signature: signature.to_vec(),
updated_at: 0, // filled in by store
},
)
.map_err(ApiError::internal)?;
if !applied {
return Err(ApiError::conflict(
"stale bundle: lamport is not newer than the stored one",
));
}
Ok(StatusCode::NO_CONTENT)
}
/// `GET /v0/account/:account_id` — fetch the device-list bundle for an account.
///
/// Returns the latest published bundle so consumers can verify the
/// account signature and decode the list of LocalIdentity keys themselves.
async fn fetch_account(
State(store): State<Arc<Store>>,
Path(account_id): Path<String>,
) -> Result<Json<FetchAccountResponse>, ApiError> {
let Some(bundle) = store.get_account(&account_id).map_err(ApiError::internal)? else {
return Err(ApiError::not_found("no account bundle for account_id"));
};
Ok(Json(FetchAccountResponse {
payload: BASE64.encode(&bundle.payload),
signature: BASE64.encode(&bundle.signature),
updated_at: bundle.updated_at,
}))
}
struct ApiError {
status: StatusCode,
message: String,
}
impl ApiError {
fn bad(msg: impl Into<String>) -> Self {
Self {
status: StatusCode::BAD_REQUEST,
message: msg.into(),
}
}
fn not_found(msg: impl Into<String>) -> Self {
Self {
status: StatusCode::NOT_FOUND,
message: msg.into(),
}
}
fn conflict(msg: impl Into<String>) -> Self {
Self {
status: StatusCode::CONFLICT,
message: msg.into(),
}
}
fn internal<E: std::fmt::Display>(err: E) -> Self {
tracing::error!("internal: {err}");
Self {
status: StatusCode::INTERNAL_SERVER_ERROR,
message: "internal error".into(),
}
}
}
impl IntoResponse for ApiError {
fn into_response(self) -> Response {
(
self.status,
Json(ErrorBody {
error: self.message,
}),
)
.into_response()
}
}

94
src/main.rs Normal file
View File

@ -0,0 +1,94 @@
//! Testnet KeyPackage Registry HTTP service.
//!
//! Throwaway service for issue #110 — replaced by λLEZ in v0.3. Intentionally
//! self-contained: depends only on axum + sqlite + ed25519, no libchat core.
//!
//! Wire:
//! POST /v0/keypackage — submit a signed keypackage bundle
//! GET /v0/keypackage/{device_id} — fetch the latest stored keypackage bundle
//! POST /v0/account — upsert a signed account device-list bundle
//! GET /v0/account/{account_id} — fetch the account device-list bundle
mod handlers;
mod store;
use std::net::SocketAddr;
use std::path::PathBuf;
use std::sync::Arc;
use std::time::Duration;
use anyhow::{Context, Result};
use clap::Parser;
use tracing_subscriber::EnvFilter;
use store::Store;
#[derive(Parser, Debug)]
#[command(name = "keypackage-registry", about = "Testnet KeyPackage Registry")]
struct Cli {
/// Address to bind the HTTP server.
#[arg(long, default_value = "0.0.0.0:8080")]
bind: SocketAddr,
/// SQLite database path.
#[arg(long, default_value = "keypackage-registry.db")]
db: PathBuf,
/// Maximum number of bundles retained per account_id.
#[arg(long, default_value_t = 100)]
max_per_identity: usize,
/// Retention window in days; older bundles are pruned.
#[arg(long, default_value_t = 30)]
retention_days: u64,
/// How often the prune task runs.
#[arg(long, default_value_t = 3600)]
prune_interval_secs: u64,
}
#[tokio::main]
async fn main() -> Result<()> {
tracing_subscriber::fmt()
.with_env_filter(
EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info")),
)
.init();
let cli = Cli::parse();
let store = Arc::new(Store::open(&cli.db).context("failed to open store")?);
let prune_store = store.clone();
let max_per_id = cli.max_per_identity;
let retention = Duration::from_secs(cli.retention_days * 24 * 3600);
let interval = Duration::from_secs(cli.prune_interval_secs);
tokio::spawn(async move {
let mut ticker = tokio::time::interval(interval);
loop {
ticker.tick().await;
if let Err(e) = prune_store.prune_key_packages(max_per_id, retention) {
tracing::warn!("prune (keypackages) failed: {e}");
}
if let Err(e) = prune_store.prune_accounts(retention) {
tracing::warn!("prune (accounts) failed: {e}");
}
}
});
let app = handlers::router(store);
let listener = tokio::net::TcpListener::bind(cli.bind)
.await
.with_context(|| format!("failed to bind {}", cli.bind))?;
tracing::info!("keypackage-registry listening on {}", cli.bind);
axum::serve(listener, app)
.with_graceful_shutdown(shutdown_signal())
.await
.context("server error")?;
Ok(())
}
async fn shutdown_signal() {
let _ = tokio::signal::ctrl_c().await;
tracing::info!("shutdown signal received");
}

309
src/store.rs Normal file
View File

@ -0,0 +1,309 @@
use std::path::Path;
use std::sync::Mutex;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use anyhow::{Context, Result};
use rusqlite::{Connection, OptionalExtension, params};
pub struct Store {
conn: Mutex<Connection>,
}
#[derive(Debug, Clone)]
pub struct StoredKeyPackageBundle {
/// The canonical signed payload, stored verbatim and returned as-is so
/// consumers verify over the exact bytes that were signed.
pub payload: Vec<u8>,
/// 64-byte Ed25519 signature over `payload`. Opaque to the server.
pub signature: Vec<u8>,
}
/// A signed bundle associating an account with its set of device (LocalIdentity)
/// public keys. The server stores exactly one blob per `account_id`; a newer
/// bundle replaces the old one only when its lamport is strictly higher (see
/// [`Store::upsert_account`]). `payload` is otherwise opaque to the server: it
/// encodes a lamport-timestamped list of device pubkeys signed by the account
/// key so that consumers can verify the full device set.
#[derive(Debug, Clone)]
pub struct StoredAccountBundle {
/// The canonical signed payload, returned verbatim so consumers can verify
/// the account signature over the exact bytes.
pub payload: Vec<u8>,
/// 64-byte Ed25519 signature over `payload` made by the account key.
pub signature: Vec<u8>,
/// Unix timestamp (ms) of the last upsert, stored for pruning.
pub updated_at: i64,
}
impl Store {
pub fn open(path: &Path) -> Result<Self> {
// Create the db's parent directory if the caller pointed at a nested
// path (e.g. `tmp/registry.db`); SQLite won't create it and errors with
// "unable to open database file" otherwise.
if let Some(parent) = path.parent()
&& !parent.as_os_str().is_empty()
{
std::fs::create_dir_all(parent)
.with_context(|| format!("create db directory {}", parent.display()))?;
}
let conn = Connection::open(path).context("open sqlite")?;
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS keypackages (
device_id TEXT NOT NULL,
received_at INTEGER NOT NULL,
payload BLOB NOT NULL,
signature BLOB NOT NULL,
PRIMARY KEY (device_id, received_at)
);
-- One row per account; newer upserts replace the existing row.
CREATE TABLE IF NOT EXISTS account_bundles (
account_id TEXT NOT NULL PRIMARY KEY,
updated_at INTEGER NOT NULL,
payload BLOB NOT NULL,
signature BLOB NOT NULL
);",
)?;
Ok(Self {
conn: Mutex::new(conn),
})
}
pub fn insert(&self, device_id: &str, bundle: &StoredKeyPackageBundle) -> Result<()> {
let received_at = now_ms() as i64;
let conn = self.conn.lock().unwrap();
conn.execute(
"INSERT INTO keypackages
(device_id, received_at, payload, signature)
VALUES (?1, ?2, ?3, ?4)",
params![device_id, received_at, bundle.payload, bundle.signature],
)?;
Ok(())
}
/// Returns the most recently received bundle for `device_id`. Scope A: the
/// chat layer consumes one bundle per device. When multi-keypackage fanout
/// lands, switch this to return a `Vec<StoredKeyPackageBundle>`.
pub fn latest(&self, device_id: &str) -> Result<Option<StoredKeyPackageBundle>> {
let conn = self.conn.lock().unwrap();
let row = conn
.query_row(
"SELECT payload, signature FROM keypackages
WHERE device_id = ?1
ORDER BY received_at DESC
LIMIT 1",
params![device_id],
|r| {
Ok(StoredKeyPackageBundle {
payload: r.get::<_, Vec<u8>>(0)?,
signature: r.get::<_, Vec<u8>>(1)?,
})
},
)
.optional()?;
Ok(row)
}
/// Upsert the signed device-list bundle for `account_id`. The server stores
/// exactly one blob per account.
///
/// Anti-replay: `lamport` is the monotonic version read from `bundle.payload`
/// (already signature-verified by the handler, so a forged value can't slip
/// past — the signature wouldn't match). The stored bundle is replaced only
/// when `lamport` is strictly greater than the one currently on file. A
/// replayed older-but-still-valid bundle therefore can't downgrade the device
/// list, and `updated_at` (the retention clock) is only bumped on a real
/// update so a replay can't keep a stale bundle alive past retention.
///
/// Returns `true` when the bundle was stored, `false` when it was rejected as
/// stale. The compare-and-swap runs under the connection lock so concurrent
/// publishes can't interleave a read with a write. The `updated_at` field of
/// `bundle` is ignored; the store stamps the row with the current time.
pub fn upsert_account(
&self,
account_id: &str,
lamport: u64,
bundle: &StoredAccountBundle,
) -> Result<bool> {
let updated_at = now_ms() as i64;
let conn = self.conn.lock().unwrap();
let existing_lamport = conn
.query_row(
"SELECT payload FROM account_bundles WHERE account_id = ?1",
params![account_id],
|r| r.get::<_, Vec<u8>>(0),
)
.optional()?
.and_then(|payload| payload_lamport(&payload));
if let Some(stored) = existing_lamport
&& lamport <= stored
{
return Ok(false);
}
conn.execute(
"INSERT INTO account_bundles (account_id, updated_at, payload, signature)
VALUES (?1, ?2, ?3, ?4)
ON CONFLICT(account_id) DO UPDATE SET
updated_at = excluded.updated_at,
payload = excluded.payload,
signature = excluded.signature",
params![account_id, updated_at, bundle.payload, bundle.signature],
)?;
Ok(true)
}
/// Returns the stored bundle for `account_id`, or `None` if unknown.
pub fn get_account(&self, account_id: &str) -> Result<Option<StoredAccountBundle>> {
let conn = self.conn.lock().unwrap();
let row = conn
.query_row(
"SELECT payload, signature, updated_at FROM account_bundles
WHERE account_id = ?1",
params![account_id],
|r| {
Ok(StoredAccountBundle {
payload: r.get::<_, Vec<u8>>(0)?,
signature: r.get::<_, Vec<u8>>(1)?,
updated_at: r.get::<_, i64>(2)?,
})
},
)
.optional()?;
Ok(row)
}
/// Drops account bundles that have not been refreshed within `retention`.
pub fn prune_accounts(&self, retention: Duration) -> Result<()> {
let cutoff_ms = now_ms().saturating_sub(retention.as_millis() as u64) as i64;
let conn = self.conn.lock().unwrap();
conn.execute(
"DELETE FROM account_bundles WHERE updated_at < ?1",
params![cutoff_ms],
)?;
Ok(())
}
/// Drops bundles older than `retention` and keeps at most
/// `max_per_identity` per `device_id` — each device's history is bounded
/// independently.
pub fn prune_key_packages(&self, max_per_identity: usize, retention: Duration) -> Result<()> {
let cutoff_ms = now_ms().saturating_sub(retention.as_millis() as u64) as i64;
let conn = self.conn.lock().unwrap();
conn.execute(
"DELETE FROM keypackages WHERE received_at < ?1",
params![cutoff_ms],
)?;
conn.execute(
"DELETE FROM keypackages
WHERE rowid IN (
SELECT rowid FROM (
SELECT rowid,
ROW_NUMBER() OVER (
PARTITION BY device_id
ORDER BY received_at DESC
) AS rn
FROM keypackages
)
WHERE rn > ?1
)",
params![max_per_identity as i64],
)?;
Ok(())
}
}
/// Domain-separation prefix on every account-device-bundle payload. Must stay in
/// sync with `account_directory::BUNDLE_DOMAIN` in the conversations crate; this
/// throwaway service deliberately has no libchat-core dependency, so the constant
/// is duplicated here rather than imported.
const BUNDLE_DOMAIN: &[u8] = b"libchat:account-device-bundle\0";
/// Extract the lamport version from a bundle payload without otherwise
/// interpreting it. The canonical layout (owned by the conversations crate's
/// `encode_bundle_payload`) is `domain | version:u8 | lamport:u64 LE | …`, so the
/// lamport sits in the 8 bytes right after the domain prefix and version byte.
/// Returns `None` when the domain prefix is absent or the payload is too short to
/// contain a header — the handler treats either as a malformed request.
pub fn payload_lamport(payload: &[u8]) -> Option<u64> {
payload
.strip_prefix(BUNDLE_DOMAIN)?
.get(1..9)
.map(|b| u64::from_le_bytes(b.try_into().expect("1..9 is 8 bytes")))
}
fn now_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_millis() as u64
}
#[cfg(test)]
mod tests {
use super::*;
/// Minimal stand-in for a real bundle payload: the domain prefix plus the
/// header fields the server reads (`version:u8 | lamport:u64 LE`), no device
/// keys needed.
fn payload_with_lamport(lamport: u64) -> Vec<u8> {
let mut p = BUNDLE_DOMAIN.to_vec();
p.push(1u8); // version
p.extend_from_slice(&lamport.to_le_bytes());
p
}
fn bundle(lamport: u64) -> StoredAccountBundle {
StoredAccountBundle {
payload: payload_with_lamport(lamport),
signature: vec![0u8; 64],
updated_at: 0,
}
}
fn upsert(store: &Store, account: &str, lamport: u64) -> bool {
store
.upsert_account(account, lamport, &bundle(lamport))
.unwrap()
}
#[test]
fn rejects_replayed_or_stale_lamport() {
let store = Store::open(Path::new(":memory:")).unwrap();
// First publish is always accepted.
assert!(upsert(&store, "acct", 5));
// A strictly higher lamport replaces it.
assert!(upsert(&store, "acct", 6));
// Re-publishing the same lamport (a replay) is rejected.
assert!(!upsert(&store, "acct", 6));
// An older lamport (a downgrade) is rejected.
assert!(!upsert(&store, "acct", 4));
// The stored bundle is still the newest one accepted.
let stored = store.get_account("acct").unwrap().unwrap();
assert_eq!(payload_lamport(&stored.payload), Some(6));
}
#[test]
fn stale_publish_does_not_refresh_retention_clock() {
let store = Store::open(Path::new(":memory:")).unwrap();
assert!(upsert(&store, "acct", 9));
let after_first = store.get_account("acct").unwrap().unwrap().updated_at;
// A rejected (stale) publish must not bump updated_at, so a replay can't
// keep a stale bundle alive past the retention window.
assert!(!upsert(&store, "acct", 9));
let after_replay = store.get_account("acct").unwrap().unwrap().updated_at;
assert_eq!(after_first, after_replay);
}
#[test]
fn payload_lamport_requires_domain_and_full_header() {
assert_eq!(payload_lamport(&payload_with_lamport(42)), Some(42));
// Missing the domain prefix → unparseable.
assert_eq!(payload_lamport(&[1u8, 0, 0, 0, 0, 0, 0, 0, 0]), None);
// Has the domain but is too short for version + u64 → unparseable.
let mut short = BUNDLE_DOMAIN.to_vec();
short.push(1u8);
assert_eq!(payload_lamport(&short), None);
}
}