mirror of
https://github.com/logos-messaging/chat-store.git
synced 2026-07-29 14:33:20 +00:00
feat: http server for account and keypackages storage (#2)
* feat: http server for account and keypackages storage * chore: rename package * feat: change to sqlx * smoke test script * chore: use account pub in code * fix dockerfile
This commit is contained in:
parent
cf09036f42
commit
4711fe8547
7
.dockerignore
Normal file
7
.dockerignore
Normal file
@ -0,0 +1,7 @@
|
||||
target
|
||||
*.db
|
||||
*.db-shm
|
||||
*.db-wal
|
||||
.git
|
||||
.gitignore
|
||||
.DS_Store
|
||||
10
.gitignore
vendored
Normal file
10
.gitignore
vendored
Normal 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
|
||||
2354
Cargo.lock
generated
Normal file
2354
Cargo.lock
generated
Normal file
File diff suppressed because it is too large
Load Diff
37
Cargo.toml
Normal file
37
Cargo.toml
Normal file
@ -0,0 +1,37 @@
|
||||
# Standalone single-crate workspace: keeps this repo independent of any parent
|
||||
# Cargo workspace it might be checked out beneath.
|
||||
[workspace]
|
||||
|
||||
[package]
|
||||
name = "chat-store"
|
||||
version = "0.1.0"
|
||||
edition = "2024"
|
||||
|
||||
[[bin]]
|
||||
name = "chat-store"
|
||||
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"
|
||||
serde = { version = "1.0", features = ["derive"] }
|
||||
serde_json = "1.0"
|
||||
# SQLite via sqlx with a bundled libsqlite3 — no sqlcipher/OpenSSL, no system libs.
|
||||
sqlx = { version = "0.8", default-features = false, features = [
|
||||
"runtime-tokio",
|
||||
"sqlite",
|
||||
"macros",
|
||||
"migrate",
|
||||
] }
|
||||
tokio = { version = "1", features = ["rt-multi-thread", "macros", "signal", "sync", "time"] }
|
||||
tracing = "0.1"
|
||||
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
|
||||
|
||||
[dev-dependencies]
|
||||
# `smoke_test` example only talks to a localhost HTTP server, so no TLS backend
|
||||
# (and no OpenSSL) is pulled in.
|
||||
reqwest = { version = "0.12", default-features = false, features = ["blocking", "json"] }
|
||||
44
Dockerfile
Normal file
44
Dockerfile
Normal file
@ -0,0 +1,44 @@
|
||||
# syntax=docker/dockerfile:1
|
||||
|
||||
########################################
|
||||
# Build stage
|
||||
########################################
|
||||
FROM rust:1-bookworm AS builder
|
||||
|
||||
WORKDIR /app
|
||||
|
||||
# Build dependencies first against a stub binary so the (slow) dependency
|
||||
# compilation — bundled libsqlite3 and the async stack — 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/chat_store* target/release/chat-store
|
||||
|
||||
# Now build the real binary; dependency artifacts above are reused. `migrations/`
|
||||
# is needed at build time too — `sqlx::migrate!()` embeds the .sql files.
|
||||
COPY src ./src
|
||||
COPY migrations ./migrations
|
||||
RUN cargo build --release --locked --bin chat-store
|
||||
|
||||
########################################
|
||||
# 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/chat-store /usr/local/bin/chat-store
|
||||
|
||||
# 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 ["chat-store"]
|
||||
CMD ["--bind", "0.0.0.0:8080", "--db", "/data/chat-store.db"]
|
||||
221
README.md
221
README.md
@ -1,3 +1,222 @@
|
||||
# Chat Store
|
||||
|
||||
For persistence of group chat users' key package.
|
||||
Persistence for group-chat users' key packages — the **chat-store** HTTP service
|
||||
(formerly `keypackage-registry`), 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_pub`**
|
||||
mapping an Account to its set of device (LocalIdentity) public keys, so clients
|
||||
can invite every LocalIdentity of an account. `account_pub` 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/chat-store # binds 0.0.0.0:8080, db ./chat-store.db
|
||||
```
|
||||
|
||||
| Flag | Default | Description |
|
||||
|------|---------|-------------|
|
||||
| `--bind <addr>` | `0.0.0.0:8080` | HTTP bind address |
|
||||
| `--db <path>` | `chat-store.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/chat-store.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_pub`** mapping an
|
||||
Account to its LocalIdentity device keys. Same trust model as keypackages: the
|
||||
server verifies `signature` over `payload` under `account_pub`'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_pub": "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_pub}`
|
||||
|
||||
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_pub`'s key, then decode the device list.
|
||||
|
||||
## Storage & retention
|
||||
|
||||
Two SQLite tables: `keypackages` keyed by `device_id`, and `account_bundles`
|
||||
(one row per `account_pub`). 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
|
||||
|
||||
The quickest end-to-end check is the bundled [`smoke_test`](examples/smoke_test.rs)
|
||||
example. It generates throwaway Ed25519 keys, signs and publishes a keypackage and
|
||||
an account bundle, fetches both back, and confirms the replay guard:
|
||||
|
||||
```bash
|
||||
# Terminal 1 — start a server with a fresh db
|
||||
cargo run -- --bind 127.0.0.1:8080 --db tmp/chat-store.db
|
||||
|
||||
# Terminal 2 — run the example against it (defaults to http://127.0.0.1:8080)
|
||||
cargo run --example smoke_test
|
||||
# or point it elsewhere:
|
||||
cargo run --example smoke_test -- http://127.0.0.1:8080
|
||||
```
|
||||
|
||||
Expected output:
|
||||
|
||||
```text
|
||||
POST /v0/keypackage -> 204 No Content (expect 204)
|
||||
GET /v0/keypackage/<id> -> 200 OK (expect 200) {"payload":...,"signature":...}
|
||||
POST /v0/account -> 204 No Content (expect 204)
|
||||
GET /v0/account/<id> -> 200 OK (expect 200) {"payload":...,"signature":...,"updated_at":...}
|
||||
POST /v0/account (replay) -> 409 Conflict (expect 409)
|
||||
```
|
||||
|
||||
You can also exercise it 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.
|
||||
|
||||
129
examples/smoke_test.rs
Normal file
129
examples/smoke_test.rs
Normal file
@ -0,0 +1,129 @@
|
||||
//! End-to-end smoke test for the chat-store HTTP API.
|
||||
//!
|
||||
//! Generates throwaway Ed25519 keys, signs a keypackage bundle and an account
|
||||
//! device-list bundle, POSTs both, then GETs them back. The server verifies the
|
||||
//! signature over the exact payload bytes, so the bundles must be properly
|
||||
//! signed — this example does that for you.
|
||||
//!
|
||||
//! Run it against a live server:
|
||||
//!
|
||||
//! ```text
|
||||
//! cargo run -- --bind 127.0.0.1:8080 --db tmp/chat-store.db # terminal 1
|
||||
//! cargo run --example smoke_test # terminal 2
|
||||
//! cargo run --example smoke_test -- http://host:port # custom target
|
||||
//! ```
|
||||
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
|
||||
use anyhow::Result;
|
||||
use base64::Engine;
|
||||
use base64::engine::general_purpose::STANDARD as BASE64;
|
||||
use ed25519_dalek::{Signer, SigningKey};
|
||||
use reqwest::blocking::Client;
|
||||
use serde_json::json;
|
||||
|
||||
/// Domain-separation prefix the server expects on every account bundle payload.
|
||||
/// Must match `BUNDLE_DOMAIN` in `src/store.rs` (and libchat's account_directory).
|
||||
const ACCOUNT_BUNDLE_DOMAIN: &[u8] = b"libchat:account-device-bundle\0";
|
||||
|
||||
fn main() -> Result<()> {
|
||||
let base = std::env::args()
|
||||
.nth(1)
|
||||
.unwrap_or_else(|| "http://127.0.0.1:8080".to_string());
|
||||
let base = base.trim_end_matches('/');
|
||||
let client = Client::new();
|
||||
|
||||
println!("Testing chat-store at {base}");
|
||||
test_keypackage(&client, base)?;
|
||||
test_account(&client, base)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn test_keypackage(client: &Client, base: &str) -> Result<()> {
|
||||
let key = signing_key(1);
|
||||
let device_id = pub_hex(&key);
|
||||
|
||||
// The keypackage payload is opaque to the server: any bytes work as long as
|
||||
// the signature matches. Real clients put `timestamp_ms_le[8] || key_package`.
|
||||
let mut payload = 0u64.to_le_bytes().to_vec();
|
||||
payload.extend_from_slice(b"hello-keypackage");
|
||||
|
||||
let resp = client
|
||||
.post(format!("{base}/v0/keypackage"))
|
||||
.json(&json!({
|
||||
"device_id": device_id,
|
||||
"payload": BASE64.encode(&payload),
|
||||
"signature": BASE64.encode(key.sign(&payload).to_bytes()),
|
||||
}))
|
||||
.send()?;
|
||||
println!("POST /v0/keypackage -> {} (expect 204)", resp.status());
|
||||
|
||||
let resp = client
|
||||
.get(format!("{base}/v0/keypackage/{device_id}"))
|
||||
.send()?;
|
||||
println!(
|
||||
"GET /v0/keypackage/<id> -> {} (expect 200) {}",
|
||||
resp.status(),
|
||||
resp.text()?
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn test_account(client: &Client, base: &str) -> Result<()> {
|
||||
let key = signing_key(2);
|
||||
let account_pub = pub_hex(&key);
|
||||
|
||||
// The account payload is NOT arbitrary: it must start with the domain prefix,
|
||||
// a version byte, and an 8-byte little-endian lamport so the server can run
|
||||
// its replay check. Device pubkeys would follow, but stay opaque to the server.
|
||||
let lamport: u64 = 1;
|
||||
let mut payload = ACCOUNT_BUNDLE_DOMAIN.to_vec();
|
||||
payload.push(1); // version
|
||||
payload.extend_from_slice(&lamport.to_le_bytes());
|
||||
|
||||
let body = json!({
|
||||
"account_pub": account_pub,
|
||||
"payload": BASE64.encode(&payload),
|
||||
"signature": BASE64.encode(key.sign(&payload).to_bytes()),
|
||||
});
|
||||
|
||||
let resp = client
|
||||
.post(format!("{base}/v0/account"))
|
||||
.json(&body)
|
||||
.send()?;
|
||||
println!("POST /v0/account -> {} (expect 204)", resp.status());
|
||||
|
||||
let resp = client.get(format!("{base}/v0/account/{account_pub}")).send()?;
|
||||
println!(
|
||||
"GET /v0/account/<id> -> {} (expect 200) {}",
|
||||
resp.status(),
|
||||
resp.text()?
|
||||
);
|
||||
|
||||
// Re-posting the same lamport must be rejected as a stale replay.
|
||||
let resp = client
|
||||
.post(format!("{base}/v0/account"))
|
||||
.json(&body)
|
||||
.send()?;
|
||||
println!("POST /v0/account (replay) -> {} (expect 409)", resp.status());
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// A throwaway Ed25519 signing key seeded from the current time (plus a salt so
|
||||
/// the keypackage and account keys differ), so repeated runs use fresh ids.
|
||||
fn signing_key(salt: u8) -> SigningKey {
|
||||
let nanos = SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.unwrap_or_default()
|
||||
.as_nanos()
|
||||
.to_le_bytes();
|
||||
let mut seed = [0u8; 32];
|
||||
for (i, b) in seed.iter_mut().enumerate() {
|
||||
*b = nanos[i % nanos.len()] ^ salt ^ (i as u8);
|
||||
}
|
||||
SigningKey::from_bytes(&seed)
|
||||
}
|
||||
|
||||
fn pub_hex(key: &SigningKey) -> String {
|
||||
hex::encode(key.verifying_key().to_bytes())
|
||||
}
|
||||
18
migrations/0001_init.sql
Normal file
18
migrations/0001_init.sql
Normal file
@ -0,0 +1,18 @@
|
||||
-- KeyPackage bundles: history per device, newest read back on fetch.
|
||||
CREATE TABLE 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)
|
||||
);
|
||||
|
||||
-- Account device-list bundles: exactly one row per account; newer upserts
|
||||
-- replace the existing row (compare-and-swap on the payload's lamport). Keyed by
|
||||
-- the hex-encoded account verifying key.
|
||||
CREATE TABLE account_bundles (
|
||||
account_pub TEXT NOT NULL PRIMARY KEY,
|
||||
updated_at INTEGER NOT NULL,
|
||||
payload BLOB NOT NULL,
|
||||
signature BLOB NOT NULL
|
||||
);
|
||||
255
src/handlers.rs
Normal file
255
src/handlers.rs
Normal file
@ -0,0 +1,255 @@
|
||||
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_pub", 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(),
|
||||
},
|
||||
)
|
||||
.await
|
||||
.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)
|
||||
.await
|
||||
.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_pub`.
|
||||
#[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_pub: 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_pub`, 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_pub)
|
||||
.ok()
|
||||
.and_then(|b| b.try_into().ok())
|
||||
.ok_or_else(|| ApiError::bad("account_pub: 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_pub: 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_pub,
|
||||
lamport,
|
||||
&StoredAccountBundle {
|
||||
payload,
|
||||
signature: signature.to_vec(),
|
||||
updated_at: 0, // filled in by store
|
||||
},
|
||||
)
|
||||
.await
|
||||
.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_pub` — 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_pub): Path<String>,
|
||||
) -> Result<Json<FetchAccountResponse>, ApiError> {
|
||||
let Some(bundle) = store
|
||||
.get_account(&account_pub)
|
||||
.await
|
||||
.map_err(ApiError::internal)?
|
||||
else {
|
||||
return Err(ApiError::not_found("no account bundle for account_pub"));
|
||||
};
|
||||
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()
|
||||
}
|
||||
}
|
||||
98
src/main.rs
Normal file
98
src/main.rs
Normal file
@ -0,0 +1,98 @@
|
||||
//! 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_pub} — 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 = "chat-store", about = "Testnet Chat Store (KeyPackage + account directory)")]
|
||||
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 = "chat-store.db")]
|
||||
db: PathBuf,
|
||||
|
||||
/// Maximum number of keypackage bundles retained per device_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)
|
||||
.await
|
||||
.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).await {
|
||||
tracing::warn!("prune (keypackages) failed: {e}");
|
||||
}
|
||||
if let Err(e) = prune_store.prune_accounts(retention).await {
|
||||
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!("chat-store 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");
|
||||
}
|
||||
316
src/store.rs
Normal file
316
src/store.rs
Normal file
@ -0,0 +1,316 @@
|
||||
use std::path::Path;
|
||||
use std::time::{Duration, SystemTime, UNIX_EPOCH};
|
||||
|
||||
use anyhow::{Context, Result};
|
||||
use sqlx::SqlitePool;
|
||||
use sqlx::sqlite::{SqliteConnectOptions, SqliteJournalMode, SqlitePoolOptions};
|
||||
|
||||
pub struct Store {
|
||||
pool: SqlitePool,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, sqlx::FromRow)]
|
||||
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_pub`; 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, sqlx::FromRow)]
|
||||
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 async fn open(path: &Path) -> Result<Self> {
|
||||
let is_memory = path == Path::new(":memory:");
|
||||
|
||||
let mut options = SqliteConnectOptions::new()
|
||||
.create_if_missing(true)
|
||||
.busy_timeout(Duration::from_secs(5));
|
||||
|
||||
if is_memory {
|
||||
options = options.filename(":memory:");
|
||||
} else {
|
||||
// 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()))?;
|
||||
}
|
||||
options = options.filename(path).journal_mode(SqliteJournalMode::Wal);
|
||||
}
|
||||
|
||||
// A shared in-memory database only lives as long as a connection is open,
|
||||
// so pin the `:memory:` pool to a single reused connection; file-backed
|
||||
// pools can fan out across requests and the prune task.
|
||||
let pool = SqlitePoolOptions::new()
|
||||
.max_connections(if is_memory { 1 } else { 5 })
|
||||
.connect_with(options)
|
||||
.await
|
||||
.context("open sqlite")?;
|
||||
|
||||
sqlx::migrate!()
|
||||
.run(&pool)
|
||||
.await
|
||||
.context("run migrations")?;
|
||||
|
||||
Ok(Self { pool })
|
||||
}
|
||||
|
||||
pub async fn insert(&self, device_id: &str, bundle: &StoredKeyPackageBundle) -> Result<()> {
|
||||
let received_at = now_ms() as i64;
|
||||
sqlx::query(
|
||||
"INSERT INTO keypackages (device_id, received_at, payload, signature)
|
||||
VALUES (?, ?, ?, ?)",
|
||||
)
|
||||
.bind(device_id)
|
||||
.bind(received_at)
|
||||
.bind(&bundle.payload)
|
||||
.bind(&bundle.signature)
|
||||
.execute(&self.pool)
|
||||
.await?;
|
||||
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 async fn latest(&self, device_id: &str) -> Result<Option<StoredKeyPackageBundle>> {
|
||||
let row = sqlx::query_as::<_, StoredKeyPackageBundle>(
|
||||
"SELECT payload, signature FROM keypackages
|
||||
WHERE device_id = ?
|
||||
ORDER BY received_at DESC
|
||||
LIMIT 1",
|
||||
)
|
||||
.bind(device_id)
|
||||
.fetch_optional(&self.pool)
|
||||
.await?;
|
||||
Ok(row)
|
||||
}
|
||||
|
||||
/// Upsert the signed device-list bundle for `account_pub`. 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 inside a transaction so a concurrent
|
||||
/// publish can't interleave the read with the write. The `updated_at` field of
|
||||
/// `bundle` is ignored; the store stamps the row with the current time.
|
||||
pub async fn upsert_account(
|
||||
&self,
|
||||
account_pub: &str,
|
||||
lamport: u64,
|
||||
bundle: &StoredAccountBundle,
|
||||
) -> Result<bool> {
|
||||
let updated_at = now_ms() as i64;
|
||||
|
||||
let mut tx = self.pool.begin().await?;
|
||||
let existing_lamport = sqlx::query_scalar::<_, Vec<u8>>(
|
||||
"SELECT payload FROM account_bundles WHERE account_pub = ?",
|
||||
)
|
||||
.bind(account_pub)
|
||||
.fetch_optional(&mut *tx)
|
||||
.await?
|
||||
.and_then(|payload| payload_lamport(&payload));
|
||||
if let Some(stored) = existing_lamport
|
||||
&& lamport <= stored
|
||||
{
|
||||
// Dropping `tx` rolls the (read-only) transaction back.
|
||||
return Ok(false);
|
||||
}
|
||||
sqlx::query(
|
||||
"INSERT INTO account_bundles (account_pub, updated_at, payload, signature)
|
||||
VALUES (?, ?, ?, ?)
|
||||
ON CONFLICT(account_pub) DO UPDATE SET
|
||||
updated_at = excluded.updated_at,
|
||||
payload = excluded.payload,
|
||||
signature = excluded.signature",
|
||||
)
|
||||
.bind(account_pub)
|
||||
.bind(updated_at)
|
||||
.bind(&bundle.payload)
|
||||
.bind(&bundle.signature)
|
||||
.execute(&mut *tx)
|
||||
.await?;
|
||||
tx.commit().await?;
|
||||
Ok(true)
|
||||
}
|
||||
|
||||
/// Returns the stored bundle for `account_pub`, or `None` if unknown.
|
||||
pub async fn get_account(&self, account_pub: &str) -> Result<Option<StoredAccountBundle>> {
|
||||
let row = sqlx::query_as::<_, StoredAccountBundle>(
|
||||
"SELECT payload, signature, updated_at FROM account_bundles
|
||||
WHERE account_pub = ?",
|
||||
)
|
||||
.bind(account_pub)
|
||||
.fetch_optional(&self.pool)
|
||||
.await?;
|
||||
Ok(row)
|
||||
}
|
||||
|
||||
/// Drops account bundles that have not been refreshed within `retention`.
|
||||
pub async fn prune_accounts(&self, retention: Duration) -> Result<()> {
|
||||
let cutoff_ms = now_ms().saturating_sub(retention.as_millis() as u64) as i64;
|
||||
sqlx::query("DELETE FROM account_bundles WHERE updated_at < ?")
|
||||
.bind(cutoff_ms)
|
||||
.execute(&self.pool)
|
||||
.await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Drops bundles older than `retention` and keeps at most
|
||||
/// `max_per_identity` per `device_id` — each device's history is bounded
|
||||
/// independently.
|
||||
pub async 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;
|
||||
sqlx::query("DELETE FROM keypackages WHERE received_at < ?")
|
||||
.bind(cutoff_ms)
|
||||
.execute(&self.pool)
|
||||
.await?;
|
||||
sqlx::query(
|
||||
"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 > ?
|
||||
)",
|
||||
)
|
||||
.bind(max_per_identity as i64)
|
||||
.execute(&self.pool)
|
||||
.await?;
|
||||
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,
|
||||
}
|
||||
}
|
||||
|
||||
async fn upsert(store: &Store, account: &str, lamport: u64) -> bool {
|
||||
store
|
||||
.upsert_account(account, lamport, &bundle(lamport))
|
||||
.await
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn rejects_replayed_or_stale_lamport() {
|
||||
let store = Store::open(Path::new(":memory:")).await.unwrap();
|
||||
|
||||
// First publish is always accepted.
|
||||
assert!(upsert(&store, "acct", 5).await);
|
||||
// A strictly higher lamport replaces it.
|
||||
assert!(upsert(&store, "acct", 6).await);
|
||||
// Re-publishing the same lamport (a replay) is rejected.
|
||||
assert!(!upsert(&store, "acct", 6).await);
|
||||
// An older lamport (a downgrade) is rejected.
|
||||
assert!(!upsert(&store, "acct", 4).await);
|
||||
|
||||
// The stored bundle is still the newest one accepted.
|
||||
let stored = store.get_account("acct").await.unwrap().unwrap();
|
||||
assert_eq!(payload_lamport(&stored.payload), Some(6));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn stale_publish_does_not_refresh_retention_clock() {
|
||||
let store = Store::open(Path::new(":memory:")).await.unwrap();
|
||||
assert!(upsert(&store, "acct", 9).await);
|
||||
let after_first = store.get_account("acct").await.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).await);
|
||||
let after_replay = store.get_account("acct").await.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);
|
||||
}
|
||||
}
|
||||
Loading…
x
Reference in New Issue
Block a user