Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
152 changes: 152 additions & 0 deletions sdk/rust/src/signal/maintain.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,152 @@
//! Idempotent key-bundle maintenance.
//!
//! Reconciles an agent's *published* pre-key bundle (on the relay) with the
//! private material in its [`SessionStore`], generating and uploading keys ONLY
//! when the relay actually needs them.
//!
//! This is the anti-orphan invariant. Blindly republishing a fresh bundle on
//! every start leaves the relay advertising one-time pre-keys whose private
//! halves a rotated or wiped store can no longer back, so a peer's X3DH against
//! that bundle fails with a MAC error and its first message is silently dropped.
//! Storing the private material locally *before* uploading, and skipping the
//! upload entirely when the store already backs a healthy published set,
//! guarantees the relay never serves a key this agent cannot answer.

use std::time::{SystemTime, UNIX_EPOCH};

use crate::api::keys::KeysApi;
use crate::error::Result;
use crate::signal::keys::{generate_pre_keys, generate_signed_pre_key, serialize_pre_key};
use crate::signal::store::SessionStore;
use crate::signer::Signer;
use crate::types::{PreKeysRequest, SignedPreKeyRequest};

/// Tuning for [`maintain_keys`].
#[derive(Debug, Clone)]
pub struct MaintainPolicy {
/// How many one-time pre-keys to (re)publish in a batch when the pool needs
/// filling.
pub one_time_batch: usize,
}

impl Default for MaintainPolicy {
fn default() -> Self {
Self { one_time_batch: 20 }
}
}

/// What [`maintain_keys`] did on a call.
#[derive(Debug, Default, Clone, PartialEq, Eq)]
pub struct KeyMaintenance {
/// Nothing needed publishing: the store backed the signed pre-key the relay
/// advertises and its one-time pool was healthy — the common restart path.
pub was_healthy: bool,
/// A fresh signed pre-key was generated and rotated on the relay, because the
/// store could not answer the id the relay was advertising.
pub rotated_signed: bool,
/// How many one-time pre-keys were generated and uploaded.
pub uploaded_one_time: usize,
}

/// Idempotently reconcile `agent_id`'s published key bundle with the relay.
///
/// Safe to call on every boot and periodically — it publishes only what the
/// relay actually needs. The two halves are maintained **independently**:
///
/// - **Signed pre-key** — rotated only when the store cannot answer the id the
/// relay advertises (`KeyHealth::signed_pre_key_key_id`). An inbound
/// `PREKEY_BUNDLE` names that exact id, so a relay serving a key we no longer
/// hold breaks first contact; a still-backed signed key is left untouched.
/// - **One-time pre-keys** — topped up only when the relay reports the pool low
/// or empty, without disturbing a healthy signed pre-key.
///
/// So a merely-depleted one-time pool refills without needless signed-key
/// churn, and a stale/unbacked signed key is repaired even while the pool is
/// full. When both are satisfied the call is a **no-op**
/// ([`KeyMaintenance::was_healthy`]) — not republishing is what prevents
/// orphaning. If the health check itself fails, both halves are (re)published
/// rather than risk leaving the agent unpublished.
///
/// `identity_key_base64` is the agent's base64 Ed25519 identity key; the relay
/// verifies each published key's signature against it.
///
/// # Errors
///
/// Returns the underlying error if key generation, a store write, the signed
/// pre-key rotation, or the one-time upload fails. A failed *health* check is
/// treated as "cannot confirm healthy" and falls through to a (re)publish rather
/// than erroring, so a transient relay blip never leaves the agent unpublished.
pub async fn maintain_keys(
keys: &KeysApi,
store: &dyn SessionStore,
signer: &dyn Signer,
agent_id: &str,
identity_key_base64: &str,
policy: &MaintainPolicy,
) -> Result<KeyMaintenance> {
// What is the relay actually serving right now? If we cannot reach it we
// cannot confirm anything is healthy, so both halves are (re)published
// rather than risk leaving the agent unpublished.
let health = keys.health(agent_id).await.ok();

let now_secs = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
let mut report = KeyMaintenance::default();

// --- Signed pre-key -----------------------------------------------------
// Rotate ONLY when the store cannot answer the signed pre-key the relay
// advertises. An inbound `PREKEY_BUNDLE` names that exact id, so a relay
// serving an id we no longer hold makes first contact fail; conversely a
// still-backed signed key is left alone even when the one-time pool needs a
// refill (the two halves are maintained independently).
let signed_pre_key_backed = match health
.as_ref()
.and_then(|h| h.signed_pre_key_key_id.as_deref())
{
Some(key_id) => store.signed_pre_key(key_id).await.ok().flatten().is_some(),
// No health, or the relay advertises no signed pre-key at all.
None => false,
};
if !signed_pre_key_backed {
let spk = generate_signed_pre_key(signer, &format!("spk_{now_secs}")).await?;
store.store_signed_pre_key(spk.clone()).await?;
keys.rotate_signed_pre_key(
agent_id,
&SignedPreKeyRequest {
identity_key: Some(identity_key_base64.to_string()),
signed_pre_key: serialize_pre_key(&spk),
},
)
.await?;
report.rotated_signed = true;
}

// --- One-time pre-keys --------------------------------------------------
// Topped up independently, driven only by the relay's own pool report. The
// private halves are stored BEFORE upload so the relay never advertises a
// key the store cannot back.
let one_time_depleted = match health.as_ref() {
Some(h) => h.low_one_time_pre_keys || h.one_time_pre_key_count <= 0,
None => true,
};
if one_time_depleted {
let one_time = generate_pre_keys(signer, now_secs, policy.one_time_batch).await?;
for key in &one_time {
store.store_pre_key(key.clone()).await?;
}
keys.upload_pre_keys(
agent_id,
&PreKeysRequest {
identity_key: Some(identity_key_base64.to_string()),
pre_keys: one_time.iter().map(serialize_pre_key).collect(),
},
)
.await?;
report.uploaded_one_time = one_time.len();
}

report.was_healthy = !report.rotated_signed && report.uploaded_one_time == 0;
Ok(report)
}
1 change: 1 addition & 0 deletions sdk/rust/src/signal/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@

pub mod crypto;
pub mod keys;
pub mod maintain;
pub mod memory_store;
pub mod ratchet;
pub mod sender_key;
Expand Down
177 changes: 177 additions & 0 deletions sdk/rust/tests/signal_maintain.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,177 @@
//! Coverage for idempotent key-bundle maintenance
//! ([`tinyplace::signal::maintain::maintain_keys`]).
//!
//! The signed pre-key and the one-time pool are maintained **independently**, so
//! each test pins one half healthy and the other degraded to prove neither
//! drags the other into a needless republish.

mod common;

use common::{all_requests, client_for, test_signer};
use serde_json::{json, Value};
use wiremock::matchers::{method, path_regex};
use wiremock::{Mock, MockServer, ResponseTemplate};

use tinyplace::signal::crypto::generate_x25519_keypair;
use tinyplace::signal::keys::generate_signed_pre_key;
use tinyplace::signal::maintain::{maintain_keys, MaintainPolicy};
use tinyplace::signal::memory_store::MemorySessionStore;
use tinyplace::signal::store::SessionStore;
use tinyplace::{LocalSigner, Signer};

/// A relay mock whose `GET /health` reports `one_time_count` one-time pre-keys
/// and advertises `signed_id` as its signed pre-key, and whose `PUT`s (rotate /
/// upload) return `null`.
async fn relay(one_time_count: i64, signed_id: Value) -> MockServer {
let server = MockServer::start().await;
// Pinned to the real health route (the id segment is the agent's base58 key),
// so a wrong path would fail the scenario instead of silently matching.
Mock::given(method("GET"))
.and(path_regex(r"^/keys/[^/]+/health$"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"agentId": "x",
"oneTimePreKeyCount": one_time_count,
"lowOneTimePreKeys": one_time_count < 5,
"signedPreKeyKeyId": signed_id,
})))
.mount(&server)
.await;
Mock::given(method("PUT"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!(null)))
.mount(&server)
.await;
server
}

/// A store already backing a signed pre-key under `key_id`.
async fn store_backing(signer: &LocalSigner, key_id: &str) -> MemorySessionStore {
let store = MemorySessionStore::new(generate_x25519_keypair());
let spk = generate_signed_pre_key(signer, key_id).await.unwrap();
store.store_signed_pre_key(spk).await.unwrap();
store
}

/// Count the `PUT`s that hit each key route.
async fn puts(server: &MockServer) -> (usize, usize) {
let requests = all_requests(server).await;
let count = |suffix: &str| {
requests
.iter()
.filter(|r| r.method.as_str() == "PUT" && r.url.path().ends_with(suffix))
.count()
};
(count("/signed-prekey"), count("/prekeys"))
}

#[tokio::test]
async fn no_ops_when_the_advertised_signed_key_is_backed_and_the_pool_is_healthy() {
let signer = test_signer();
let server = relay(20, json!("spk_known")).await;
let client = client_for(&server);
let store = store_backing(&signer, "spk_known").await;

let report = maintain_keys(
&client.keys,
&store,
signer.as_ref(),
&signer.agent_id(),
"identity-key",
&MaintainPolicy::default(),
)
.await
.unwrap();

assert!(report.was_healthy);
assert!(!report.rotated_signed);
assert_eq!(report.uploaded_one_time, 0);
assert_eq!(
puts(&server).await,
(0, 0),
"a healthy agent publishes nothing"
);
}

#[tokio::test]
async fn rotates_only_the_signed_key_when_the_relay_advertises_one_we_cannot_back() {
// The relay serves `spk_missing`, which the store does not hold — an inbound
// PREKEY_BUNDLE naming it would fail to decrypt, so maintenance must repair
// it. The one-time pool is healthy and must be left alone.
let signer = test_signer();
let server = relay(20, json!("spk_missing")).await;
let client = client_for(&server);
let store = store_backing(&signer, "spk_other").await;

let report = maintain_keys(
&client.keys,
&store,
signer.as_ref(),
&signer.agent_id(),
"identity-key",
&MaintainPolicy::default(),
)
.await
.unwrap();

assert!(!report.was_healthy);
assert!(report.rotated_signed);
assert_eq!(
report.uploaded_one_time, 0,
"a healthy pool is not disturbed"
);
assert_eq!(puts(&server).await, (1, 0));
}

#[tokio::test]
async fn tops_up_only_the_pool_when_low_leaving_a_good_signed_key_alone() {
// Peers have drawn the pool DOWN but not to zero: the relay still reports a
// positive count and only raises `lowOneTimePreKeys`. This pins the low-flag
// branch specifically (the empty-pool `count <= 0` branch is covered by the
// first-boot case), while the advertised signed pre-key is still backed — so
// refill must happen without any signed-key churn.
let signer = test_signer();
let server = relay(4, json!("spk_known")).await;
let client = client_for(&server);
let store = store_backing(&signer, "spk_known").await;

let report = maintain_keys(
&client.keys,
&store,
signer.as_ref(),
&signer.agent_id(),
"identity-key",
&MaintainPolicy::default(),
)
.await
.unwrap();

assert!(!report.was_healthy);
assert!(!report.rotated_signed, "signed key must not be rotated");
assert_eq!(report.uploaded_one_time, 20);
assert_eq!(puts(&server).await, (0, 1));
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

#[tokio::test]
async fn publishes_both_halves_on_first_boot() {
// Nothing published yet: the relay advertises no signed pre-key and an empty
// pool, and the store is empty.
let signer = test_signer();
let server = relay(0, json!(null)).await;
let client = client_for(&server);
let store = MemorySessionStore::new(generate_x25519_keypair());

let report = maintain_keys(
&client.keys,
&store,
signer.as_ref(),
&signer.agent_id(),
"identity-key",
&MaintainPolicy::default(),
)
.await
.unwrap();

assert!(!report.was_healthy);
assert!(report.rotated_signed);
assert_eq!(report.uploaded_one_time, 20);
assert_eq!(puts(&server).await, (1, 1));
}
Loading