A lexicon-driven AppView for ATProto.
0

Configure Feed

Select the types of activity you want to include in your feed.

feat: refactor permissioned spaces to better align with Dan's implementation

Signed-off-by: Trezy <tre@trezy.com>

Signed-off-by: Trezy <tre@trezy.com>

+986 -254
+4
migrations/postgres/20260511000000_add_space_did.sql
··· 1 + ALTER TABLE spaces ADD COLUMN did TEXT; 2 + UPDATE spaces SET did = owner_did; 3 + ALTER TABLE spaces ALTER COLUMN did SET NOT NULL; 4 + CREATE UNIQUE INDEX idx_spaces_did_type_skey ON spaces(did, type_nsid, skey);
+1
migrations/postgres/20260511100000_add_space_revision.sql
··· 1 + ALTER TABLE spaces ADD COLUMN revision TEXT;
+3
migrations/sqlite/20260511000000_add_space_did.sql
··· 1 + ALTER TABLE spaces ADD COLUMN did TEXT; 2 + UPDATE spaces SET did = owner_did; 3 + CREATE UNIQUE INDEX idx_spaces_did_type_skey ON spaces(did, type_nsid, skey);
+1
migrations/sqlite/20260511100000_add_space_revision.sql
··· 1 + ALTER TABLE spaces ADD COLUMN revision TEXT;
+6 -4
src/lua/atproto_api.rs
··· 280 280 let space = crate::spaces::db::get_space_by_address( 281 281 &state.db, 282 282 state.db_backend, 283 - &uri.owner_did, 283 + &uri.did, 284 284 &uri.type_nsid, 285 285 &uri.skey, 286 286 ) ··· 321 321 let space = crate::spaces::db::get_space_by_address( 322 322 &state.db, 323 323 state.db_backend, 324 - &uri.owner_did, 324 + &uri.did, 325 325 &uri.type_nsid, 326 326 &uri.skey, 327 327 ) ··· 361 361 let space = crate::spaces::db::get_space_by_address( 362 362 &state.db, 363 363 state.db_backend, 364 - &uri.owner_did, 364 + &uri.did, 365 365 &uri.type_nsid, 366 366 &uri.skey, 367 367 ) ··· 416 416 let space = crate::spaces::db::get_space_by_address( 417 417 &state.db, 418 418 state.db_backend, 419 - &uri.owner_did, 419 + &uri.did, 420 420 &uri.type_nsid, 421 421 &uri.skey, 422 422 ) ··· 433 433 &state.db, 434 434 state.db_backend, 435 435 &space.id, 436 + None, 436 437 collection.as_deref(), 437 438 limit.min(100), 438 439 cursor.as_deref(), 440 + false, 439 441 ) 440 442 .await 441 443 .map_err(|e| mlua::Error::runtime(format!("record query failed: {e}")))?;
+8 -8
src/lua/context.rs
··· 5 5 /// Optional space context passed to Lua scripts when the request is space-scoped. 6 6 #[derive(Debug, Clone)] 7 7 pub struct SpaceContext { 8 - pub space_uri: String, 8 + pub space: String, 9 9 pub space_id: String, 10 + pub did: String, 10 11 pub owner_did: String, 11 12 pub type_nsid: String, 12 13 pub skey: String, ··· 17 18 match space { 18 19 Some(ctx) => { 19 20 let table = lua.create_table()?; 20 - table.set("space_uri", ctx.space_uri.as_str())?; 21 + table.set("space", ctx.space.as_str())?; 21 22 table.set("space_id", ctx.space_id.as_str())?; 23 + table.set("did", ctx.did.as_str())?; 22 24 table.set("owner_did", ctx.owner_did.as_str())?; 23 25 table.set("type_nsid", ctx.type_nsid.as_str())?; 24 26 table.set("skey", ctx.skey.as_str())?; ··· 269 271 let lua = create_sandbox().unwrap(); 270 272 let params = HashMap::new(); 271 273 let space = SpaceContext { 272 - space_uri: "ats://did:plc:owner/com.example.forum/main".into(), 274 + space: "ats://did:plc:owner/com.example.forum/main".into(), 273 275 space_id: "space-123".into(), 276 + did: "did:plc:owner".into(), 274 277 owner_did: "did:plc:owner".into(), 275 278 type_nsid: "com.example.forum".into(), 276 279 skey: "main".into(), ··· 288 291 let globals = lua.globals(); 289 292 let space_table: mlua::Table = globals.get("space").unwrap(); 290 293 assert_eq!( 291 - space_table.get::<String>("space_uri").unwrap(), 294 + space_table.get::<String>("space").unwrap(), 292 295 "ats://did:plc:owner/com.example.forum/main" 293 296 ); 294 297 assert_eq!(space_table.get::<String>("space_id").unwrap(), "space-123"); 295 - assert_eq!( 296 - space_table.get::<String>("owner_did").unwrap(), 297 - "did:plc:owner" 298 - ); 298 + assert_eq!(space_table.get::<String>("did").unwrap(), "did:plc:owner"); 299 299 } 300 300 301 301 #[test]
+1 -1
src/lua/mod.rs
··· 5 5 mod http_api; 6 6 mod record; 7 7 pub(crate) mod sandbox; 8 - mod tid; 8 + pub(crate) mod tid; 9 9 mod xrpc_api; 10 10 11 11 #[allow(unused_imports)]
+6 -32
src/spaces/auth.rs
··· 9 9 use crate::error::AppError; 10 10 use crate::plugin::encryption::{decrypt, encrypt}; 11 11 use crate::spaces::credential::{ 12 - DEFAULT_CREDENTIAL_TTL_SECS, SpaceCredentialClaims, sign_credential, verify_credential, 12 + DEFAULT_CREDENTIAL_TTL_SECS, SpaceCredentialClaims, sign_credential, 13 13 }; 14 14 use crate::spaces::types::{AccessMode, Space}; 15 15 ··· 37 37 let exp = now + DEFAULT_CREDENTIAL_TTL_SECS; 38 38 39 39 let claims = SpaceCredentialClaims { 40 - iss: space.owner_did.clone(), 40 + iss: space.did.clone(), 41 41 sub: subject_did.to_string(), 42 - space: format!("{}/{}/{}", space.owner_did, space.type_nsid, space.skey), 42 + space: format!("ats://{}/{}/{}", space.did, space.type_nsid, space.skey), 43 43 scope: "read".into(), 44 44 iat: now, 45 45 exp, ··· 55 55 .unwrap_or_default(); 56 56 57 57 Ok(IssuedCredential { token, expires_at }) 58 - } 59 - 60 - pub async fn refresh_credential( 61 - pool: &sqlx::AnyPool, 62 - backend: DatabaseBackend, 63 - encryption_key: &[u8; 32], 64 - space: &Space, 65 - current_token: &str, 66 - ) -> Result<IssuedCredential, AppError> { 67 - let public_jwk = get_public_key(pool, backend, encryption_key, space).await?; 68 - let claims = verify_credential(current_token, &public_jwk)?; 69 - 70 - issue_credential(pool, backend, encryption_key, space, &claims.sub, None).await 71 58 } 72 59 73 60 pub fn check_app_access(space: &Space, client_id: Option<&str>) -> Result<(), AppError> { ··· 148 135 149 136 sqlx::query(&insert_sql) 150 137 .bind(Uuid::new_v4().to_string()) 151 - .bind(&space.owner_did) 138 + .bind(&space.did) 152 139 .bind(&space.id) 153 140 .bind(&encrypted_signing) 154 141 .bind(&encrypted_rotation) ··· 159 146 .map_err(|e| AppError::Internal(format!("failed to store space signing key: {e}")))?; 160 147 161 148 Ok(keypair.private_jwk) 162 - } 163 - 164 - async fn get_public_key( 165 - pool: &sqlx::AnyPool, 166 - backend: DatabaseBackend, 167 - encryption_key: &[u8; 32], 168 - space: &Space, 169 - ) -> Result<serde_json::Value, AppError> { 170 - let private_jwk = get_or_create_signing_key(pool, backend, encryption_key, space).await?; 171 - Ok(serde_json::json!({ 172 - "kty": "EC", 173 - "crv": "P-256", 174 - "x": private_jwk["x"], 175 - "y": private_jwk["y"], 176 - })) 177 149 } 178 150 179 151 struct SpaceKeypair { ··· 252 224 fn test_space(access_mode: AccessMode) -> Space { 253 225 Space { 254 226 id: "test-space".into(), 227 + did: "did:plc:owner".into(), 255 228 owner_did: "did:plc:owner".into(), 256 229 type_nsid: "com.example.forum".into(), 257 230 skey: "main".into(), ··· 262 235 app_denylist: None, 263 236 managing_app_did: None, 264 237 config: SpaceConfig::default(), 238 + revision: None, 265 239 created_at: String::new(), 266 240 updated_at: String::new(), 267 241 }
+106
src/spaces/credential.rs
··· 7 7 use crate::profile; 8 8 9 9 pub const DEFAULT_CREDENTIAL_TTL_SECS: u64 = 4 * 60 * 60; // 4 hours 10 + pub const GRANT_TTL_SECS: u64 = 5 * 60; // 5 minutes 11 + 12 + #[derive(Debug, Clone, Serialize, Deserialize)] 13 + pub struct MemberGrantClaims { 14 + pub sub: String, 15 + pub space: String, 16 + pub scope: String, 17 + pub iat: u64, 18 + pub exp: u64, 19 + } 20 + 21 + pub fn sign_grant(claims: &MemberGrantClaims, secret: &[u8; 32]) -> Result<String, AppError> { 22 + let header = jsonwebtoken::Header::new(jsonwebtoken::Algorithm::HS256); 23 + let key = jsonwebtoken::EncodingKey::from_secret(secret); 24 + jsonwebtoken::encode(&header, claims, &key) 25 + .map_err(|e| AppError::Internal(format!("failed to sign member grant: {e}"))) 26 + } 27 + 28 + pub fn verify_grant(token: &str, secret: &[u8; 32]) -> Result<MemberGrantClaims, AppError> { 29 + let key = jsonwebtoken::DecodingKey::from_secret(secret); 30 + let mut validation = jsonwebtoken::Validation::new(jsonwebtoken::Algorithm::HS256); 31 + validation.required_spec_claims.clear(); 32 + validation.validate_exp = false; 33 + let data = jsonwebtoken::decode::<MemberGrantClaims>(token, &key, &validation) 34 + .map_err(|e| AppError::Auth(format!("invalid member grant: {e}")))?; 35 + 36 + let now = std::time::SystemTime::now() 37 + .duration_since(std::time::UNIX_EPOCH) 38 + .unwrap() 39 + .as_secs(); 40 + 41 + if now > data.claims.exp { 42 + return Err(AppError::Auth("member grant has expired".into())); 43 + } 44 + 45 + Ok(data.claims) 46 + } 10 47 11 48 #[derive(Debug, Clone, Serialize, Deserialize)] 12 49 pub struct SpaceCredentialClaims { ··· 280 317 let keypair = generate_dpop_keypair().unwrap(); 281 318 let result = verify_credential("not-a-jwt", &keypair.public_jwk); 282 319 assert!(result.is_err()); 320 + } 321 + 322 + fn test_secret() -> [u8; 32] { 323 + [0xAB; 32] 324 + } 325 + 326 + #[test] 327 + fn grant_sign_and_verify_roundtrip() { 328 + let secret = test_secret(); 329 + let now = std::time::SystemTime::now() 330 + .duration_since(std::time::UNIX_EPOCH) 331 + .unwrap() 332 + .as_secs(); 333 + let claims = MemberGrantClaims { 334 + sub: "did:plc:member".into(), 335 + space: "ats://did:plc:space/com.example.forum/main".into(), 336 + scope: "read".into(), 337 + iat: now, 338 + exp: now + GRANT_TTL_SECS, 339 + }; 340 + 341 + let token = sign_grant(&claims, &secret).unwrap(); 342 + let verified = verify_grant(&token, &secret).unwrap(); 343 + 344 + assert_eq!(verified.sub, claims.sub); 345 + assert_eq!(verified.space, claims.space); 346 + assert_eq!(verified.scope, claims.scope); 347 + } 348 + 349 + #[test] 350 + fn grant_rejects_wrong_secret() { 351 + let secret1 = [0xAB; 32]; 352 + let secret2 = [0xCD; 32]; 353 + let now = std::time::SystemTime::now() 354 + .duration_since(std::time::UNIX_EPOCH) 355 + .unwrap() 356 + .as_secs(); 357 + let claims = MemberGrantClaims { 358 + sub: "did:plc:member".into(), 359 + space: "ats://did:plc:space/com.example.forum/main".into(), 360 + scope: "read".into(), 361 + iat: now, 362 + exp: now + GRANT_TTL_SECS, 363 + }; 364 + 365 + let token = sign_grant(&claims, &secret1).unwrap(); 366 + let result = verify_grant(&token, &secret2); 367 + assert!(result.is_err()); 368 + } 369 + 370 + #[test] 371 + fn grant_rejects_expired() { 372 + let secret = test_secret(); 373 + let now = std::time::SystemTime::now() 374 + .duration_since(std::time::UNIX_EPOCH) 375 + .unwrap() 376 + .as_secs(); 377 + let claims = MemberGrantClaims { 378 + sub: "did:plc:member".into(), 379 + space: "ats://did:plc:space/com.example.forum/main".into(), 380 + scope: "read".into(), 381 + iat: now - 600, 382 + exp: now - 300, 383 + }; 384 + 385 + let token = sign_grant(&claims, &secret).unwrap(); 386 + let result = verify_grant(&token, &secret); 387 + assert!(result.is_err()); 388 + assert!(result.unwrap_err().to_string().contains("expired")); 283 389 } 284 390 }
+242 -67
src/spaces/db.rs
··· 24 24 .map(|v| serde_json::to_string(v).unwrap_or_default()); 25 25 26 26 let sql = adapt_sql( 27 - "INSERT INTO spaces (id, owner_did, type_nsid, skey, display_name, description, access_mode, app_allowlist, app_denylist, managing_app_did, config, created_at, updated_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", 27 + "INSERT INTO spaces (id, did, owner_did, type_nsid, skey, display_name, description, access_mode, app_allowlist, app_denylist, managing_app_did, config, created_at, updated_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", 28 28 backend, 29 29 ); 30 30 31 31 sqlx::query(&sql) 32 32 .bind(&space.id) 33 + .bind(&space.did) 33 34 .bind(&space.owner_did) 34 35 .bind(&space.type_nsid) 35 36 .bind(&space.skey) ··· 55 56 id: &str, 56 57 ) -> Result<Option<Space>, AppError> { 57 58 let sql = adapt_sql( 58 - "SELECT id, owner_did, type_nsid, skey, display_name, description, access_mode, app_allowlist, app_denylist, managing_app_did, config, created_at, updated_at FROM spaces WHERE id = ?", 59 + "SELECT id, did, owner_did, type_nsid, skey, display_name, description, access_mode, app_allowlist, app_denylist, managing_app_did, config, revision, created_at, updated_at FROM spaces WHERE id = ?", 59 60 backend, 60 61 ); 61 62 ··· 71 72 pub async fn get_space_by_address( 72 73 pool: &sqlx::AnyPool, 73 74 backend: DatabaseBackend, 74 - owner_did: &str, 75 + did: &str, 75 76 type_nsid: &str, 76 77 skey: &str, 77 78 ) -> Result<Option<Space>, AppError> { 78 79 let sql = adapt_sql( 79 - "SELECT id, owner_did, type_nsid, skey, display_name, description, access_mode, app_allowlist, app_denylist, managing_app_did, config, created_at, updated_at FROM spaces WHERE owner_did = ? AND type_nsid = ? AND skey = ?", 80 + "SELECT id, did, owner_did, type_nsid, skey, display_name, description, access_mode, app_allowlist, app_denylist, managing_app_did, config, revision, created_at, updated_at FROM spaces WHERE did = ? AND type_nsid = ? AND skey = ?", 80 81 backend, 81 82 ); 82 83 83 84 let row: Option<SpaceRow> = sqlx::query_as(&sql) 84 - .bind(owner_did) 85 + .bind(did) 85 86 .bind(type_nsid) 86 87 .bind(skey) 87 88 .fetch_optional(pool) ··· 94 95 pub async fn list_spaces_by_owner( 95 96 pool: &sqlx::AnyPool, 96 97 backend: DatabaseBackend, 97 - owner_did: &str, 98 + did: &str, 98 99 ) -> Result<Vec<Space>, AppError> { 99 100 let sql = adapt_sql( 100 - "SELECT id, owner_did, type_nsid, skey, display_name, description, access_mode, app_allowlist, app_denylist, managing_app_did, config, created_at, updated_at FROM spaces WHERE owner_did = ? ORDER BY created_at DESC", 101 + "SELECT id, did, owner_did, type_nsid, skey, display_name, description, access_mode, app_allowlist, app_denylist, managing_app_did, config, revision, created_at, updated_at FROM spaces WHERE owner_did = ? ORDER BY created_at DESC", 101 102 backend, 102 103 ); 103 104 104 105 let rows: Vec<SpaceRow> = sqlx::query_as(&sql) 105 - .bind(owner_did) 106 + .bind(did) 106 107 .fetch_all(pool) 107 108 .await 108 109 .map_err(|e| AppError::Internal(format!("failed to list spaces: {e}")))?; ··· 110 111 rows.into_iter().map(parse_space_row).collect() 111 112 } 112 113 114 + pub struct SpaceView { 115 + pub uri: String, 116 + pub is_owner: bool, 117 + } 118 + 119 + pub async fn list_spaces_for_user( 120 + pool: &sqlx::AnyPool, 121 + backend: DatabaseBackend, 122 + did: &str, 123 + limit: i64, 124 + cursor: Option<&str>, 125 + ) -> Result<Vec<SpaceView>, AppError> { 126 + let (sql, has_cursor) = if cursor.is_some() { 127 + ( 128 + adapt_sql( 129 + "SELECT s.did, s.owner_did, s.type_nsid, s.skey, sm.created_at FROM space_members sm JOIN spaces s ON s.id = sm.space_id WHERE sm.member_did = ? AND sm.created_at > ? ORDER BY sm.created_at ASC LIMIT ?", 130 + backend, 131 + ), 132 + true, 133 + ) 134 + } else { 135 + ( 136 + adapt_sql( 137 + "SELECT s.did, s.owner_did, s.type_nsid, s.skey, sm.created_at FROM space_members sm JOIN spaces s ON s.id = sm.space_id WHERE sm.member_did = ? ORDER BY sm.created_at ASC LIMIT ?", 138 + backend, 139 + ), 140 + false, 141 + ) 142 + }; 143 + 144 + let mut query = sqlx::query_as::<_, (String, String, String, String, String)>(&sql).bind(did); 145 + if has_cursor { 146 + query = query.bind(cursor.unwrap()); 147 + } 148 + query = query.bind(limit); 149 + 150 + let rows = query 151 + .fetch_all(pool) 152 + .await 153 + .map_err(|e| AppError::Internal(format!("failed to list spaces for user: {e}")))?; 154 + 155 + Ok(rows 156 + .into_iter() 157 + .map( 158 + |(space_did, owner_did, type_nsid, skey, _created_at)| SpaceView { 159 + uri: format!("ats://{}/{}/{}", space_did, type_nsid, skey), 160 + is_owner: owner_did == did, 161 + }, 162 + ) 163 + .collect()) 164 + } 165 + 113 166 pub async fn update_space( 114 167 pool: &sqlx::AnyPool, 115 168 backend: DatabaseBackend, ··· 170 223 String, 171 224 String, 172 225 String, 226 + String, 173 227 Option<String>, 174 228 Option<String>, 175 229 String, ··· 177 231 Option<String>, 178 232 Option<String>, 179 233 String, 234 + Option<String>, 180 235 String, 181 236 String, 182 237 ); 183 238 184 239 fn parse_space_row(r: SpaceRow) -> Result<Space, AppError> { 185 - let access_mode = AccessMode::parse(&r.6) 186 - .ok_or_else(|| AppError::Internal(format!("invalid access_mode: {}", r.6)))?; 240 + let access_mode = AccessMode::parse(&r.7) 241 + .ok_or_else(|| AppError::Internal(format!("invalid access_mode: {}", r.7)))?; 187 242 let app_allowlist: Option<Vec<String>> = 188 - r.7.as_deref() 243 + r.8.as_deref() 189 244 .map(serde_json::from_str) 190 245 .transpose() 191 246 .map_err(|e| AppError::Internal(format!("invalid app_allowlist: {e}")))?; 192 247 let app_denylist: Option<Vec<String>> = 193 - r.8.as_deref() 248 + r.9.as_deref() 194 249 .map(serde_json::from_str) 195 250 .transpose() 196 251 .map_err(|e| AppError::Internal(format!("invalid app_denylist: {e}")))?; 197 - let config: SpaceConfig = serde_json::from_str(&r.10) 252 + let config: SpaceConfig = serde_json::from_str(&r.11) 198 253 .map_err(|e| AppError::Internal(format!("invalid space config: {e}")))?; 199 254 200 255 Ok(Space { 201 256 id: r.0, 202 - owner_did: r.1, 203 - type_nsid: r.2, 204 - skey: r.3, 205 - display_name: r.4, 206 - description: r.5, 257 + did: r.1, 258 + owner_did: r.2, 259 + type_nsid: r.3, 260 + skey: r.4, 261 + display_name: r.5, 262 + description: r.6, 207 263 access_mode, 208 264 app_allowlist, 209 265 app_denylist, 210 - managing_app_did: r.9, 266 + managing_app_did: r.10, 211 267 config, 212 - created_at: r.11, 213 - updated_at: r.12, 268 + revision: r.12, 269 + created_at: r.13, 270 + updated_at: r.14, 214 271 }) 215 272 } 216 273 ··· 232 289 sqlx::query(&sql) 233 290 .bind(&member.id) 234 291 .bind(&member.space_id) 235 - .bind(&member.member_did) 292 + .bind(&member.did) 236 293 .bind(member.access.as_str()) 237 294 .bind(member.is_delegation as i32) 238 295 .bind(&member.granted_by) ··· 248 305 pool: &sqlx::AnyPool, 249 306 backend: DatabaseBackend, 250 307 space_id: &str, 251 - member_did: &str, 308 + did: &str, 252 309 ) -> Result<bool, AppError> { 253 310 let sql = adapt_sql( 254 311 "DELETE FROM space_members WHERE space_id = ? AND member_did = ?", ··· 257 314 258 315 let result = sqlx::query(&sql) 259 316 .bind(space_id) 260 - .bind(member_did) 317 + .bind(did) 261 318 .execute(pool) 262 319 .await 263 320 .map_err(|e| AppError::Internal(format!("failed to remove member: {e}")))?; ··· 269 326 pool: &sqlx::AnyPool, 270 327 backend: DatabaseBackend, 271 328 space_id: &str, 272 - member_did: &str, 329 + did: &str, 273 330 ) -> Result<Option<SpaceMember>, AppError> { 274 331 let sql = adapt_sql( 275 332 "SELECT id, space_id, member_did, access, is_delegation, granted_by, created_at FROM space_members WHERE space_id = ? AND member_did = ?", ··· 278 335 279 336 let row: Option<MemberRow> = sqlx::query_as(&sql) 280 337 .bind(space_id) 281 - .bind(member_did) 338 + .bind(did) 282 339 .fetch_optional(pool) 283 340 .await 284 341 .map_err(|e| AppError::Internal(format!("failed to get member: {e}")))?; ··· 308 365 pub async fn list_spaces_for_member( 309 366 pool: &sqlx::AnyPool, 310 367 backend: DatabaseBackend, 311 - member_did: &str, 368 + did: &str, 312 369 ) -> Result<Vec<SpaceMember>, AppError> { 313 370 let sql = adapt_sql( 314 371 "SELECT id, space_id, member_did, access, is_delegation, granted_by, created_at FROM space_members WHERE member_did = ? ORDER BY created_at ASC", ··· 316 373 ); 317 374 318 375 let rows: Vec<MemberRow> = sqlx::query_as(&sql) 319 - .bind(member_did) 376 + .bind(did) 320 377 .fetch_all(pool) 321 378 .await 322 379 .map_err(|e| AppError::Internal(format!("failed to list spaces for member: {e}")))?; ··· 333 390 Ok(SpaceMember { 334 391 id: r.0, 335 392 space_id: r.1, 336 - member_did: r.2, 393 + did: r.2, 337 394 access, 338 395 is_delegation: r.4 != 0, 339 396 granted_by: r.5, ··· 422 479 row.map(parse_record_row).transpose() 423 480 } 424 481 482 + #[allow(clippy::too_many_arguments)] 425 483 pub async fn list_space_records( 426 484 pool: &sqlx::AnyPool, 427 485 backend: DatabaseBackend, 428 486 space_id: &str, 487 + repo: Option<&str>, 429 488 collection: Option<&str>, 430 489 limit: i64, 431 490 cursor: Option<&str>, 491 + reverse: bool, 432 492 ) -> Result<Vec<SpaceRecord>, AppError> { 433 - let (sql, has_collection, has_cursor) = match (collection, cursor) { 434 - (Some(_), Some(_)) => ( 435 - adapt_sql( 436 - "SELECT uri, space_id, author_did, collection, rkey, record, cid, indexed_at FROM space_records WHERE space_id = ? AND collection = ? AND indexed_at > ? ORDER BY indexed_at ASC LIMIT ?", 437 - backend, 438 - ), 439 - true, 440 - true, 441 - ), 442 - (Some(_), None) => ( 443 - adapt_sql( 444 - "SELECT uri, space_id, author_did, collection, rkey, record, cid, indexed_at FROM space_records WHERE space_id = ? AND collection = ? ORDER BY indexed_at ASC LIMIT ?", 445 - backend, 446 - ), 447 - true, 448 - false, 449 - ), 450 - (None, Some(_)) => ( 451 - adapt_sql( 452 - "SELECT uri, space_id, author_did, collection, rkey, record, cid, indexed_at FROM space_records WHERE space_id = ? AND indexed_at > ? ORDER BY indexed_at ASC LIMIT ?", 453 - backend, 454 - ), 455 - false, 456 - true, 457 - ), 458 - (None, None) => ( 459 - adapt_sql( 460 - "SELECT uri, space_id, author_did, collection, rkey, record, cid, indexed_at FROM space_records WHERE space_id = ? ORDER BY indexed_at ASC LIMIT ?", 461 - backend, 462 - ), 463 - false, 464 - false, 465 - ), 493 + let mut conditions = vec!["space_id = ?".to_string()]; 494 + if repo.is_some() { 495 + conditions.push("author_did = ?".to_string()); 496 + } 497 + if collection.is_some() { 498 + conditions.push("collection = ?".to_string()); 499 + } 500 + let (cursor_op, order) = if reverse { 501 + ("indexed_at < ?", "DESC") 502 + } else { 503 + ("indexed_at > ?", "ASC") 466 504 }; 505 + if cursor.is_some() { 506 + conditions.push(cursor_op.to_string()); 507 + } 508 + 509 + let where_clause = conditions.join(" AND "); 510 + let raw = format!( 511 + "SELECT uri, space_id, author_did, collection, rkey, record, cid, indexed_at FROM space_records WHERE {} ORDER BY indexed_at {} LIMIT ?", 512 + where_clause, order 513 + ); 514 + let sql = adapt_sql(&raw, backend); 467 515 468 516 let mut query = sqlx::query_as::<_, RecordRow>(&sql).bind(space_id); 469 - 470 - if has_collection { 471 - query = query.bind(collection.unwrap()); 517 + if let Some(r) = repo { 518 + query = query.bind(r); 519 + } 520 + if let Some(c) = collection { 521 + query = query.bind(c); 472 522 } 473 - if has_cursor { 474 - query = query.bind(cursor.unwrap()); 523 + if let Some(cur) = cursor { 524 + query = query.bind(cur); 475 525 } 476 526 query = query.bind(limit); 477 527 ··· 483 533 rows.into_iter().map(parse_record_row).collect() 484 534 } 485 535 536 + pub async fn insert_space_record( 537 + pool: &sqlx::AnyPool, 538 + backend: DatabaseBackend, 539 + record: &SpaceRecord, 540 + ) -> Result<(), AppError> { 541 + let now = now_rfc3339(); 542 + let record_json = serde_json::to_string(&record.record) 543 + .map_err(|e| AppError::Internal(format!("failed to serialize record: {e}")))?; 544 + 545 + let sql = adapt_sql( 546 + "INSERT INTO space_records (uri, space_id, author_did, collection, rkey, record, cid, indexed_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?)", 547 + backend, 548 + ); 549 + 550 + sqlx::query(&sql) 551 + .bind(&record.uri) 552 + .bind(&record.space_id) 553 + .bind(&record.author_did) 554 + .bind(&record.collection) 555 + .bind(&record.rkey) 556 + .bind(&record_json) 557 + .bind(&record.cid) 558 + .bind(&now) 559 + .execute(pool) 560 + .await 561 + .map_err(|e| { 562 + let msg = e.to_string(); 563 + if msg.contains("UNIQUE") || msg.contains("duplicate") || msg.contains("unique") { 564 + AppError::Conflict("Record already exists".into()) 565 + } else { 566 + AppError::Internal(format!("failed to create space record: {e}")) 567 + } 568 + })?; 569 + 570 + Ok(()) 571 + } 572 + 573 + pub async fn upsert_space_record_with_swap( 574 + pool: &sqlx::AnyPool, 575 + backend: DatabaseBackend, 576 + record: &SpaceRecord, 577 + swap_cid: &str, 578 + ) -> Result<(), AppError> { 579 + let now = now_rfc3339(); 580 + let record_json = serde_json::to_string(&record.record) 581 + .map_err(|e| AppError::Internal(format!("failed to serialize record: {e}")))?; 582 + 583 + let sql = adapt_sql( 584 + "UPDATE space_records SET record = ?, cid = ?, indexed_at = ? WHERE uri = ? AND cid = ?", 585 + backend, 586 + ); 587 + 588 + let result = sqlx::query(&sql) 589 + .bind(&record_json) 590 + .bind(&record.cid) 591 + .bind(&now) 592 + .bind(&record.uri) 593 + .bind(swap_cid) 594 + .execute(pool) 595 + .await 596 + .map_err(|e| AppError::Internal(format!("failed to update space record: {e}")))?; 597 + 598 + if result.rows_affected() == 0 { 599 + let existing = get_space_record(pool, backend, &record.uri).await?; 600 + if existing.is_some() { 601 + return Err(AppError::Conflict("Record CID mismatch".into())); 602 + } 603 + return Err(AppError::NotFound("Record not found".into())); 604 + } 605 + 606 + Ok(()) 607 + } 608 + 486 609 pub async fn delete_space_record( 487 610 pool: &sqlx::AnyPool, 488 611 backend: DatabaseBackend, ··· 497 620 .map_err(|e| AppError::Internal(format!("failed to delete space record: {e}")))?; 498 621 499 622 Ok(result.rows_affected() > 0) 623 + } 624 + 625 + pub async fn delete_space_record_with_swap( 626 + pool: &sqlx::AnyPool, 627 + backend: DatabaseBackend, 628 + uri: &str, 629 + swap_cid: &str, 630 + ) -> Result<bool, AppError> { 631 + let sql = adapt_sql( 632 + "DELETE FROM space_records WHERE uri = ? AND cid = ?", 633 + backend, 634 + ); 635 + 636 + let result = sqlx::query(&sql) 637 + .bind(uri) 638 + .bind(swap_cid) 639 + .execute(pool) 640 + .await 641 + .map_err(|e| AppError::Internal(format!("failed to delete space record: {e}")))?; 642 + 643 + if result.rows_affected() == 0 { 644 + let existing = get_space_record(pool, backend, uri).await?; 645 + if existing.is_some() { 646 + return Err(AppError::Conflict("Record CID mismatch".into())); 647 + } 648 + return Err(AppError::NotFound("Record not found".into())); 649 + } 650 + 651 + Ok(true) 652 + } 653 + 654 + pub async fn update_space_revision( 655 + pool: &sqlx::AnyPool, 656 + backend: DatabaseBackend, 657 + space_id: &str, 658 + revision: &str, 659 + ) -> Result<(), AppError> { 660 + let now = now_rfc3339(); 661 + let sql = adapt_sql( 662 + "UPDATE spaces SET revision = ?, updated_at = ? WHERE id = ?", 663 + backend, 664 + ); 665 + 666 + sqlx::query(&sql) 667 + .bind(revision) 668 + .bind(&now) 669 + .bind(space_id) 670 + .execute(pool) 671 + .await 672 + .map_err(|e| AppError::Internal(format!("failed to update space revision: {e}")))?; 673 + 674 + Ok(()) 500 675 } 501 676 502 677 type RecordRow = (
+5 -6
src/spaces/members.rs
··· 77 77 .await?; 78 78 } 79 79 } else { 80 - merge_access(resolved, &member.member_did, member.access); 80 + merge_access(resolved, &member.did, member.access); 81 81 } 82 82 } 83 83 ··· 93 93 backend: DatabaseBackend, 94 94 member: &SpaceMember, 95 95 ) -> Result<Option<String>, AppError> { 96 - if member.member_did.starts_with("ats://") { 97 - let uri = SpaceUri::parse(&member.member_did)?; 96 + if member.did.starts_with("ats://") { 97 + let uri = SpaceUri::parse(&member.did)?; 98 98 let space = 99 - db::get_space_by_address(pool, backend, &uri.owner_did, &uri.type_nsid, &uri.skey) 100 - .await?; 99 + db::get_space_by_address(pool, backend, &uri.did, &uri.type_nsid, &uri.skey).await?; 101 100 Ok(space.map(|s| s.id)) 102 101 } else { 103 - let space = db::get_space(pool, backend, &member.member_did).await?; 102 + let space = db::get_space(pool, backend, &member.did).await?; 104 103 Ok(space.map(|s| s.id)) 105 104 } 106 105 }
+12 -16
src/spaces/mod.rs
··· 10 10 11 11 /// A parsed `ats://` URI for addressing permissioned data. 12 12 /// 13 - /// Full form: `ats://<space-owner-did>/<space-type-nsid>/<skey>/<user-did>/<collection-nsid>/<rkey>` 14 - /// Space-only form: `ats://<space-owner-did>/<space-type-nsid>/<skey>` 13 + /// Full form: `ats://<space-did>/<space-type-nsid>/<skey>/<user-did>/<collection-nsid>/<rkey>` 14 + /// Space-only form: `ats://<space-did>/<space-type-nsid>/<skey>` 15 15 #[derive(Debug, Clone, PartialEq, Eq, Hash)] 16 16 pub struct SpaceUri { 17 - pub owner_did: String, 17 + pub did: String, 18 18 pub type_nsid: String, 19 19 pub skey: String, 20 20 pub user_did: Option<String>, ··· 32 32 33 33 if parts.len() < 3 { 34 34 return Err(AppError::BadRequest( 35 - "SpaceUri requires at least owner_did/type_nsid/skey".into(), 35 + "SpaceUri requires at least did/type_nsid/skey".into(), 36 36 )); 37 37 } 38 38 ··· 42 42 )); 43 43 } 44 44 45 - let owner_did = parts[0].to_string(); 45 + let did = parts[0].to_string(); 46 46 let type_nsid = parts[1].to_string(); 47 47 let skey = parts[2].to_string(); 48 48 ··· 61 61 }; 62 62 63 63 Ok(SpaceUri { 64 - owner_did, 64 + did, 65 65 type_nsid, 66 66 skey, 67 67 user_did, ··· 71 71 } 72 72 73 73 pub fn space_uri(&self) -> String { 74 - format!("ats://{}/{}/{}", self.owner_did, self.type_nsid, self.skey) 74 + format!("ats://{}/{}/{}", self.did, self.type_nsid, self.skey) 75 75 } 76 76 77 77 pub fn is_record_uri(&self) -> bool { ··· 85 85 86 86 impl fmt::Display for SpaceUri { 87 87 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { 88 - write!( 89 - f, 90 - "ats://{}/{}/{}", 91 - self.owner_did, self.type_nsid, self.skey 92 - )?; 88 + write!(f, "ats://{}/{}/{}", self.did, self.type_nsid, self.skey)?; 93 89 if let (Some(user), Some(col), Some(rkey)) = (&self.user_did, &self.collection, &self.rkey) 94 90 { 95 91 write!(f, "/{}/{}/{}", user, col, rkey)?; ··· 105 101 #[test] 106 102 fn parse_space_uri() { 107 103 let uri = SpaceUri::parse("ats://did:plc:abc123/com.example.forum/main").unwrap(); 108 - assert_eq!(uri.owner_did, "did:plc:abc123"); 104 + assert_eq!(uri.did, "did:plc:abc123"); 109 105 assert_eq!(uri.type_nsid, "com.example.forum"); 110 106 assert_eq!(uri.skey, "main"); 111 107 assert!(uri.is_space_uri()); ··· 119 115 "ats://did:plc:abc123/com.example.forum/main/did:plc:user1/com.example.forum.post/3k2abc", 120 116 ) 121 117 .unwrap(); 122 - assert_eq!(uri.owner_did, "did:plc:abc123"); 118 + assert_eq!(uri.did, "did:plc:abc123"); 123 119 assert_eq!(uri.type_nsid, "com.example.forum"); 124 120 assert_eq!(uri.skey, "main"); 125 121 assert_eq!(uri.user_did.as_deref(), Some("did:plc:user1")); ··· 132 128 #[test] 133 129 fn display_space_uri() { 134 130 let uri = SpaceUri { 135 - owner_did: "did:plc:abc123".into(), 131 + did: "did:plc:abc123".into(), 136 132 type_nsid: "com.example.forum".into(), 137 133 skey: "main".into(), 138 134 user_did: None, ··· 148 144 #[test] 149 145 fn display_record_uri() { 150 146 let uri = SpaceUri { 151 - owner_did: "did:plc:abc123".into(), 147 + did: "did:plc:abc123".into(), 152 148 type_nsid: "com.example.forum".into(), 153 149 skey: "main".into(), 154 150 user_did: Some("did:plc:user1".into()),
+587 -119
src/spaces/routes.rs
··· 11 11 use crate::auth::XrpcClaims; 12 12 use crate::db::{adapt_sql, now_rfc3339}; 13 13 use crate::error::AppError; 14 + use crate::lua::tid::generate_tid; 14 15 use crate::spaces::types::*; 15 16 use crate::spaces::{SpaceUri, db, members}; 16 17 ··· 21 22 #[derive(Deserialize)] 22 23 #[serde(rename_all = "camelCase")] 23 24 struct CreateSpaceInput { 25 + #[serde(rename = "type")] 24 26 type_nsid: String, 25 27 skey: String, 26 28 display_name: Option<String>, ··· 33 35 #[derive(Deserialize)] 34 36 #[serde(rename_all = "camelCase")] 35 37 struct SpaceUriQuery { 36 - space_uri: String, 38 + space: String, 37 39 } 38 40 39 41 #[derive(Deserialize)] 40 42 #[serde(rename_all = "camelCase")] 41 43 struct ListSpacesQuery { 42 - owner_did: Option<String>, 44 + did: Option<String>, 45 + limit: Option<i64>, 46 + cursor: Option<String>, 43 47 } 44 48 45 49 #[derive(Deserialize)] 46 50 #[serde(rename_all = "camelCase")] 47 51 struct DeleteSpaceInput { 48 - space_uri: String, 52 + space: String, 49 53 } 50 54 51 55 #[derive(Deserialize)] 52 56 #[serde(rename_all = "camelCase")] 53 57 struct UpdateSpaceInput { 54 - space_uri: String, 58 + space: String, 55 59 display_name: Option<Option<String>>, 56 60 description: Option<Option<String>>, 57 61 access_mode: Option<AccessMode>, ··· 64 68 #[derive(Deserialize)] 65 69 #[serde(rename_all = "camelCase")] 66 70 struct PutRecordInput { 67 - space_uri: String, 71 + space: String, 68 72 collection: String, 69 73 rkey: String, 70 74 record: serde_json::Value, 75 + swap_record: Option<String>, 71 76 } 72 77 73 78 #[derive(Deserialize)] 74 79 #[serde(rename_all = "camelCase")] 75 80 struct DeleteRecordInput { 76 - space_uri: String, 81 + space: String, 77 82 collection: String, 78 83 rkey: String, 84 + swap_record: Option<String>, 79 85 } 80 86 81 87 #[derive(Deserialize)] 82 88 #[serde(rename_all = "camelCase")] 83 89 struct GetRecordQuery { 84 - space_uri: String, 90 + space: String, 85 91 collection: String, 86 92 rkey: String, 87 93 } ··· 89 95 #[derive(Deserialize)] 90 96 #[serde(rename_all = "camelCase")] 91 97 struct ListRecordsQuery { 92 - space_uri: String, 98 + space: String, 99 + repo: Option<String>, 93 100 collection: Option<String>, 94 101 limit: Option<i64>, 95 102 cursor: Option<String>, 103 + reverse: Option<bool>, 96 104 } 97 105 98 106 #[derive(Deserialize)] 99 107 #[serde(rename_all = "camelCase")] 100 108 struct AddMemberInput { 101 - space_uri: String, 102 - member_did: String, 109 + space: String, 110 + did: String, 103 111 access: Option<SpaceAccess>, 104 112 is_delegation: Option<bool>, 105 113 } ··· 107 115 #[derive(Deserialize)] 108 116 #[serde(rename_all = "camelCase")] 109 117 struct RemoveMemberInput { 110 - space_uri: String, 111 - member_did: String, 118 + space: String, 119 + did: String, 112 120 } 113 121 114 122 #[derive(Deserialize)] 115 123 #[serde(rename_all = "camelCase")] 116 124 struct CreateInviteInput { 117 - space_uri: String, 125 + space: String, 118 126 access: Option<SpaceAccess>, 119 127 max_uses: Option<i64>, 120 128 expires_at: Option<String>, ··· 129 137 #[derive(Deserialize)] 130 138 #[serde(rename_all = "camelCase")] 131 139 struct RevokeInviteInput { 132 - space_uri: String, 140 + space: String, 133 141 invite_id: String, 134 142 } 135 143 136 144 #[derive(Deserialize)] 137 145 #[serde(rename_all = "camelCase")] 138 - struct GetCredentialInput { 139 - space_uri: String, 146 + struct GetMemberGrantInput { 147 + space: String, 148 + } 149 + 150 + #[derive(Deserialize)] 151 + #[serde(rename_all = "camelCase")] 152 + struct GetSpaceCredentialInput { 153 + grant: String, 154 + } 155 + 156 + #[derive(Deserialize)] 157 + #[serde(rename_all = "camelCase")] 158 + struct CreateRecordInput { 159 + space: String, 160 + collection: String, 161 + record: serde_json::Value, 140 162 } 141 163 142 164 #[derive(Deserialize)] 143 165 #[serde(rename_all = "camelCase")] 144 - struct RefreshCredentialInput { 145 - space_uri: String, 146 - credential: String, 166 + struct ApplyWritesInput { 167 + space: String, 168 + swap_commit: Option<String>, 169 + writes: Vec<WriteOp>, 170 + } 171 + 172 + #[derive(Deserialize)] 173 + #[serde(tag = "action", rename_all = "camelCase")] 174 + enum WriteOp { 175 + Create { 176 + collection: String, 177 + rkey: Option<String>, 178 + value: serde_json::Value, 179 + }, 180 + Update { 181 + collection: String, 182 + rkey: String, 183 + value: serde_json::Value, 184 + #[serde(rename = "swapRecord")] 185 + swap_record: Option<String>, 186 + }, 187 + Delete { 188 + collection: String, 189 + rkey: String, 190 + #[serde(rename = "swapRecord")] 191 + swap_record: Option<String>, 192 + }, 147 193 } 148 194 149 195 // --------------------------------------------------------------------------- ··· 161 207 .route(&format!("/xrpc/{NS}.space.delete"), post(delete_space)) 162 208 .route(&format!("/xrpc/{NS}.space.update"), post(update_space)) 163 209 // Records 210 + .route( 211 + &format!("/xrpc/{NS}.space.createRecord"), 212 + post(create_record), 213 + ) 164 214 .route(&format!("/xrpc/{NS}.space.putRecord"), post(put_record)) 165 215 .route( 166 216 &format!("/xrpc/{NS}.space.deleteRecord"), 167 217 post(delete_record), 168 218 ) 219 + .route(&format!("/xrpc/{NS}.space.applyWrites"), post(apply_writes)) 169 220 .route(&format!("/xrpc/{NS}.space.getRecord"), get(get_record)) 170 221 .route(&format!("/xrpc/{NS}.space.listRecords"), get(list_records)) 171 222 // Members ··· 191 242 .route(&format!("/xrpc/{NS}.space.invite.list"), get(list_invites)) 192 243 // Credentials 193 244 .route( 194 - &format!("/xrpc/{NS}.space.getCredential"), 195 - post(get_credential), 245 + &format!("/xrpc/{NS}.space.getMemberGrant"), 246 + post(get_member_grant), 196 247 ) 197 248 .route( 198 - &format!("/xrpc/{NS}.space.refreshCredential"), 199 - post(refresh_credential), 249 + &format!("/xrpc/{NS}.space.getSpaceCredential"), 250 + post(get_space_credential), 200 251 ) 201 252 } 202 253 ··· 216 267 db::get_space_by_address( 217 268 &state.db, 218 269 state.db_backend, 219 - &uri.owner_did, 270 + &uri.did, 220 271 &uri.type_nsid, 221 272 &uri.skey, 222 273 ) ··· 257 308 space_credential: Option<&str>, 258 309 ) -> Result<SpaceAccess, AppError> { 259 310 if let Some(token) = space_credential { 260 - let space_uri = format!( 261 - "ats://{}/{}/{}", 262 - space.owner_did, space.type_nsid, space.skey 263 - ); 311 + let space_uri = format!("ats://{}/{}/{}", space.did, space.type_nsid, space.skey); 264 312 match crate::spaces::credential::verify_external_credential( 265 313 token, 266 314 &state.http, ··· 319 367 let did = claims.did().to_string(); 320 368 321 369 if input.type_nsid.is_empty() || input.skey.is_empty() { 322 - return Err(AppError::BadRequest( 323 - "type_nsid and skey are required".into(), 324 - )); 370 + return Err(AppError::BadRequest("type and skey are required".into())); 325 371 } 326 372 327 373 let existing = db::get_space_by_address( ··· 340 386 341 387 let space = Space { 342 388 id: Uuid::new_v4().to_string(), 389 + did: did.clone(), 343 390 owner_did: did.clone(), 344 391 type_nsid: input.type_nsid, 345 392 skey: input.skey, ··· 350 397 app_denylist: None, 351 398 managing_app_did: input.managing_app_did, 352 399 config: input.config.unwrap_or_default(), 400 + revision: None, 353 401 created_at: now_rfc3339(), 354 402 updated_at: now_rfc3339(), 355 403 }; ··· 360 408 let member = SpaceMember { 361 409 id: Uuid::new_v4().to_string(), 362 410 space_id: space.id.clone(), 363 - member_did: did.clone(), 411 + did: did.clone(), 364 412 access: SpaceAccess::Write, 365 413 is_delegation: false, 366 414 granted_by: Some(did), ··· 368 416 }; 369 417 db::add_member(&state.db, state.db_backend, &member).await?; 370 418 371 - let space_uri = format!( 372 - "ats://{}/{}/{}", 373 - space.owner_did, space.type_nsid, space.skey 374 - ); 419 + let space_uri = format!("ats://{}/{}/{}", space.did, space.type_nsid, space.skey); 375 420 let body = serde_json::json!({ 376 - "spaceUri": space_uri, 377 - "space": space, 421 + "uri": space_uri, 378 422 }); 379 423 380 424 let mut response = Json(body).into_response(); ··· 387 431 xrpc_claims: XrpcClaims, 388 432 Query(query): Query<SpaceUriQuery>, 389 433 ) -> Result<Json<serde_json::Value>, AppError> { 390 - let space = resolve_space(&state, &query.space_uri).await?; 434 + let space = resolve_space(&state, &query.space).await?; 391 435 392 436 // If the space's membership is not public, require auth + membership 393 437 if !space.config.membership_public { ··· 400 444 } 401 445 } 402 446 403 - let space_uri = format!( 404 - "ats://{}/{}/{}", 405 - space.owner_did, space.type_nsid, space.skey 406 - ); 447 + let space_uri = format!("ats://{}/{}/{}", space.did, space.type_nsid, space.skey); 407 448 Ok(Json(serde_json::json!({ 408 - "spaceUri": space_uri, 449 + "uri": space_uri, 409 450 "space": space, 410 451 }))) 411 452 } ··· 416 457 Query(query): Query<ListSpacesQuery>, 417 458 ) -> Result<Json<serde_json::Value>, AppError> { 418 459 let claims = require_auth(&xrpc_claims)?; 419 - let did = claims.did().to_string(); 460 + let did = query.did.unwrap_or_else(|| claims.did().to_string()); 461 + let limit = query.limit.unwrap_or(50).min(100); 420 462 421 - let owner = query.owner_did.as_deref().unwrap_or(&did); 422 - let spaces = db::list_spaces_by_owner(&state.db, state.db_backend, owner).await?; 463 + let views = db::list_spaces_for_user( 464 + &state.db, 465 + state.db_backend, 466 + &did, 467 + limit, 468 + query.cursor.as_deref(), 469 + ) 470 + .await?; 423 471 424 - let spaces_with_uris: Vec<serde_json::Value> = spaces 472 + let cursor = if views.len() as i64 == limit { 473 + views.last().map(|v| v.uri.clone()) 474 + } else { 475 + None 476 + }; 477 + 478 + let spaces_json: Vec<serde_json::Value> = views 425 479 .into_iter() 426 - .map(|s| { 427 - let uri = format!("ats://{}/{}/{}", s.owner_did, s.type_nsid, s.skey); 428 - serde_json::json!({ "spaceUri": uri, "space": s }) 480 + .map(|v| { 481 + serde_json::json!({ 482 + "uri": v.uri, 483 + "isOwner": v.is_owner, 484 + }) 429 485 }) 430 486 .collect(); 431 487 432 - Ok(Json(serde_json::json!({ "spaces": spaces_with_uris }))) 488 + Ok(Json(serde_json::json!({ 489 + "spaces": spaces_json, 490 + "cursor": cursor, 491 + }))) 433 492 } 434 493 435 494 async fn delete_space( ··· 438 497 Json(input): Json<DeleteSpaceInput>, 439 498 ) -> Result<Json<serde_json::Value>, AppError> { 440 499 let claims = require_auth(&xrpc_claims)?; 441 - let space = resolve_space(&state, &input.space_uri).await?; 500 + let space = resolve_space(&state, &input.space).await?; 442 501 require_space_admin(&state, &space, claims.did()).await?; 443 502 444 503 db::delete_space(&state.db, state.db_backend, &space.id).await?; ··· 452 511 Json(input): Json<UpdateSpaceInput>, 453 512 ) -> Result<Json<serde_json::Value>, AppError> { 454 513 let claims = require_auth(&xrpc_claims)?; 455 - let mut space = resolve_space(&state, &input.space_uri).await?; 514 + let mut space = resolve_space(&state, &input.space).await?; 456 515 require_space_admin(&state, &space, claims.did()).await?; 457 516 458 517 if let Some(name) = input.display_name { ··· 479 538 480 539 db::update_space(&state.db, state.db_backend, &space).await?; 481 540 482 - let space_uri = format!( 483 - "ats://{}/{}/{}", 484 - space.owner_did, space.type_nsid, space.skey 485 - ); 541 + let space_uri = format!("ats://{}/{}/{}", space.did, space.type_nsid, space.skey); 486 542 Ok(Json(serde_json::json!({ 487 - "spaceUri": space_uri, 543 + "uri": space_uri, 488 544 "space": space, 489 545 }))) 490 546 } ··· 493 549 // Record handlers 494 550 // --------------------------------------------------------------------------- 495 551 552 + async fn create_record( 553 + State(state): State<AppState>, 554 + xrpc_claims: XrpcClaims, 555 + headers: HeaderMap, 556 + Json(input): Json<CreateRecordInput>, 557 + ) -> Result<Response, AppError> { 558 + let claims = require_auth(&xrpc_claims)?; 559 + let did = claims.did().to_string(); 560 + let space = resolve_space(&state, &input.space).await?; 561 + let cred = extract_space_credential(&headers); 562 + require_membership(&state, &space, &did, true, cred.as_deref()).await?; 563 + 564 + let rkey = generate_tid(); 565 + let cid = content_cid(&input.record); 566 + let record_uri = format!( 567 + "ats://{}/{}/{}/{}/{}/{}", 568 + space.did, space.type_nsid, space.skey, did, input.collection, rkey 569 + ); 570 + 571 + let record = SpaceRecord { 572 + uri: record_uri.clone(), 573 + space_id: space.id.clone(), 574 + author_did: did, 575 + collection: input.collection, 576 + rkey, 577 + record: input.record, 578 + cid: cid.clone(), 579 + indexed_at: now_rfc3339(), 580 + }; 581 + 582 + db::insert_space_record(&state.db, state.db_backend, &record).await?; 583 + 584 + let rev = generate_tid(); 585 + db::update_space_revision(&state.db, state.db_backend, &space.id, &rev).await?; 586 + 587 + let body = serde_json::json!({ 588 + "uri": record_uri, 589 + "cid": cid, 590 + }); 591 + 592 + let mut response = Json(body).into_response(); 593 + *response.status_mut() = StatusCode::CREATED; 594 + Ok(response) 595 + } 596 + 496 597 async fn put_record( 497 598 State(state): State<AppState>, 498 599 xrpc_claims: XrpcClaims, ··· 501 602 ) -> Result<Response, AppError> { 502 603 let claims = require_auth(&xrpc_claims)?; 503 604 let did = claims.did().to_string(); 504 - let space = resolve_space(&state, &input.space_uri).await?; 605 + let space = resolve_space(&state, &input.space).await?; 505 606 let cred = extract_space_credential(&headers); 506 607 require_membership(&state, &space, &did, true, cred.as_deref()).await?; 507 608 508 609 let cid = content_cid(&input.record); 509 610 let record_uri = format!( 510 611 "ats://{}/{}/{}/{}/{}/{}", 511 - space.owner_did, space.type_nsid, space.skey, did, input.collection, input.rkey 612 + space.did, space.type_nsid, space.skey, did, input.collection, input.rkey 512 613 ); 513 614 514 615 let record = SpaceRecord { 515 616 uri: record_uri.clone(), 516 - space_id: space.id, 617 + space_id: space.id.clone(), 517 618 author_did: did, 518 619 collection: input.collection, 519 620 rkey: input.rkey, ··· 522 623 indexed_at: now_rfc3339(), 523 624 }; 524 625 525 - db::upsert_space_record(&state.db, state.db_backend, &record).await?; 626 + if let Some(swap_cid) = input.swap_record { 627 + db::upsert_space_record_with_swap(&state.db, state.db_backend, &record, &swap_cid).await?; 628 + } else { 629 + db::upsert_space_record(&state.db, state.db_backend, &record).await?; 630 + } 631 + 632 + let rev = generate_tid(); 633 + db::update_space_revision(&state.db, state.db_backend, &space.id, &rev).await?; 526 634 527 635 let body = serde_json::json!({ 528 636 "uri": record_uri, ··· 541 649 ) -> Result<Json<serde_json::Value>, AppError> { 542 650 let claims = require_auth(&xrpc_claims)?; 543 651 let did = claims.did().to_string(); 544 - let space = resolve_space(&state, &input.space_uri).await?; 652 + let space = resolve_space(&state, &input.space).await?; 545 653 546 654 let record_uri = format!( 547 655 "ats://{}/{}/{}/{}/{}/{}", 548 - space.owner_did, space.type_nsid, space.skey, did, input.collection, input.rkey 656 + space.did, space.type_nsid, space.skey, did, input.collection, input.rkey 549 657 ); 550 658 551 - let record = db::get_space_record(&state.db, state.db_backend, &record_uri).await?; 552 - match record { 553 - Some(r) if r.author_did != did => { 554 - return Err(AppError::Forbidden( 555 - "You can only delete your own records".into(), 556 - )); 557 - } 558 - None => { 559 - return Err(AppError::NotFound("Record not found".into())); 659 + if let Some(swap_cid) = input.swap_record { 660 + db::delete_space_record_with_swap(&state.db, state.db_backend, &record_uri, &swap_cid) 661 + .await?; 662 + } else { 663 + let record = db::get_space_record(&state.db, state.db_backend, &record_uri).await?; 664 + match record { 665 + Some(r) if r.author_did != did => { 666 + return Err(AppError::Forbidden( 667 + "You can only delete your own records".into(), 668 + )); 669 + } 670 + None => { 671 + return Err(AppError::NotFound("Record not found".into())); 672 + } 673 + _ => {} 560 674 } 561 - _ => {} 675 + db::delete_space_record(&state.db, state.db_backend, &record_uri).await?; 562 676 } 563 677 564 - db::delete_space_record(&state.db, state.db_backend, &record_uri).await?; 678 + let rev = generate_tid(); 679 + db::update_space_revision(&state.db, state.db_backend, &space.id, &rev).await?; 565 680 566 681 Ok(Json(serde_json::json!({ "success": true }))) 567 682 } 568 683 684 + async fn apply_writes( 685 + State(state): State<AppState>, 686 + xrpc_claims: XrpcClaims, 687 + headers: HeaderMap, 688 + Json(input): Json<ApplyWritesInput>, 689 + ) -> Result<Json<serde_json::Value>, AppError> { 690 + let claims = require_auth(&xrpc_claims)?; 691 + let did = claims.did().to_string(); 692 + let space = resolve_space(&state, &input.space).await?; 693 + let cred = extract_space_credential(&headers); 694 + require_membership(&state, &space, &did, true, cred.as_deref()).await?; 695 + 696 + if let Some(ref expected_rev) = input.swap_commit { 697 + match &space.revision { 698 + Some(current_rev) if current_rev != expected_rev => { 699 + return Err(AppError::Conflict("swapCommit mismatch".into())); 700 + } 701 + None if !expected_rev.is_empty() => { 702 + return Err(AppError::Conflict("swapCommit mismatch".into())); 703 + } 704 + _ => {} 705 + } 706 + } 707 + 708 + let mut results = Vec::with_capacity(input.writes.len()); 709 + 710 + for op in input.writes { 711 + match op { 712 + WriteOp::Create { 713 + collection, 714 + rkey, 715 + value, 716 + } => { 717 + let rkey = rkey.unwrap_or_else(generate_tid); 718 + let cid = content_cid(&value); 719 + let record_uri = format!( 720 + "ats://{}/{}/{}/{}/{}/{}", 721 + space.did, space.type_nsid, space.skey, did, collection, rkey 722 + ); 723 + let record = SpaceRecord { 724 + uri: record_uri.clone(), 725 + space_id: space.id.clone(), 726 + author_did: did.clone(), 727 + collection, 728 + rkey, 729 + record: value, 730 + cid: cid.clone(), 731 + indexed_at: now_rfc3339(), 732 + }; 733 + db::insert_space_record(&state.db, state.db_backend, &record).await?; 734 + results.push(serde_json::json!({ 735 + "uri": record_uri, 736 + "cid": cid, 737 + })); 738 + } 739 + WriteOp::Update { 740 + collection, 741 + rkey, 742 + value, 743 + swap_record, 744 + } => { 745 + let cid = content_cid(&value); 746 + let record_uri = format!( 747 + "ats://{}/{}/{}/{}/{}/{}", 748 + space.did, space.type_nsid, space.skey, did, collection, rkey 749 + ); 750 + let record = SpaceRecord { 751 + uri: record_uri.clone(), 752 + space_id: space.id.clone(), 753 + author_did: did.clone(), 754 + collection, 755 + rkey, 756 + record: value, 757 + cid: cid.clone(), 758 + indexed_at: now_rfc3339(), 759 + }; 760 + if let Some(swap_cid) = swap_record { 761 + db::upsert_space_record_with_swap( 762 + &state.db, 763 + state.db_backend, 764 + &record, 765 + &swap_cid, 766 + ) 767 + .await?; 768 + } else { 769 + db::upsert_space_record(&state.db, state.db_backend, &record).await?; 770 + } 771 + results.push(serde_json::json!({ 772 + "uri": record_uri, 773 + "cid": cid, 774 + })); 775 + } 776 + WriteOp::Delete { 777 + collection, 778 + rkey, 779 + swap_record, 780 + } => { 781 + let record_uri = format!( 782 + "ats://{}/{}/{}/{}/{}/{}", 783 + space.did, space.type_nsid, space.skey, did, collection, rkey 784 + ); 785 + if let Some(swap_cid) = swap_record { 786 + db::delete_space_record_with_swap( 787 + &state.db, 788 + state.db_backend, 789 + &record_uri, 790 + &swap_cid, 791 + ) 792 + .await?; 793 + } else { 794 + db::delete_space_record(&state.db, state.db_backend, &record_uri).await?; 795 + } 796 + results.push(serde_json::json!({})); 797 + } 798 + } 799 + } 800 + 801 + let rev = generate_tid(); 802 + db::update_space_revision(&state.db, state.db_backend, &space.id, &rev).await?; 803 + 804 + Ok(Json(serde_json::json!({ 805 + "results": results, 806 + }))) 807 + } 808 + 569 809 async fn get_record( 570 810 State(state): State<AppState>, 571 811 xrpc_claims: XrpcClaims, ··· 573 813 Query(query): Query<GetRecordQuery>, 574 814 ) -> Result<Json<serde_json::Value>, AppError> { 575 815 let claims = require_auth(&xrpc_claims)?; 576 - let space = resolve_space(&state, &query.space_uri).await?; 816 + let space = resolve_space(&state, &query.space).await?; 577 817 let cred = extract_space_credential(&headers); 578 818 require_membership(&state, &space, claims.did(), false, cred.as_deref()).await?; 579 819 ··· 589 829 590 830 Ok(Json(serde_json::json!({ 591 831 "uri": record.uri, 592 - "space": query.space_uri, 593 - "collection": record.collection, 594 - "record": record.record, 595 832 "cid": record.cid, 833 + "value": record.record, 596 834 }))) 597 835 } 598 836 ··· 603 841 Query(query): Query<ListRecordsQuery>, 604 842 ) -> Result<Json<serde_json::Value>, AppError> { 605 843 let claims = require_auth(&xrpc_claims)?; 606 - let space = resolve_space(&state, &query.space_uri).await?; 844 + let space = resolve_space(&state, &query.space).await?; 607 845 let cred = extract_space_credential(&headers); 608 846 require_membership(&state, &space, claims.did(), false, cred.as_deref()).await?; 609 847 848 + let repo = query.repo.as_deref().or_else(|| { 849 + if cred.is_some() { 850 + None 851 + } else { 852 + Some(claims.did()) 853 + } 854 + }); 855 + 610 856 let limit = query.limit.unwrap_or(50).min(100); 857 + let reverse = query.reverse.unwrap_or(false); 611 858 let records = db::list_space_records( 612 859 &state.db, 613 860 state.db_backend, 614 861 &space.id, 862 + repo, 615 863 query.collection.as_deref(), 616 864 limit, 617 865 query.cursor.as_deref(), 866 + reverse, 618 867 ) 619 868 .await?; 620 869 ··· 624 873 .into_iter() 625 874 .map(|r| { 626 875 serde_json::json!({ 627 - "uri": r.uri, 628 - "space": query.space_uri, 629 876 "collection": r.collection, 630 - "record": r.record, 877 + "rkey": r.rkey, 631 878 "cid": r.cid, 632 879 }) 633 880 }) ··· 649 896 headers: HeaderMap, 650 897 Query(query): Query<SpaceUriQuery>, 651 898 ) -> Result<Json<serde_json::Value>, AppError> { 652 - let space = resolve_space(&state, &query.space_uri).await?; 899 + let space = resolve_space(&state, &query.space).await?; 653 900 654 901 if !space.config.membership_public { 655 902 let claims = require_auth(&xrpc_claims)?; ··· 668 915 Json(input): Json<AddMemberInput>, 669 916 ) -> Result<Response, AppError> { 670 917 let claims = require_auth(&xrpc_claims)?; 671 - let space = resolve_space(&state, &input.space_uri).await?; 918 + let space = resolve_space(&state, &input.space).await?; 672 919 require_space_admin(&state, &space, claims.did()).await?; 673 920 674 - let existing = 675 - db::get_member(&state.db, state.db_backend, &space.id, &input.member_did).await?; 921 + let existing = db::get_member(&state.db, state.db_backend, &space.id, &input.did).await?; 676 922 if existing.is_some() { 677 923 return Err(AppError::Conflict( 678 924 "Member already exists in this space".into(), ··· 682 928 let member = SpaceMember { 683 929 id: Uuid::new_v4().to_string(), 684 930 space_id: space.id, 685 - member_did: input.member_did, 931 + did: input.did, 686 932 access: input.access.unwrap_or(SpaceAccess::Read), 687 933 is_delegation: input.is_delegation.unwrap_or(false), 688 934 granted_by: Some(claims.did().to_string()), ··· 702 948 Json(input): Json<RemoveMemberInput>, 703 949 ) -> Result<Json<serde_json::Value>, AppError> { 704 950 let claims = require_auth(&xrpc_claims)?; 705 - let space = resolve_space(&state, &input.space_uri).await?; 951 + let space = resolve_space(&state, &input.space).await?; 706 952 require_space_admin(&state, &space, claims.did()).await?; 707 953 708 - let removed = 709 - db::remove_member(&state.db, state.db_backend, &space.id, &input.member_did).await?; 954 + let removed = db::remove_member(&state.db, state.db_backend, &space.id, &input.did).await?; 710 955 711 956 if !removed { 712 957 return Err(AppError::NotFound("Member not found in this space".into())); ··· 725 970 Json(input): Json<CreateInviteInput>, 726 971 ) -> Result<Response, AppError> { 727 972 let claims = require_auth(&xrpc_claims)?; 728 - let space = resolve_space(&state, &input.space_uri).await?; 973 + let space = resolve_space(&state, &input.space).await?; 729 974 require_space_admin(&state, &space, claims.did()).await?; 730 975 731 976 let mut token_bytes = [0u8; 24]; ··· 802 1047 let member = SpaceMember { 803 1048 id: Uuid::new_v4().to_string(), 804 1049 space_id: invite.space_id.clone(), 805 - member_did: did, 1050 + did, 806 1051 access: invite.access, 807 1052 is_delegation: false, 808 1053 granted_by: Some(invite.created_by.clone()), ··· 813 1058 db::increment_invite_uses(&state.db, state.db_backend, &invite.id).await?; 814 1059 815 1060 let space = db::get_space(&state.db, state.db_backend, &invite.space_id).await?; 816 - let space_uri = space.map(|s| format!("ats://{}/{}/{}", s.owner_did, s.type_nsid, s.skey)); 1061 + let space_uri = space.map(|s| format!("ats://{}/{}/{}", s.did, s.type_nsid, s.skey)); 817 1062 818 1063 let mut response = Json(serde_json::json!({ 819 - "spaceUri": space_uri, 1064 + "uri": space_uri, 820 1065 "access": member.access, 821 1066 })) 822 1067 .into_response(); ··· 830 1075 Json(input): Json<RevokeInviteInput>, 831 1076 ) -> Result<Json<serde_json::Value>, AppError> { 832 1077 let claims = require_auth(&xrpc_claims)?; 833 - let space = resolve_space(&state, &input.space_uri).await?; 1078 + let space = resolve_space(&state, &input.space).await?; 834 1079 require_space_admin(&state, &space, claims.did()).await?; 835 1080 836 1081 let revoked = db::revoke_invite(&state.db, state.db_backend, &input.invite_id).await?; ··· 847 1092 Query(query): Query<SpaceUriQuery>, 848 1093 ) -> Result<Json<serde_json::Value>, AppError> { 849 1094 let claims = require_auth(&xrpc_claims)?; 850 - let space = resolve_space(&state, &query.space_uri).await?; 1095 + let space = resolve_space(&state, &query.space).await?; 851 1096 require_space_admin(&state, &space, claims.did()).await?; 852 1097 853 1098 let invites = db::list_invites(&state.db, state.db_backend, &space.id).await?; ··· 875 1120 // Credential handlers 876 1121 // --------------------------------------------------------------------------- 877 1122 878 - async fn get_credential( 1123 + async fn get_member_grant( 879 1124 State(state): State<AppState>, 880 1125 xrpc_claims: XrpcClaims, 881 - Json(input): Json<GetCredentialInput>, 1126 + Json(input): Json<GetMemberGrantInput>, 882 1127 ) -> Result<Json<serde_json::Value>, AppError> { 883 1128 let claims = require_auth(&xrpc_claims)?; 884 1129 let did = claims.did().to_string(); 885 - let space = resolve_space(&state, &input.space_uri).await?; 1130 + let space = resolve_space(&state, &input.space).await?; 886 1131 887 1132 require_membership(&state, &space, &did, false, None).await?; 888 1133 ··· 890 1135 AppError::Internal("TOKEN_ENCRYPTION_KEY is required for space credentials".into()) 891 1136 })?; 892 1137 893 - let client_id = claims.client_key().map(|k| k.to_string()); 894 - let issued = crate::spaces::auth::issue_credential( 895 - &state.db, 896 - state.db_backend, 897 - encryption_key, 898 - &space, 899 - &did, 900 - client_id.as_deref(), 901 - ) 902 - .await?; 1138 + let now = std::time::SystemTime::now() 1139 + .duration_since(std::time::UNIX_EPOCH) 1140 + .unwrap() 1141 + .as_secs(); 1142 + let exp = now + crate::spaces::credential::GRANT_TTL_SECS; 1143 + 1144 + let space_uri = format!("ats://{}/{}/{}", space.did, space.type_nsid, space.skey); 1145 + let grant_claims = crate::spaces::credential::MemberGrantClaims { 1146 + sub: did, 1147 + space: space_uri, 1148 + scope: "read".into(), 1149 + iat: now, 1150 + exp, 1151 + }; 1152 + 1153 + let grant = crate::spaces::credential::sign_grant(&grant_claims, encryption_key)?; 1154 + 1155 + let expires_at = chrono::DateTime::from_timestamp(exp as i64, 0) 1156 + .map(|dt| dt.to_rfc3339()) 1157 + .unwrap_or_default(); 903 1158 904 1159 Ok(Json(serde_json::json!({ 905 - "credential": issued.token, 906 - "expiresAt": issued.expires_at, 1160 + "grant": grant, 1161 + "expiresAt": expires_at, 907 1162 }))) 908 1163 } 909 1164 910 - async fn refresh_credential( 1165 + async fn get_space_credential( 911 1166 State(state): State<AppState>, 912 1167 xrpc_claims: XrpcClaims, 913 - Json(input): Json<RefreshCredentialInput>, 1168 + Json(input): Json<GetSpaceCredentialInput>, 914 1169 ) -> Result<Json<serde_json::Value>, AppError> { 915 - let _claims = require_auth(&xrpc_claims)?; 916 - let space = resolve_space(&state, &input.space_uri).await?; 1170 + let claims = require_auth(&xrpc_claims)?; 917 1171 918 1172 let encryption_key = state.config.token_encryption_key.as_ref().ok_or_else(|| { 919 1173 AppError::Internal("TOKEN_ENCRYPTION_KEY is required for space credentials".into()) 920 1174 })?; 921 1175 922 - let issued = crate::spaces::auth::refresh_credential( 1176 + let grant_claims = crate::spaces::credential::verify_grant(&input.grant, encryption_key)?; 1177 + 1178 + let space = resolve_space(&state, &grant_claims.space).await?; 1179 + 1180 + let client_id = claims.client_key().map(|k| k.to_string()); 1181 + let issued = crate::spaces::auth::issue_credential( 923 1182 &state.db, 924 1183 state.db_backend, 925 1184 encryption_key, 926 1185 &space, 927 - &input.credential, 1186 + &grant_claims.sub, 1187 + client_id.as_deref(), 928 1188 ) 929 1189 .await?; 930 1190 ··· 933 1193 "expiresAt": issued.expires_at, 934 1194 }))) 935 1195 } 1196 + 1197 + #[cfg(test)] 1198 + mod tests { 1199 + use super::*; 1200 + use serde_json::json; 1201 + 1202 + #[test] 1203 + fn content_cid_deterministic() { 1204 + let record = json!({"text": "hello"}); 1205 + let cid1 = content_cid(&record); 1206 + let cid2 = content_cid(&record); 1207 + assert_eq!(cid1, cid2); 1208 + assert!(cid1.starts_with("bafyrei")); 1209 + } 1210 + 1211 + #[test] 1212 + fn content_cid_changes_for_different_records() { 1213 + let a = content_cid(&json!({"text": "hello"})); 1214 + let b = content_cid(&json!({"text": "world"})); 1215 + assert_ne!(a, b); 1216 + } 1217 + 1218 + #[test] 1219 + fn deserialize_create_record_input() { 1220 + let input: CreateRecordInput = serde_json::from_value(json!({ 1221 + "space": "ats://did:plc:abc/com.example.forum/main", 1222 + "collection": "com.example.forum.post", 1223 + "record": { "text": "hello" } 1224 + })) 1225 + .unwrap(); 1226 + assert_eq!(input.space, "ats://did:plc:abc/com.example.forum/main"); 1227 + assert_eq!(input.collection, "com.example.forum.post"); 1228 + assert_eq!(input.record["text"], "hello"); 1229 + } 1230 + 1231 + #[test] 1232 + fn deserialize_put_record_with_swap() { 1233 + let input: PutRecordInput = serde_json::from_value(json!({ 1234 + "space": "ats://did:plc:abc/com.example.forum/main", 1235 + "collection": "com.example.forum.post", 1236 + "rkey": "3k2abc", 1237 + "record": { "text": "updated" }, 1238 + "swapRecord": "bafyrei123" 1239 + })) 1240 + .unwrap(); 1241 + assert_eq!(input.swap_record.as_deref(), Some("bafyrei123")); 1242 + } 1243 + 1244 + #[test] 1245 + fn deserialize_put_record_without_swap() { 1246 + let input: PutRecordInput = serde_json::from_value(json!({ 1247 + "space": "ats://did:plc:abc/com.example.forum/main", 1248 + "collection": "com.example.forum.post", 1249 + "rkey": "3k2abc", 1250 + "record": { "text": "hello" } 1251 + })) 1252 + .unwrap(); 1253 + assert_eq!(input.swap_record, None); 1254 + } 1255 + 1256 + #[test] 1257 + fn deserialize_delete_record_with_swap() { 1258 + let input: DeleteRecordInput = serde_json::from_value(json!({ 1259 + "space": "ats://did:plc:abc/com.example.forum/main", 1260 + "collection": "com.example.forum.post", 1261 + "rkey": "3k2abc", 1262 + "swapRecord": "bafyrei456" 1263 + })) 1264 + .unwrap(); 1265 + assert_eq!(input.swap_record.as_deref(), Some("bafyrei456")); 1266 + } 1267 + 1268 + #[test] 1269 + fn deserialize_write_op_create() { 1270 + let op: WriteOp = serde_json::from_value(json!({ 1271 + "action": "create", 1272 + "collection": "com.example.forum.post", 1273 + "value": { "text": "new post" } 1274 + })) 1275 + .unwrap(); 1276 + match op { 1277 + WriteOp::Create { 1278 + collection, 1279 + rkey, 1280 + value, 1281 + } => { 1282 + assert_eq!(collection, "com.example.forum.post"); 1283 + assert_eq!(rkey, None); 1284 + assert_eq!(value["text"], "new post"); 1285 + } 1286 + _ => panic!("expected Create"), 1287 + } 1288 + } 1289 + 1290 + #[test] 1291 + fn deserialize_write_op_create_with_rkey() { 1292 + let op: WriteOp = serde_json::from_value(json!({ 1293 + "action": "create", 1294 + "collection": "com.example.forum.post", 1295 + "rkey": "custom-key", 1296 + "value": { "text": "new post" } 1297 + })) 1298 + .unwrap(); 1299 + match op { 1300 + WriteOp::Create { rkey, .. } => { 1301 + assert_eq!(rkey.as_deref(), Some("custom-key")); 1302 + } 1303 + _ => panic!("expected Create"), 1304 + } 1305 + } 1306 + 1307 + #[test] 1308 + fn deserialize_write_op_update() { 1309 + let op: WriteOp = serde_json::from_value(json!({ 1310 + "action": "update", 1311 + "collection": "com.example.forum.post", 1312 + "rkey": "3k2abc", 1313 + "value": { "text": "updated" }, 1314 + "swapRecord": "bafyrei789" 1315 + })) 1316 + .unwrap(); 1317 + match op { 1318 + WriteOp::Update { 1319 + collection, 1320 + rkey, 1321 + swap_record, 1322 + .. 1323 + } => { 1324 + assert_eq!(collection, "com.example.forum.post"); 1325 + assert_eq!(rkey, "3k2abc"); 1326 + assert_eq!(swap_record.as_deref(), Some("bafyrei789")); 1327 + } 1328 + _ => panic!("expected Update"), 1329 + } 1330 + } 1331 + 1332 + #[test] 1333 + fn deserialize_write_op_delete() { 1334 + let op: WriteOp = serde_json::from_value(json!({ 1335 + "action": "delete", 1336 + "collection": "com.example.forum.post", 1337 + "rkey": "3k2abc" 1338 + })) 1339 + .unwrap(); 1340 + match op { 1341 + WriteOp::Delete { 1342 + collection, 1343 + rkey, 1344 + swap_record, 1345 + } => { 1346 + assert_eq!(collection, "com.example.forum.post"); 1347 + assert_eq!(rkey, "3k2abc"); 1348 + assert_eq!(swap_record, None); 1349 + } 1350 + _ => panic!("expected Delete"), 1351 + } 1352 + } 1353 + 1354 + #[test] 1355 + fn deserialize_apply_writes_input() { 1356 + let input: ApplyWritesInput = serde_json::from_value(json!({ 1357 + "space": "ats://did:plc:abc/com.example.forum/main", 1358 + "swapCommit": "tid123", 1359 + "writes": [ 1360 + { 1361 + "action": "create", 1362 + "collection": "com.example.forum.post", 1363 + "value": { "text": "post 1" } 1364 + }, 1365 + { 1366 + "action": "delete", 1367 + "collection": "com.example.forum.post", 1368 + "rkey": "old-key" 1369 + } 1370 + ] 1371 + })) 1372 + .unwrap(); 1373 + assert_eq!(input.space, "ats://did:plc:abc/com.example.forum/main"); 1374 + assert_eq!(input.swap_commit.as_deref(), Some("tid123")); 1375 + assert_eq!(input.writes.len(), 2); 1376 + } 1377 + 1378 + #[test] 1379 + fn deserialize_apply_writes_without_swap_commit() { 1380 + let input: ApplyWritesInput = serde_json::from_value(json!({ 1381 + "space": "ats://did:plc:abc/com.example.forum/main", 1382 + "writes": [ 1383 + { 1384 + "action": "create", 1385 + "collection": "com.example.forum.post", 1386 + "value": { "text": "post" } 1387 + } 1388 + ] 1389 + })) 1390 + .unwrap(); 1391 + assert_eq!(input.swap_commit, None); 1392 + } 1393 + 1394 + #[test] 1395 + fn deserialize_write_op_rejects_unknown_action() { 1396 + let result = serde_json::from_value::<WriteOp>(json!({ 1397 + "action": "unknown", 1398 + "collection": "test", 1399 + "rkey": "key" 1400 + })); 1401 + assert!(result.is_err()); 1402 + } 1403 + }
+4 -1
src/spaces/types.rs
··· 72 72 #[derive(Debug, Clone, Serialize, Deserialize)] 73 73 pub struct Space { 74 74 pub id: String, 75 + pub did: String, 75 76 pub owner_did: String, 77 + #[serde(rename = "type")] 76 78 pub type_nsid: String, 77 79 pub skey: String, 78 80 pub display_name: Option<String>, ··· 82 84 pub app_denylist: Option<Vec<String>>, 83 85 pub managing_app_did: Option<String>, 84 86 pub config: SpaceConfig, 87 + pub revision: Option<String>, 85 88 pub created_at: String, 86 89 pub updated_at: String, 87 90 } ··· 100 103 pub struct SpaceMember { 101 104 pub id: String, 102 105 pub space_id: String, 103 - pub member_did: String, 106 + pub did: String, 104 107 pub access: SpaceAccess, 105 108 pub is_delegation: bool, 106 109 pub granted_by: Option<String>,