Skip to main content

max / synckit

18.0 KB · 482 lines History Blame Raw
1 //! Group sync: list the caller's groups, fetch their sealed GCK grant, and
2 //! push/pull a group's shared changelog under its Group Content Key.
3 //!
4 //! The GCK is supplied by the caller, the `SyncStore` resolves it from the grant
5 //! via the identity key (a later slice); these methods only encrypt/decrypt with
6 //! it. Group entries bind `(group_id, table, row_id)` as AEAD associated data, so
7 //! a ciphertext cannot be relocated across groups. Design: wiki
8 //! synckit-multiscope-design.
9
10 use bytes::Bytes;
11 use tracing::instrument;
12 use uuid::Uuid;
13
14 use crate::{
15 crypto,
16 error::Result,
17 identity::{IdentityKeypair, IdentityPublicKey, generate_group_key, seal_gck_to_member},
18 ids::{DeviceId, GroupId, UserId},
19 types::{
20 ChangeEntry, GroupGrant, GroupMember, PullRequest, PullResponse, PulledChange,
21 PushResponse, SyncGroup, WirePushRequest,
22 },
23 };
24
25 use super::SyncKitClient;
26 use super::helpers::{Idempotency, check_response};
27
28 /// Request body for `POST /groups`. The id is client-generated so the admin's
29 /// grant can be sealed bound to it before the group exists.
30 #[derive(serde::Serialize)]
31 struct CreateGroupBody<'a> {
32 id: GroupId,
33 name: &'a str,
34 admin_sealed_gck: String,
35 admin_pubkey: String,
36 }
37
38 /// Request body for `POST /groups/{id}/members`.
39 #[derive(serde::Serialize)]
40 struct AddMemberBody<'a> {
41 member_email: &'a str,
42 sealed_gck: String,
43 member_pubkey: &'a str,
44 #[serde(skip_serializing_if = "Option::is_none")]
45 role: Option<&'a str>,
46 }
47
48 impl SyncKitClient {
49 /// List the groups the authenticated user belongs to within this app.
50 #[instrument(skip(self))]
51 pub async fn list_groups(&self) -> Result<Vec<SyncGroup>> {
52 let token = self.require_token()?;
53 self.retry_request_json(Idempotency::ReadOnly, || {
54 let req = self.http.get(self.endpoints.groups()).bearer_auth(&token);
55 async move { check_response(req.send().await?).await }
56 })
57 .await
58 }
59
60 /// This user's identity public key (base64), derived from their master key.
61 /// A member shares it with a group admin so the admin can seal the GCK to it
62 /// (the paste-a-public-key model). Requires the master key to be loaded.
63 pub fn my_identity_public_key(&self) -> Result<String> {
64 let master = self.require_master_key()?;
65 Ok(IdentityKeypair::from_master_key(&master)
66 .public_key()
67 .to_base64())
68 }
69
70 /// Create a group with the caller as its admin and first member. Mints a
71 /// fresh Group Content Key, seals it to the caller's own identity public key
72 /// (derived from the master key), and registers the group.
73 ///
74 /// The group id is generated client-side: the admin's grant is sealed with the
75 /// group id bound as associated data, so the id must be known before the seal,
76 /// before the server round-trip. Generation starts at 1.
77 #[instrument(skip(self))]
78 pub async fn create_group(&self, name: &str) -> Result<SyncGroup> {
79 let token = self.require_token()?;
80 let master = self.require_master_key()?;
81 let identity = IdentityKeypair::from_master_key(&master);
82 let group_id = GroupId::new(Uuid::new_v4());
83 let gck = generate_group_key();
84 let sealed = seal_gck_to_member(&gck, &identity.public_key(), &group_id.to_string(), 1)?;
85
86 let body = Bytes::from(serde_json::to_vec(&CreateGroupBody {
87 id: group_id,
88 name,
89 admin_sealed_gck: sealed,
90 admin_pubkey: identity.public_key().to_base64(),
91 })?);
92 let url = self.endpoints.groups().to_string();
93
94 // The client-chosen id makes this safe to retry: a replay hits the PK and
95 // the group already exists, but the sealed grant is identical.
96 self.retry_request_json(Idempotency::Keyed, || {
97 let req = self
98 .http
99 .post(&url)
100 .bearer_auth(&token)
101 .header("content-type", "application/json")
102 .body(body.clone());
103 async move { check_response(req.send().await?).await }
104 })
105 .await
106 }
107
108 /// Add a member to a group by account email, sealing the group's current GCK
109 /// to the member's identity public key (obtained out of band, the member
110 /// pastes their public key from [`my_identity_public_key`]). Admin only.
111 #[instrument(skip(self, member_pubkey_b64))]
112 pub async fn add_member(
113 &self,
114 group_id: GroupId,
115 member_email: &str,
116 member_pubkey_b64: &str,
117 ) -> Result<()> {
118 let token = self.require_token()?;
119 // Resolve the group's current GCK (and its generation) from our own grant.
120 let grant = self.group_grant(group_id).await?;
121 let master = self.require_master_key()?;
122 let gck = Self::open_group_grant(&grant, &master, group_id)?;
123
124 let member_pubkey = IdentityPublicKey::from_base64(member_pubkey_b64)?;
125 let sealed = seal_gck_to_member(
126 &gck,
127 &member_pubkey,
128 &group_id.to_string(),
129 grant.gck_version,
130 )?;
131
132 let body = Bytes::from(serde_json::to_vec(&AddMemberBody {
133 member_email,
134 sealed_gck: sealed,
135 member_pubkey: member_pubkey_b64,
136 role: None,
137 })?);
138 let url = self.endpoints.group_members(group_id);
139
140 self.retry_request(Idempotency::Keyed, || {
141 let req = self
142 .http
143 .post(&url)
144 .bearer_auth(&token)
145 .header("content-type", "application/json")
146 .body(body.clone());
147 async move { check_response(req.send().await?).await }
148 })
149 .await?;
150 Ok(())
151 }
152
153 /// List a group's members (id, role, joined-at). Admin only (the server gates
154 /// it); use to show who can be removed.
155 #[instrument(skip(self))]
156 pub async fn list_members(&self, group_id: GroupId) -> Result<Vec<GroupMember>> {
157 let token = self.require_token()?;
158 let url = self.endpoints.group_members(group_id);
159 self.retry_request_json(Idempotency::ReadOnly, || {
160 let req = self.http.get(&url).bearer_auth(&token);
161 async move { check_response(req.send().await?).await }
162 })
163 .await
164 }
165
166 /// Remove a member from a group by user id. Admin only. The server membership
167 /// ACL revokes the member's group read/write access immediately.
168 ///
169 /// Forward secrecy for writes made *after* removal requires rotating the GCK
170 /// (re-mint, re-seal to the remaining members, bump the generation), which the
171 /// server does not yet expose, a later slice. This method performs the
172 /// membership revocation only; data the member already pulled is in their hands.
173 #[instrument(skip(self))]
174 pub async fn remove_member(&self, group_id: GroupId, member: UserId) -> Result<()> {
175 let token = self.require_token()?;
176 let url = self.endpoints.group_member(group_id, member);
177 self.retry_request(Idempotency::Keyed, || {
178 let req = self.http.delete(&url).bearer_auth(&token);
179 async move { check_response(req.send().await?).await }
180 })
181 .await?;
182 Ok(())
183 }
184
185 /// Fetch the caller's own sealed GCK grant for a group. Open it with the
186 /// member's identity private key to recover the GCK.
187 #[instrument(skip(self))]
188 pub async fn group_grant(&self, group_id: GroupId) -> Result<GroupGrant> {
189 let token = self.require_token()?;
190 let url = self.endpoints.group_grant(group_id);
191 self.retry_request_json(Idempotency::ReadOnly, || {
192 let req = self.http.get(&url).bearer_auth(&token);
193 async move { check_response(req.send().await?).await }
194 })
195 .await
196 }
197
198 /// Push encrypted changes to a group's shared changelog under its GCK.
199 /// Returns the server cursor after the push.
200 #[instrument(skip(self, gck, changes))]
201 pub async fn group_push(
202 &self,
203 group_id: GroupId,
204 gck: &[u8; 32],
205 device_id: DeviceId,
206 changes: Vec<ChangeEntry>,
207 ) -> Result<i64> {
208 let token = self.require_token()?;
209 let group_str = group_id.to_string();
210 let wire_changes = changes
211 .into_iter()
212 .map(|c| Self::encrypt_group_change_with_key(&group_str, c, gck))
213 .collect::<Result<Vec<_>>>()?;
214
215 let body = Bytes::from(serde_json::to_vec(&WirePushRequest {
216 device_id,
217 batch_id: Uuid::new_v4(),
218 changes: wire_changes,
219 })?);
220 let url = self.endpoints.group_push(group_id);
221
222 let push_resp: PushResponse = self
223 .retry_request_json(Idempotency::Keyed, || {
224 let req = self
225 .http
226 .post(&url)
227 .bearer_auth(&token)
228 .header("content-type", "application/json")
229 .body(body.clone());
230 async move { check_response(req.send().await?).await }
231 })
232 .await?;
233 Ok(push_resp.cursor)
234 }
235
236 /// Pull a group's changes since `cursor`, decrypting under its GCK. Returns
237 /// `(changes, new_cursor, has_more)` with per-row device/seq metadata, ready
238 /// for conflict resolution. Group pull has no master-key-rotation window (a
239 /// group's key rotates server-side, re-encrypted in place), so the decrypt is
240 /// a single-key pass.
241 #[instrument(skip(self, gck))]
242 pub async fn group_pull_rich(
243 &self,
244 group_id: GroupId,
245 gck: &[u8; 32],
246 device_id: DeviceId,
247 cursor: i64,
248 ) -> Result<(Vec<PulledChange>, i64, bool)> {
249 let token = self.require_token()?;
250 let body = Bytes::from(serde_json::to_vec(&PullRequest { device_id, cursor })?);
251 let url = self.endpoints.group_pull(group_id);
252
253 let pull_resp: PullResponse = self
254 .retry_request_json(Idempotency::ReadOnly, || {
255 let req = self
256 .http
257 .post(&url)
258 .bearer_auth(&token)
259 .header("content-type", "application/json")
260 .body(body.clone());
261 async move { check_response(req.send().await?).await }
262 })
263 .await?;
264
265 let group_str = group_id.to_string();
266 let changes = pull_resp
267 .changes
268 .into_iter()
269 .map(|c| Self::decrypt_group_change_to_pulled(&group_str, c, gck))
270 .collect::<Result<Vec<_>>>()?;
271 Ok((changes, pull_resp.cursor, pull_resp.has_more))
272 }
273
274 /// Resolve a group's decrypted Group Content Key for `gck_version`, fetching
275 /// and opening the sealed grant on a cache miss.
276 ///
277 /// Fast path: a cached GCK at the requested version is returned without a
278 /// network call. On a miss (or a version mismatch, the group's key rotated),
279 /// the sealed grant is fetched, opened with the identity keypair derived from
280 /// this client's master key, and cached under the grant's actual version. The
281 /// master key and identity secret never leave the client, the `SyncStore`
282 /// only ever receives the resolved GCK.
283 // Consumed by `sync_now`'s scope iteration in slice 4.
284 #[allow(dead_code)]
285 pub(crate) async fn group_content_key(
286 &self,
287 group_id: GroupId,
288 gck_version: i32,
289 ) -> Result<crypto::ZeroizeOnDrop> {
290 if let Some((v, gck)) = self.gck_cache.read().get(&group_id)
291 && *v == gck_version
292 {
293 return Ok(crypto::ZeroizeOnDrop(**gck));
294 }
295
296 let grant = self.group_grant(group_id).await?;
297 let master = self.require_master_key()?;
298 let gck = Self::open_group_grant(&grant, &master, group_id)?;
299 let out = crypto::ZeroizeOnDrop(*gck);
300 self.gck_cache
301 .write()
302 .insert(group_id, (grant.gck_version, gck));
303 Ok(out)
304 }
305
306 /// Open a sealed grant with the identity keypair derived from `master_key`.
307 /// Pure (no I/O), so the derive-and-open path is unit-testable.
308 fn open_group_grant(
309 grant: &GroupGrant,
310 master_key: &[u8; 32],
311 group_id: GroupId,
312 ) -> Result<crypto::ZeroizeOnDrop> {
313 let identity = crate::identity::IdentityKeypair::from_master_key(master_key);
314 let gck = crate::identity::open_gck_grant(
315 &grant.sealed_gck,
316 &identity,
317 &group_id.to_string(),
318 grant.gck_version,
319 )?;
320 Ok(crypto::ZeroizeOnDrop(gck))
321 }
322
323 /// Drop a group's cached GCK, forcing the next
324 /// [`group_content_key`](Self::group_content_key) to re-fetch its grant. Call
325 /// on a decrypt failure (a stale key after a rotation this client missed).
326 // Consumed by `sync_now`'s scope iteration in slice 4.
327 #[allow(dead_code)]
328 pub(crate) fn invalidate_gck(&self, group_id: GroupId) {
329 self.gck_cache.write().remove(&group_id);
330 }
331 }
332
333 #[cfg(test)]
334 mod tests {
335 use super::*;
336 use crate::error::SyncKitError;
337 use crate::ids::DeviceId;
338 use crate::types::{ChangeOp, Hlc, PullChangeEntry, WireChangeEntry};
339 use chrono::Utc;
340
341 fn insert(table: &str, row: &str, data: serde_json::Value) -> ChangeEntry {
342 ChangeEntry {
343 table: table.into(),
344 op: ChangeOp::Insert,
345 row_id: row.into(),
346 timestamp: Utc::now(),
347 hlc: Hlc::zero(DeviceId::nil()),
348 data: Some(data),
349 extra: serde_json::Map::default(),
350 }
351 }
352
353 fn to_pull(wire: WireChangeEntry, device: DeviceId, seq: i64) -> PullChangeEntry {
354 PullChangeEntry {
355 seq,
356 device_id: device,
357 table: wire.table,
358 op: wire.op,
359 row_id: wire.row_id,
360 timestamp: wire.timestamp,
361 data: wire.data,
362 key_id: None,
363 }
364 }
365
366 #[test]
367 fn group_change_encrypt_decrypt_roundtrip_preserves_data_and_hlc() {
368 let gck = crypto::generate_master_key();
369 let device = DeviceId::new(uuid::Uuid::new_v4());
370 let hlc = Hlc {
371 wall_ms: 7,
372 counter: 1,
373 node: device,
374 };
375 let mut e = insert("tasks", "r1", serde_json::json!({ "title": "shared" }));
376 e.hlc = hlc;
377
378 let wire = SyncKitClient::encrypt_group_change_with_key("grp-1", e, &gck).unwrap();
379 let pulled =
380 SyncKitClient::decrypt_group_change_to_pulled("grp-1", to_pull(wire, device, 3), &gck)
381 .unwrap();
382
383 assert_eq!(pulled.seq, 3);
384 assert_eq!(pulled.device_id, device);
385 assert_eq!(
386 pulled.entry.data.unwrap(),
387 serde_json::json!({ "title": "shared" })
388 );
389 assert_eq!(pulled.entry.hlc, hlc);
390 }
391
392 #[test]
393 fn group_ciphertext_cannot_be_opened_under_another_group() {
394 let gck = crypto::generate_master_key();
395 let wire = SyncKitClient::encrypt_group_change_with_key(
396 "grp-A",
397 insert("t", "r", serde_json::json!(1)),
398 &gck,
399 )
400 .unwrap();
401 // Same GCK, but a different group id in the AAD: the open fails closed.
402 let err = SyncKitClient::decrypt_group_change_to_pulled(
403 "grp-B",
404 to_pull(wire, DeviceId::nil(), 1),
405 &gck,
406 )
407 .unwrap_err();
408 assert!(matches!(err, SyncKitError::DecryptionFailed));
409 }
410
411 #[test]
412 fn wrong_gck_fails_closed() {
413 let gck = crypto::generate_master_key();
414 let other = crypto::generate_master_key();
415 let wire = SyncKitClient::encrypt_group_change_with_key(
416 "grp-1",
417 insert("t", "r", serde_json::json!(1)),
418 &gck,
419 )
420 .unwrap();
421 let err = SyncKitClient::decrypt_group_change_to_pulled(
422 "grp-1",
423 to_pull(wire, DeviceId::nil(), 1),
424 &other,
425 )
426 .unwrap_err();
427 assert!(matches!(err, SyncKitError::DecryptionFailed));
428 }
429
430 #[test]
431 fn create_group_admin_self_grant_is_openable() {
432 // Mirrors create_group's sealing path: a client-generated group id, the
433 // GCK sealed to the admin's own derived identity at generation 1. The
434 // admin's device must be able to open its own grant back to the same GCK,
435 // or it would be locked out of the group it just created.
436 let master = crypto::generate_master_key();
437 let identity = crate::identity::IdentityKeypair::from_master_key(&master);
438 let group = GroupId::new(uuid::Uuid::new_v4());
439 let gck = crate::identity::generate_group_key();
440 let sealed = crate::identity::seal_gck_to_member(
441 &gck,
442 &identity.public_key(),
443 &group.to_string(),
444 1,
445 )
446 .unwrap();
447 let grant = GroupGrant {
448 sealed_gck: sealed,
449 gck_version: 1,
450 };
451 let opened = SyncKitClient::open_group_grant(&grant, &master, group).unwrap();
452 assert_eq!(*opened, gck);
453 }
454
455 #[test]
456 fn open_group_grant_recovers_gck_via_derived_identity() {
457 // The admin seals the GCK to the member's derived public key; the member's
458 // client re-derives the identity from its master key and opens the grant.
459 let group = GroupId::new(uuid::Uuid::new_v4());
460 let master_key = crypto::generate_master_key();
461 let member = crate::identity::IdentityKeypair::from_master_key(&master_key);
462 let gck = crate::identity::generate_group_key();
463 let sealed =
464 crate::identity::seal_gck_to_member(&gck, &member.public_key(), &group.to_string(), 4)
465 .unwrap();
466 let grant = GroupGrant {
467 sealed_gck: sealed,
468 gck_version: 4,
469 };
470
471 let opened = SyncKitClient::open_group_grant(&grant, &master_key, group).unwrap();
472 assert_eq!(*opened, gck);
473
474 // A wrong-version grant (AAD mismatch) fails closed.
475 let bad = GroupGrant {
476 sealed_gck: grant.sealed_gck.clone(),
477 gck_version: 5,
478 };
479 assert!(SyncKitClient::open_group_grant(&bad, &master_key, group).is_err());
480 }
481 }
482