//! Group sync: list the caller's groups, fetch their sealed GCK grant, and //! push/pull a group's shared changelog under its Group Content Key. //! //! The GCK is supplied by the caller, the `SyncStore` resolves it from the grant //! via the identity key (a later slice); these methods only encrypt/decrypt with //! it. Group entries bind `(group_id, table, row_id)` as AEAD associated data, so //! a ciphertext cannot be relocated across groups. Design: wiki //! synckit-multiscope-design. use bytes::Bytes; use tracing::instrument; use uuid::Uuid; use crate::{ crypto, error::Result, identity::{IdentityKeypair, IdentityPublicKey, generate_group_key, seal_gck_to_member}, ids::{DeviceId, GroupId, UserId}, types::{ ChangeEntry, GroupGrant, GroupMember, PullRequest, PullResponse, PulledChange, PushResponse, SyncGroup, WirePushRequest, }, }; use super::SyncKitClient; use super::helpers::{Idempotency, check_response}; /// Request body for `POST /groups`. The id is client-generated so the admin's /// grant can be sealed bound to it before the group exists. #[derive(serde::Serialize)] struct CreateGroupBody<'a> { id: GroupId, name: &'a str, admin_sealed_gck: String, admin_pubkey: String, } /// Request body for `POST /groups/{id}/members`. #[derive(serde::Serialize)] struct AddMemberBody<'a> { member_email: &'a str, sealed_gck: String, member_pubkey: &'a str, #[serde(skip_serializing_if = "Option::is_none")] role: Option<&'a str>, } impl SyncKitClient { /// List the groups the authenticated user belongs to within this app. #[instrument(skip(self))] pub async fn list_groups(&self) -> Result> { let token = self.require_token()?; self.retry_request_json(Idempotency::ReadOnly, || { let req = self.http.get(self.endpoints.groups()).bearer_auth(&token); async move { check_response(req.send().await?).await } }) .await } /// This user's identity public key (base64), derived from their master key. /// A member shares it with a group admin so the admin can seal the GCK to it /// (the paste-a-public-key model). Requires the master key to be loaded. pub fn my_identity_public_key(&self) -> Result { let master = self.require_master_key()?; Ok(IdentityKeypair::from_master_key(&master) .public_key() .to_base64()) } /// Create a group with the caller as its admin and first member. Mints a /// fresh Group Content Key, seals it to the caller's own identity public key /// (derived from the master key), and registers the group. /// /// The group id is generated client-side: the admin's grant is sealed with the /// group id bound as associated data, so the id must be known before the seal, /// before the server round-trip. Generation starts at 1. #[instrument(skip(self))] pub async fn create_group(&self, name: &str) -> Result { let token = self.require_token()?; let master = self.require_master_key()?; let identity = IdentityKeypair::from_master_key(&master); let group_id = GroupId::new(Uuid::new_v4()); let gck = generate_group_key(); let sealed = seal_gck_to_member(&gck, &identity.public_key(), &group_id.to_string(), 1)?; let body = Bytes::from(serde_json::to_vec(&CreateGroupBody { id: group_id, name, admin_sealed_gck: sealed, admin_pubkey: identity.public_key().to_base64(), })?); let url = self.endpoints.groups().to_string(); // The client-chosen id makes this safe to retry: a replay hits the PK and // the group already exists, but the sealed grant is identical. self.retry_request_json(Idempotency::Keyed, || { let req = self .http .post(&url) .bearer_auth(&token) .header("content-type", "application/json") .body(body.clone()); async move { check_response(req.send().await?).await } }) .await } /// Add a member to a group by account email, sealing the group's current GCK /// to the member's identity public key (obtained out of band, the member /// pastes their public key from [`my_identity_public_key`]). Admin only. #[instrument(skip(self, member_pubkey_b64))] pub async fn add_member( &self, group_id: GroupId, member_email: &str, member_pubkey_b64: &str, ) -> Result<()> { let token = self.require_token()?; // Resolve the group's current GCK (and its generation) from our own grant. let grant = self.group_grant(group_id).await?; let master = self.require_master_key()?; let gck = Self::open_group_grant(&grant, &master, group_id)?; let member_pubkey = IdentityPublicKey::from_base64(member_pubkey_b64)?; let sealed = seal_gck_to_member( &gck, &member_pubkey, &group_id.to_string(), grant.gck_version, )?; let body = Bytes::from(serde_json::to_vec(&AddMemberBody { member_email, sealed_gck: sealed, member_pubkey: member_pubkey_b64, role: None, })?); let url = self.endpoints.group_members(group_id); self.retry_request(Idempotency::Keyed, || { let req = self .http .post(&url) .bearer_auth(&token) .header("content-type", "application/json") .body(body.clone()); async move { check_response(req.send().await?).await } }) .await?; Ok(()) } /// List a group's members (id, role, joined-at). Admin only (the server gates /// it); use to show who can be removed. #[instrument(skip(self))] pub async fn list_members(&self, group_id: GroupId) -> Result> { let token = self.require_token()?; let url = self.endpoints.group_members(group_id); self.retry_request_json(Idempotency::ReadOnly, || { let req = self.http.get(&url).bearer_auth(&token); async move { check_response(req.send().await?).await } }) .await } /// Remove a member from a group by user id. Admin only. The server membership /// ACL revokes the member's group read/write access immediately. /// /// Forward secrecy for writes made *after* removal requires rotating the GCK /// (re-mint, re-seal to the remaining members, bump the generation), which the /// server does not yet expose, a later slice. This method performs the /// membership revocation only; data the member already pulled is in their hands. #[instrument(skip(self))] pub async fn remove_member(&self, group_id: GroupId, member: UserId) -> Result<()> { let token = self.require_token()?; let url = self.endpoints.group_member(group_id, member); self.retry_request(Idempotency::Keyed, || { let req = self.http.delete(&url).bearer_auth(&token); async move { check_response(req.send().await?).await } }) .await?; Ok(()) } /// Fetch the caller's own sealed GCK grant for a group. Open it with the /// member's identity private key to recover the GCK. #[instrument(skip(self))] pub async fn group_grant(&self, group_id: GroupId) -> Result { let token = self.require_token()?; let url = self.endpoints.group_grant(group_id); self.retry_request_json(Idempotency::ReadOnly, || { let req = self.http.get(&url).bearer_auth(&token); async move { check_response(req.send().await?).await } }) .await } /// Push encrypted changes to a group's shared changelog under its GCK. /// Returns the server cursor after the push. #[instrument(skip(self, gck, changes))] pub async fn group_push( &self, group_id: GroupId, gck: &[u8; 32], device_id: DeviceId, changes: Vec, ) -> Result { let token = self.require_token()?; let group_str = group_id.to_string(); let wire_changes = changes .into_iter() .map(|c| Self::encrypt_group_change_with_key(&group_str, c, gck)) .collect::>>()?; let body = Bytes::from(serde_json::to_vec(&WirePushRequest { device_id, batch_id: Uuid::new_v4(), changes: wire_changes, })?); let url = self.endpoints.group_push(group_id); let push_resp: PushResponse = self .retry_request_json(Idempotency::Keyed, || { let req = self .http .post(&url) .bearer_auth(&token) .header("content-type", "application/json") .body(body.clone()); async move { check_response(req.send().await?).await } }) .await?; Ok(push_resp.cursor) } /// Pull a group's changes since `cursor`, decrypting under its GCK. Returns /// `(changes, new_cursor, has_more)` with per-row device/seq metadata, ready /// for conflict resolution. Group pull has no master-key-rotation window (a /// group's key rotates server-side, re-encrypted in place), so the decrypt is /// a single-key pass. #[instrument(skip(self, gck))] pub async fn group_pull_rich( &self, group_id: GroupId, gck: &[u8; 32], device_id: DeviceId, cursor: i64, ) -> Result<(Vec, i64, bool)> { let token = self.require_token()?; let body = Bytes::from(serde_json::to_vec(&PullRequest { device_id, cursor })?); let url = self.endpoints.group_pull(group_id); let pull_resp: PullResponse = self .retry_request_json(Idempotency::ReadOnly, || { let req = self .http .post(&url) .bearer_auth(&token) .header("content-type", "application/json") .body(body.clone()); async move { check_response(req.send().await?).await } }) .await?; let group_str = group_id.to_string(); let changes = pull_resp .changes .into_iter() .map(|c| Self::decrypt_group_change_to_pulled(&group_str, c, gck)) .collect::>>()?; Ok((changes, pull_resp.cursor, pull_resp.has_more)) } /// Resolve a group's decrypted Group Content Key for `gck_version`, fetching /// and opening the sealed grant on a cache miss. /// /// Fast path: a cached GCK at the requested version is returned without a /// network call. On a miss (or a version mismatch, the group's key rotated), /// the sealed grant is fetched, opened with the identity keypair derived from /// this client's master key, and cached under the grant's actual version. The /// master key and identity secret never leave the client, the `SyncStore` /// only ever receives the resolved GCK. // Consumed by `sync_now`'s scope iteration in slice 4. #[allow(dead_code)] pub(crate) async fn group_content_key( &self, group_id: GroupId, gck_version: i32, ) -> Result { if let Some((v, gck)) = self.gck_cache.read().get(&group_id) && *v == gck_version { return Ok(crypto::ZeroizeOnDrop(**gck)); } let grant = self.group_grant(group_id).await?; let master = self.require_master_key()?; let gck = Self::open_group_grant(&grant, &master, group_id)?; let out = crypto::ZeroizeOnDrop(*gck); self.gck_cache .write() .insert(group_id, (grant.gck_version, gck)); Ok(out) } /// Open a sealed grant with the identity keypair derived from `master_key`. /// Pure (no I/O), so the derive-and-open path is unit-testable. fn open_group_grant( grant: &GroupGrant, master_key: &[u8; 32], group_id: GroupId, ) -> Result { let identity = crate::identity::IdentityKeypair::from_master_key(master_key); let gck = crate::identity::open_gck_grant( &grant.sealed_gck, &identity, &group_id.to_string(), grant.gck_version, )?; Ok(crypto::ZeroizeOnDrop(gck)) } /// Drop a group's cached GCK, forcing the next /// [`group_content_key`](Self::group_content_key) to re-fetch its grant. Call /// on a decrypt failure (a stale key after a rotation this client missed). // Consumed by `sync_now`'s scope iteration in slice 4. #[allow(dead_code)] pub(crate) fn invalidate_gck(&self, group_id: GroupId) { self.gck_cache.write().remove(&group_id); } } #[cfg(test)] mod tests { use super::*; use crate::error::SyncKitError; use crate::ids::DeviceId; use crate::types::{ChangeOp, Hlc, PullChangeEntry, WireChangeEntry}; use chrono::Utc; fn insert(table: &str, row: &str, data: serde_json::Value) -> ChangeEntry { ChangeEntry { table: table.into(), op: ChangeOp::Insert, row_id: row.into(), timestamp: Utc::now(), hlc: Hlc::zero(DeviceId::nil()), data: Some(data), extra: serde_json::Map::default(), } } fn to_pull(wire: WireChangeEntry, device: DeviceId, seq: i64) -> PullChangeEntry { PullChangeEntry { seq, device_id: device, table: wire.table, op: wire.op, row_id: wire.row_id, timestamp: wire.timestamp, data: wire.data, key_id: None, } } #[test] fn group_change_encrypt_decrypt_roundtrip_preserves_data_and_hlc() { let gck = crypto::generate_master_key(); let device = DeviceId::new(uuid::Uuid::new_v4()); let hlc = Hlc { wall_ms: 7, counter: 1, node: device, }; let mut e = insert("tasks", "r1", serde_json::json!({ "title": "shared" })); e.hlc = hlc; let wire = SyncKitClient::encrypt_group_change_with_key("grp-1", e, &gck).unwrap(); let pulled = SyncKitClient::decrypt_group_change_to_pulled("grp-1", to_pull(wire, device, 3), &gck) .unwrap(); assert_eq!(pulled.seq, 3); assert_eq!(pulled.device_id, device); assert_eq!( pulled.entry.data.unwrap(), serde_json::json!({ "title": "shared" }) ); assert_eq!(pulled.entry.hlc, hlc); } #[test] fn group_ciphertext_cannot_be_opened_under_another_group() { let gck = crypto::generate_master_key(); let wire = SyncKitClient::encrypt_group_change_with_key( "grp-A", insert("t", "r", serde_json::json!(1)), &gck, ) .unwrap(); // Same GCK, but a different group id in the AAD: the open fails closed. let err = SyncKitClient::decrypt_group_change_to_pulled( "grp-B", to_pull(wire, DeviceId::nil(), 1), &gck, ) .unwrap_err(); assert!(matches!(err, SyncKitError::DecryptionFailed)); } #[test] fn wrong_gck_fails_closed() { let gck = crypto::generate_master_key(); let other = crypto::generate_master_key(); let wire = SyncKitClient::encrypt_group_change_with_key( "grp-1", insert("t", "r", serde_json::json!(1)), &gck, ) .unwrap(); let err = SyncKitClient::decrypt_group_change_to_pulled( "grp-1", to_pull(wire, DeviceId::nil(), 1), &other, ) .unwrap_err(); assert!(matches!(err, SyncKitError::DecryptionFailed)); } #[test] fn create_group_admin_self_grant_is_openable() { // Mirrors create_group's sealing path: a client-generated group id, the // GCK sealed to the admin's own derived identity at generation 1. The // admin's device must be able to open its own grant back to the same GCK, // or it would be locked out of the group it just created. let master = crypto::generate_master_key(); let identity = crate::identity::IdentityKeypair::from_master_key(&master); let group = GroupId::new(uuid::Uuid::new_v4()); let gck = crate::identity::generate_group_key(); let sealed = crate::identity::seal_gck_to_member( &gck, &identity.public_key(), &group.to_string(), 1, ) .unwrap(); let grant = GroupGrant { sealed_gck: sealed, gck_version: 1, }; let opened = SyncKitClient::open_group_grant(&grant, &master, group).unwrap(); assert_eq!(*opened, gck); } #[test] fn open_group_grant_recovers_gck_via_derived_identity() { // The admin seals the GCK to the member's derived public key; the member's // client re-derives the identity from its master key and opens the grant. let group = GroupId::new(uuid::Uuid::new_v4()); let master_key = crypto::generate_master_key(); let member = crate::identity::IdentityKeypair::from_master_key(&master_key); let gck = crate::identity::generate_group_key(); let sealed = crate::identity::seal_gck_to_member(&gck, &member.public_key(), &group.to_string(), 4) .unwrap(); let grant = GroupGrant { sealed_gck: sealed, gck_version: 4, }; let opened = SyncKitClient::open_group_grant(&grant, &master_key, group).unwrap(); assert_eq!(*opened, gck); // A wrong-version grant (AAD mismatch) fails closed. let bad = GroupGrant { sealed_gck: grant.sealed_gck.clone(), gck_version: 5, }; assert!(SyncKitClient::open_group_grant(&bad, &master_key, group).is_err()); } }