atproto Thingiverse but good
12

Configure Feed

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

PM-32: migrate projection to the native Data value type

Orual (Jun 29, 2026, 7:06 PM EDT) d6a6a08d da5cf4aa

+356 -175
+9 -1
src/appview/proxy.rs
··· 34 34 use jacquard_axum::oauth::ExtractOAuthSession; 35 35 use jacquard_common::error::{ClientError, ClientErrorKind, HttpError}; 36 36 use jacquard_common::types::string::Did; 37 + use jacquard_common::types::value::to_data; 37 38 use jacquard_common::xrpc::{Response as XrpcResponse, XrpcClient, XrpcResp}; 38 39 use polymodel_api::com_atproto::identity::resolve_handle::ResolveHandleRequest; 39 40 use polymodel_api::com_atproto::repo::{ ··· 329 330 let Some((rkey, cid)) = parse_record_ref(output_buffer) else { 330 331 return; 331 332 }; 333 + let record_data = match to_data(record) { 334 + Ok(d) => d, 335 + Err(e) => { 336 + tracing::warn!(error = %e, %collection, "proxied record to_data failed"); 337 + return; 338 + } 339 + }; 332 340 let event = ProjectionEvent { 333 341 seq: 0, 334 342 did: did.as_ref().to_string(), 335 343 collection: collection.to_string(), 336 344 rkey, 337 345 action: action.to_string(), 338 - record: Some(record.clone()), 346 + record: Some(record_data), 339 347 cid: Some(cid), 340 348 }; 341 349 // Serialize against PM-28's decision-making writes (which read local state
+6 -4
src/appview/tests.rs
··· 955 955 956 956 use jacquard_common::deps::bytes::Bytes; 957 957 use jacquard_common::error::{AuthError, ClientError}; 958 + use jacquard_common::types::value::{Data, to_data}; 958 959 use jacquard_common::xrpc::Response as XrpcResponse; 959 960 use polymodel_api::com_atproto::repo::get_record::GetRecordResponse; 960 - use serde_json::Value; 961 961 962 962 use crate::indexing::projection::{ 963 963 ProjectionEvent, ProjectionInput, project_eager_record, project_event, ··· 1124 1124 assert!(body.contains("fresh"), "{body}"); 1125 1125 } 1126 1126 1127 - fn poly_event(action: &str, record: Option<Value>) -> ProjectionEvent { 1127 + fn poly_event(action: &str, record: Option<Data>) -> ProjectionEvent { 1128 1128 ProjectionEvent { 1129 1129 seq: 1, 1130 1130 did: DID_A.to_string(), ··· 1140 1140 async fn eager_write_marks_pending_and_firehose_clears_it() { 1141 1141 let pool = pool().await; 1142 1142 let uri = format!("at://{DID_A}/space.polymodel.library.thing/t1"); 1143 - let record = 1144 - json!({"name": "x", "license": "CC-BY-4.0", "createdAt": "2024-01-01T00:00:00.000Z"}); 1143 + let record = to_data( 1144 + &json!({"name": "x", "license": "CC-BY-4.0", "createdAt": "2024-01-01T00:00:00.000Z"}), 1145 + ) 1146 + .unwrap(); 1145 1147 project_eager_record(&pool, &poly_event("create", Some(record.clone()))) 1146 1148 .await 1147 1149 .unwrap();
+9 -14
src/appview/views.rs
··· 1 1 //! SQLite read queries and generated-view builders for the read endpoints. 2 2 //! 3 3 //! Records come from the `record_json` column (written by the projection); the 4 - //! required `record` field on thing/model/part views is built from it via 5 - //! [`to_data`]. Author handles come from `identities` (refreshed by identity 6 - //! events) and profile display fields from `profiles`. Authored image/blob fields 7 - //! are deserialized through generated lexicon types and resolved to Bluesky CDN 8 - //! URLs in appview response fields; no fake CDN/media data is synthesized. 4 + //! required `record` field on thing/model/part views is built from it. Author 5 + //! handles come from `identities` (refreshed by identity events) and profile 6 + //! display fields from `profiles`. Authored image/blob fields are deserialized 7 + //! through generated lexicon types and resolved to Bluesky CDN URLs in appview 8 + //! response fields; no fake CDN/media data is synthesized. 9 9 //! 10 10 //! Cursor contract: clients treat `cursor` as opaque. 11 11 //! - Chronological feeds (`getAuthorThings`, `getFeed?algorithm=recent`) keyset ··· 25 25 use jacquard::oauth::client::OAuthSession; 26 26 use jacquard_common::deps::smol_str::SmolStr; 27 27 use jacquard_common::types::string::{AtUri, Datetime, Did, Handle, UriValue}; 28 - use jacquard_common::types::{ 29 - blob::BlobRef, 30 - value::{Data, to_data}, 31 - }; 28 + use jacquard_common::types::{blob::BlobRef, value::Data}; 32 29 use polymodel_api::space_polymodel::library::{ 33 30 Actor, FeedItem, FeedView, File, Image, ImageView, ModelView, PartView, ThingView, 34 31 ThingViewBasic, ViewerState, model::Model, part::Part, thing::Thing, ··· 67 64 /// Build the required `record: Data` field from a stored `record_json` blob. 68 65 fn record_data(record_json: Option<&str>) -> AppResult<Data> { 69 66 let raw = record_json.ok_or_else(|| internal("content row missing record_json"))?; 70 - let value: serde_json::Value = 71 - serde_json::from_str(raw).map_err(|e| internal(format!("record_json parse: {e}")))?; 72 - to_data(&value).map_err(|e| internal(format!("record encode: {e}"))) 67 + serde_json::from_str::<Data>(raw).map_err(|e| internal(format!("record_json parse: {e}"))) 73 68 } 74 69 75 70 fn file_data(file_json: Option<&str>) -> AppResult<File> { ··· 1154 1149 }; 1155 1150 let stats = profile_stats(state, did.as_ref()).await?; 1156 1151 let handle_str = resolve_handle(state, did).await?; 1157 - let record_value: serde_json::Value = serde_json::from_str(&row.record_json) 1152 + let record: Data = serde_json::from_str(&row.record_json) 1158 1153 .map_err(|e| internal(format!("profile record_json parse: {e}")))?; 1159 1154 Ok(Some(polymodel_api::space_polymodel::actor::ProfileView { 1160 1155 did: did.clone(), 1161 1156 handle: s(handle_str), 1162 - record: to_data(&record_value).map_err(|e| internal(format!("record encode: {e}")))?, 1157 + record, 1163 1158 display_name: row.display_name.map(s), 1164 1159 description: row.description.map(s), 1165 1160 avatar: avatar_url(did.as_ref(), row.avatar_json.as_deref())?,
+2 -2
src/appview/writes.rs
··· 22 22 use jacquard_common::types::ident::AtIdentifier; 23 23 use jacquard_common::types::recordkey::{RecordKey, Rkey}; 24 24 use jacquard_common::types::string::{AtUri, Cid, Datetime, Did}; 25 + use jacquard_common::types::value::to_data; 25 26 use polymodel_api::app_bsky::actor::profile::Profile as BskyProfile; 26 27 use polymodel_api::com_atproto::repo::strong_ref::StrongRef; 27 28 use polymodel_api::space_polymodel::actor::ProfileView; ··· 1255 1256 record: &R, 1256 1257 action: &str, 1257 1258 ) -> AppResult<()> { 1258 - let record = 1259 - serde_json::to_value(record).map_err(|e| internal(format!("record serialize: {e}")))?; 1259 + let record = to_data(record).map_err(|e| internal(format!("record serialize: {e}")))?; 1260 1260 let event = ProjectionEvent { 1261 1261 seq: 0, 1262 1262 did: actor.as_ref().to_string(),
+330 -154
src/indexing/projection.rs
··· 10 10 //! subject of the existing row, delete it, and only decrement if a row was 11 11 //! actually removed. 12 12 13 - use serde_json::Value; 13 + use jacquard_common::types::value::{Data, DataDeserializerError, from_data, to_data}; 14 + use polymodel_api::space_polymodel::{ 15 + actor::profile::Profile, 16 + graph::{like::Like, listitem::Listitem, save::Save, tag::Tag}, 17 + library::{model::Model, part::Part, thing::Thing}, 18 + }; 14 19 use sqlx::{Sqlite, SqlitePool, Transaction}; 15 20 16 21 /// A repo-record projection event, decoupled from Hydrant's internal types so ··· 24 29 pub rkey: String, 25 30 /// `"create"`, `"update"`, or `"delete"`. 26 31 pub action: String, 27 - /// The record body for create/update events; `None` for delete events. 28 - pub record: Option<Value>, 32 + /// The record body for create/update events, carried as the native AT 33 + /// Protocol value type; `None` for delete events. 34 + pub record: Option<Data>, 29 35 pub cid: Option<String>, 30 36 } 31 37 ··· 44 50 /// Account status events are not projected and are dropped by 45 51 /// [`ProjectionInput::from_hydrant`]. 46 52 #[derive(Debug, Clone)] 53 + #[allow(clippy::large_enum_variant)] 47 54 pub enum ProjectionInput { 48 55 Record(ProjectionEvent), 49 56 Identity(IdentityProjection), ··· 63 70 collection: rec.collection.to_string(), 64 71 rkey: rec.rkey.to_string(), 65 72 action: rec.action.to_string(), 66 - record: rec.record.clone(), 73 + record: rec.record.as_ref().and_then(|v| { 74 + to_data(v) 75 + .map_err(|e| tracing::warn!(error = %e, "to_data conversion failed")) 76 + .ok() 77 + }), 67 78 cid: rec.cid.as_ref().map(|c| c.to_string()), 68 79 }))); 69 80 } ··· 143 154 let (Some(cid), Some(record)) = (event.cid.clone(), event.record.as_ref()) else { 144 155 return Ok(()); 145 156 }; 146 - let value = record.to_string(); 157 + let value = serde_json::to_string(record).unwrap_or_default(); 147 158 let now = now_nanos(); 148 159 sqlx::query!( 149 160 "INSERT INTO pending_writes (uri, cid, value, written_at) VALUES (?, ?, ?, ?) \ ··· 282 293 } 283 294 284 295 // --------------------------------------------------------------------------- 285 - // Field extraction helpers (records arrive as serde_json::Value). 296 + // Typed record decode from borrowed Data. 286 297 // --------------------------------------------------------------------------- 287 298 288 - fn s_str(v: &Value, key: &str) -> Option<String> { 289 - v.get(key).and_then(|x| x.as_str()).map(String::from) 290 - } 291 - 292 - /// URI from a single strongRef field, e.g. `derivedFrom: { uri, cid }`. 293 - fn ref_uri(v: &Value, key: &str) -> Option<String> { 294 - v.get(key) 295 - .and_then(|r| r.get("uri")) 296 - .and_then(|u| u.as_str()) 297 - .map(String::from) 298 - } 299 - 300 - /// URIs from an array of strongRefs, e.g. `models: [{ uri, cid }, ...]`. 301 - fn ref_uris(v: &Value, key: &str) -> Vec<String> { 302 - v.get(key) 303 - .and_then(|x| x.as_array()) 304 - .map(|arr| { 305 - arr.iter() 306 - .filter_map(|el| el.get("uri").and_then(|u| u.as_str()).map(String::from)) 307 - .collect() 308 - }) 309 - .unwrap_or_default() 310 - } 311 - 312 - fn str_array(v: &Value, key: &str) -> Vec<String> { 313 - v.get(key) 314 - .and_then(|x| x.as_array()) 315 - .map(|arr| { 316 - arr.iter() 317 - .filter_map(|el| el.as_str().map(String::from)) 318 - .collect() 319 - }) 320 - .unwrap_or_default() 321 - } 322 - 323 - fn parse_nanos(v: &Value, key: &str) -> Option<i64> { 324 - chrono::DateTime::parse_from_rfc3339(v.get(key)?.as_str()?) 325 - .ok() 326 - .and_then(|dt| dt.timestamp_nanos_opt()) 299 + /// Decode a typed record from a borrowed `&Data` body, injecting the 300 + /// collection NSID as the `$type` tag that all polymodel record structs are 301 + /// internally tagged on. Wire bodies and test fixtures may omit `$type`; 302 + /// injecting it at this single decode chokepoint — through which every record 303 + /// flows (firehose, eager writes, proxy, fixtures) — lets `from_data` succeed 304 + /// uniformly. The tag is always inserted/replaced; this is safe because 305 + /// dispatch is collection-keyed, so the injected value always matches the 306 + /// target struct's `#[serde(rename)]`. 307 + fn decode<T>(data: &Data, collection: &str) -> Result<T, DataDeserializerError> 308 + where 309 + T: for<'de> serde::Deserialize<'de>, 310 + { 311 + let mut tagged = data.clone(); 312 + if let Some(obj) = tagged.as_object_mut() { 313 + obj.0.insert( 314 + "$type".into(), 315 + to_data(&collection.to_string()).expect("serializing a string cannot fail"), 316 + ); 317 + } 318 + from_data(&tagged) 327 319 } 328 320 329 321 fn now_nanos() -> i64 { ··· 344 336 tx: &mut Transaction<'_, Sqlite>, 345 337 e: &ProjectionEvent, 346 338 ) -> anyhow::Result<()> { 347 - let Some(rec) = &e.record else { 339 + let Some(data) = &e.record else { 348 340 return Ok(()); 341 + }; 342 + let rec = match decode::<Thing>(data, &e.collection) { 343 + Ok(r) => r, 344 + Err(err) => { 345 + tracing::warn!(error = %err, "thing record decode failed"); 346 + return Ok(()); 347 + } 349 348 }; 350 349 let uri = at_uri(&e.did, &e.collection, &e.rkey); 351 - // Required content fields: skip the event gracefully if absent (malformed 352 - // record) rather than storing a blank row. 353 350 let Some(cid) = e.cid.clone() else { 354 351 return Ok(()); 355 352 }; 356 - let Some(name) = s_str(rec, "name") else { 357 - return Ok(()); 358 - }; 359 - let summary = s_str(rec, "summary"); 360 - let license = s_str(rec, "license"); 361 - let instructions = s_str(rec, "instructions"); 362 - let cover_json = rec.get("cover").map(|v| v.to_string()); 363 - let derived_from_uri = ref_uri(rec, "derivedFrom"); 353 + let name = rec.name.to_string(); 354 + let summary = rec.summary.as_ref().map(|s| s.to_string()); 355 + let license = rec.license.to_string(); 356 + let instructions = rec.instructions.as_ref().map(|lines| { 357 + lines 358 + .iter() 359 + .map(|s| s.to_string()) 360 + .collect::<Vec<_>>() 361 + .join("\n") 362 + }); 363 + let cover_json = data 364 + .as_object() 365 + .and_then(|o| o.0.get("cover")) 366 + .and_then(|v| serde_json::to_string(v).ok()); 367 + let derived_from_uri = rec 368 + .derived_from 369 + .as_ref() 370 + .map(|r| r.uri.as_str().to_string()); 364 371 365 - let tags = str_array(rec, "tags"); 372 + let tags = rec 373 + .tags 374 + .as_ref() 375 + .map(|v| v.iter().map(|s| s.to_string()).collect::<Vec<_>>()) 376 + .unwrap_or_default(); 366 377 let tags_json = if tags.is_empty() { 367 378 None 368 379 } else { ··· 374 385 Some(tags.join(" ")) 375 386 }; 376 387 377 - let created_at = parse_nanos(rec, "createdAt").unwrap_or_else(now_nanos); 388 + let created_at = rec 389 + .created_at 390 + .as_ref() 391 + .timestamp_nanos_opt() 392 + .unwrap_or_else(now_nanos); 378 393 let indexed_at = now_nanos(); 379 - // The full raw record, so read endpoints can populate the required `record` 380 - // field on thingView directly from SQLite (and the write path gets 381 - // read-your-own-write, since Hydrant has no public single-record write API). 382 - let record_json = rec.to_string(); 394 + let record_json = serde_json::to_string(data).unwrap_or_default(); 383 395 384 396 // UPSERT (not INSERT OR REPLACE): REPLACE performs DELETE+INSERT with a new 385 397 // rowid, which skips the AFTER UPDATE FTS trigger and (with recursive_triggers ··· 421 433 sqlx::query!("DELETE FROM thing_models WHERE thing_uri = ?", uri) 422 434 .execute(&mut **tx) 423 435 .await?; 424 - for (pos, model_uri) in ref_uris(rec, "models").into_iter().enumerate() { 436 + for (pos, model_ref) in rec.models.iter().flatten().enumerate() { 425 437 sqlx::query!( 426 438 "INSERT OR REPLACE INTO thing_models (thing_uri, model_uri, position) \ 427 439 VALUES (?, ?, ?)", 428 440 uri, 429 - model_uri, 441 + model_ref.uri.as_str(), 430 442 pos as i64, 431 443 ) 432 444 .execute(&mut **tx) ··· 470 482 tx: &mut Transaction<'_, Sqlite>, 471 483 e: &ProjectionEvent, 472 484 ) -> anyhow::Result<()> { 473 - let Some(rec) = &e.record else { 485 + let Some(data) = &e.record else { 474 486 return Ok(()); 475 487 }; 488 + let rec = match decode::<Model>(data, &e.collection) { 489 + Ok(r) => r, 490 + Err(err) => { 491 + tracing::warn!(error = %err, "model record decode failed"); 492 + return Ok(()); 493 + } 494 + }; 476 495 let uri = at_uri(&e.did, &e.collection, &e.rkey); 477 496 let Some(cid) = e.cid.clone() else { 478 497 return Ok(()); 479 498 }; 480 - let Some(name) = s_str(rec, "name") else { 481 - return Ok(()); 482 - }; 483 - let summary = s_str(rec, "summary"); 484 - let created_at = parse_nanos(rec, "createdAt").unwrap_or_else(now_nanos); 499 + let name = rec.name.to_string(); 500 + let summary = rec.summary.as_ref().map(|s| s.to_string()); 501 + let created_at = rec 502 + .created_at 503 + .as_ref() 504 + .timestamp_nanos_opt() 505 + .unwrap_or_else(now_nanos); 485 506 let indexed_at = now_nanos(); 486 - let record_json = rec.to_string(); 507 + let record_json = serde_json::to_string(data).unwrap_or_default(); 487 508 488 509 sqlx::query!( 489 510 "INSERT INTO models \ ··· 509 530 sqlx::query!("DELETE FROM model_parts WHERE model_uri = ?", uri) 510 531 .execute(&mut **tx) 511 532 .await?; 512 - for (pos, part_uri) in ref_uris(rec, "parts").into_iter().enumerate() { 533 + for (pos, part_ref) in rec.parts.iter().enumerate() { 513 534 sqlx::query!( 514 535 "INSERT OR REPLACE INTO model_parts (model_uri, part_uri, position) \ 515 536 VALUES (?, ?, ?)", 516 537 uri, 517 - part_uri, 538 + part_ref.uri.as_str(), 518 539 pos as i64, 519 540 ) 520 541 .execute(&mut **tx) ··· 545 566 tx: &mut Transaction<'_, Sqlite>, 546 567 e: &ProjectionEvent, 547 568 ) -> anyhow::Result<()> { 548 - let Some(rec) = &e.record else { 569 + let Some(data) = &e.record else { 549 570 return Ok(()); 550 571 }; 572 + let rec = match decode::<Part>(data, &e.collection) { 573 + Ok(r) => r, 574 + Err(err) => { 575 + tracing::warn!(error = %err, "part record decode failed"); 576 + return Ok(()); 577 + } 578 + }; 551 579 let uri = at_uri(&e.did, &e.collection, &e.rkey); 552 580 let Some(cid) = e.cid.clone() else { 553 581 return Ok(()); 554 582 }; 555 - let Some(name) = s_str(rec, "name") else { 556 - return Ok(()); 557 - }; 558 - let format = s_str(rec, "format"); 559 - let file_json = rec.get("file").map(|v| v.to_string()); 560 - let created_at = parse_nanos(rec, "createdAt").unwrap_or_else(now_nanos); 583 + let name = rec.name.to_string(); 584 + let format = rec.format.as_ref().map(|s| s.to_string()); 585 + let file_json = data 586 + .as_object() 587 + .and_then(|o| o.0.get("file")) 588 + .and_then(|v| serde_json::to_string(v).ok()); 589 + let created_at = rec 590 + .created_at 591 + .as_ref() 592 + .timestamp_nanos_opt() 593 + .unwrap_or_else(now_nanos); 561 594 let indexed_at = now_nanos(); 562 - let record_json = rec.to_string(); 595 + let record_json = serde_json::to_string(data).unwrap_or_default(); 563 596 564 597 sqlx::query!( 565 598 "INSERT INTO parts \ ··· 613 646 tx: &mut Transaction<'_, Sqlite>, 614 647 e: &ProjectionEvent, 615 648 ) -> anyhow::Result<()> { 616 - let Some(rec) = &e.record else { 649 + let Some(data) = &e.record else { 617 650 return Ok(()); 618 651 }; 619 - let display_name = s_str(rec, "displayName"); 620 - let description = s_str(rec, "description"); 621 - let avatar_json = rec.get("avatar").map(|v| v.to_string()); 622 - let default_license = s_str(rec, "defaultLicense"); 623 - let pronouns = s_str(rec, "pronouns"); 624 - let printers_json = rec 625 - .get("printers") 626 - .and_then(|v| v.as_array()) 627 - .map(|a| serde_json::to_string(a).unwrap_or_else(|_| "[]".to_string())); 628 - let links_json = rec 629 - .get("links") 630 - .and_then(|v| v.as_array()) 631 - .map(|a| serde_json::to_string(a).unwrap_or_else(|_| "[]".to_string())); 632 - let record_json = rec.to_string(); 652 + let rec = match decode::<Profile>(data, &e.collection) { 653 + Ok(r) => r, 654 + Err(err) => { 655 + tracing::warn!(error = %err, "profile record decode failed"); 656 + return Ok(()); 657 + } 658 + }; 659 + let display_name = rec.display_name.as_ref().map(|s| s.to_string()); 660 + let description = rec.description.as_ref().map(|s| s.to_string()); 661 + let avatar_json = data 662 + .as_object() 663 + .and_then(|o| o.0.get("avatar")) 664 + .and_then(|v| serde_json::to_string(v).ok()); 665 + let default_license = rec.default_license.as_ref().map(|s| s.to_string()); 666 + let pronouns = rec.pronouns.as_ref().map(|s| s.to_string()); 667 + let printers_json = data 668 + .as_object() 669 + .and_then(|o| o.0.get("printers")) 670 + .and_then(|v| serde_json::to_string(v).ok()); 671 + let links_json = data 672 + .as_object() 673 + .and_then(|o| o.0.get("links")) 674 + .and_then(|v| serde_json::to_string(v).ok()); 675 + let record_json = serde_json::to_string(data).unwrap_or_default(); 633 676 let indexed_at = now_nanos(); 634 677 635 678 sqlx::query!( ··· 703 746 tx: &mut Transaction<'_, Sqlite>, 704 747 e: &ProjectionEvent, 705 748 ) -> anyhow::Result<()> { 706 - let Some(rec) = &e.record else { 749 + let Some(data) = &e.record else { 707 750 return Ok(()); 708 751 }; 709 - let Some(subject) = ref_uri(rec, "subject") else { 710 - return Ok(()); 752 + let rec = match decode::<Like>(data, &e.collection) { 753 + Ok(r) => r, 754 + Err(err) => { 755 + tracing::warn!(error = %err, "like record decode failed"); 756 + return Ok(()); 757 + } 711 758 }; 759 + let subject = rec.subject.uri.as_str().to_string(); 712 760 let Some(cid) = e.cid.clone() else { 713 761 return Ok(()); 714 762 }; 715 - let created_at = parse_nanos(rec, "createdAt").unwrap_or_else(now_nanos); 763 + let created_at = rec 764 + .created_at 765 + .as_ref() 766 + .timestamp_nanos_opt() 767 + .unwrap_or_else(now_nanos); 716 768 717 769 let res = sqlx::query!( 718 770 "INSERT OR IGNORE INTO likes (did, rkey, cid, subject_uri, created_at) VALUES (?, ?, ?, ?, ?)", ··· 775 827 tx: &mut Transaction<'_, Sqlite>, 776 828 e: &ProjectionEvent, 777 829 ) -> anyhow::Result<()> { 778 - let Some(rec) = &e.record else { 830 + let Some(data) = &e.record else { 779 831 return Ok(()); 780 832 }; 781 - let Some(subject) = ref_uri(rec, "subject") else { 782 - return Ok(()); 833 + let rec = match decode::<Save>(data, &e.collection) { 834 + Ok(r) => r, 835 + Err(err) => { 836 + tracing::warn!(error = %err, "save record decode failed"); 837 + return Ok(()); 838 + } 783 839 }; 840 + let subject = rec.subject.uri.as_str().to_string(); 784 841 let Some(cid) = e.cid.clone() else { 785 842 return Ok(()); 786 843 }; 787 - let note = s_str(rec, "note"); 788 - let created_at = parse_nanos(rec, "createdAt").unwrap_or_else(now_nanos); 844 + let note = rec.note.as_ref().map(|s| s.to_string()); 845 + let created_at = rec 846 + .created_at 847 + .as_ref() 848 + .timestamp_nanos_opt() 849 + .unwrap_or_else(now_nanos); 789 850 790 851 let res = sqlx::query!( 791 852 "INSERT OR IGNORE INTO saves (did, rkey, cid, subject_uri, note, created_at) \ ··· 850 911 tx: &mut Transaction<'_, Sqlite>, 851 912 e: &ProjectionEvent, 852 913 ) -> anyhow::Result<()> { 853 - let Some(rec) = &e.record else { 914 + let Some(data) = &e.record else { 854 915 return Ok(()); 855 916 }; 856 - let Some(subject) = ref_uri(rec, "subject") else { 857 - return Ok(()); 917 + let rec = match decode::<Tag>(data, &e.collection) { 918 + Ok(r) => r, 919 + Err(err) => { 920 + tracing::warn!(error = %err, "tag record decode failed"); 921 + return Ok(()); 922 + } 858 923 }; 859 - let Some(tag) = s_str(rec, "tag") else { 860 - return Ok(()); 861 - }; 924 + let subject = rec.subject.uri.as_str().to_string(); 925 + let tag = rec.tag.to_string(); 862 926 let Some(cid) = e.cid.clone() else { 863 927 return Ok(()); 864 928 }; 865 - let created_at = parse_nanos(rec, "createdAt").unwrap_or_else(now_nanos); 929 + let created_at = rec 930 + .created_at 931 + .as_ref() 932 + .timestamp_nanos_opt() 933 + .unwrap_or_else(now_nanos); 866 934 867 935 let res = sqlx::query!( 868 936 "INSERT OR IGNORE INTO tags (did, rkey, cid, subject_uri, tag, created_at) \ ··· 923 991 tx: &mut Transaction<'_, Sqlite>, 924 992 e: &ProjectionEvent, 925 993 ) -> anyhow::Result<()> { 926 - let Some(rec) = &e.record else { 994 + let Some(data) = &e.record else { 927 995 return Ok(()); 928 996 }; 929 - let Some(list_uri) = ref_uri(rec, "list") else { 930 - return Ok(()); 997 + let rec = match decode::<Listitem>(data, &e.collection) { 998 + Ok(r) => r, 999 + Err(err) => { 1000 + tracing::warn!(error = %err, "listitem record decode failed"); 1001 + return Ok(()); 1002 + } 931 1003 }; 932 - let Some(subject_uri) = ref_uri(rec, "subject") else { 933 - return Ok(()); 934 - }; 1004 + let list_uri = rec.list.as_str().to_string(); 1005 + let subject_uri = rec.subject.uri.as_str().to_string(); 935 1006 let Some(cid) = e.cid.clone() else { 936 1007 return Ok(()); 937 1008 }; 938 - let created_at = parse_nanos(rec, "createdAt").unwrap_or_else(now_nanos); 1009 + let created_at = rec 1010 + .created_at 1011 + .as_ref() 1012 + .timestamp_nanos_opt() 1013 + .unwrap_or_else(now_nanos); 939 1014 940 1015 // No denormalized count for listitems; list membership is queried directly. 941 1016 sqlx::query!( ··· 1010 1085 collection: "space.polymodel.library.thing".to_string(), 1011 1086 rkey: rkey.to_string(), 1012 1087 action: "create".to_string(), 1013 - record: Some(json!({ 1088 + record: Some(to_data(&json!({ 1014 1089 "name": name, 1015 1090 "summary": "a thing summary", 1016 1091 "license": "MIT", "tags": ["cube", "calibration"], 1017 1092 "createdAt": "2024-01-15T12:00:00.000Z", 1018 - "models": models.iter().map(|m| json!({"uri": m, "cid": "c"})).collect::<Vec<_>>(), 1019 - })), 1093 + "models": models.iter().map(|m| json!({"uri": format!("at://{did}/space.polymodel.library.model/{m}"), "cid": "c"})).collect::<Vec<_>>(), 1094 + })).unwrap()), 1020 1095 cid: Some("thingcid".to_string()), 1021 1096 } 1022 1097 } ··· 1028 1103 collection: "space.polymodel.graph.like".to_string(), 1029 1104 rkey: rkey.to_string(), 1030 1105 action: "create".to_string(), 1031 - record: Some(json!({ 1032 - "subject": {"uri": subject, "cid": "c"}, 1033 - "createdAt": "2024-02-01T00:00:00.000Z", 1034 - })), 1106 + record: Some( 1107 + to_data(&json!({ 1108 + "subject": {"uri": subject, "cid": "c"}, 1109 + "createdAt": "2024-02-01T00:00:00.000Z", 1110 + })) 1111 + .unwrap(), 1112 + ), 1035 1113 cid: Some("likecid".to_string()), 1036 1114 } 1037 1115 } ··· 1125 1203 let db = test_db().await; 1126 1204 let did = "did:plc:abc"; 1127 1205 let uri = "at://did:plc:abc/space.polymodel.library.thing/r1"; 1128 - let e = thing_event(10, did, "r1", "Calibration Cube", &["at://m1", "at://m2"]); 1206 + let e = thing_event(10, did, "r1", "Calibration Cube", &["m1", "m2"]); 1129 1207 project_event(&db, &ProjectionInput::Record(e)) 1130 1208 .await 1131 1209 .expect("project thing"); ··· 1150 1228 // Create with two models. 1151 1229 project_event( 1152 1230 &db, 1153 - &ProjectionInput::Record(thing_event(1, did, "r1", "Cube", &["at://m1", "at://m2"])), 1231 + &ProjectionInput::Record(thing_event(1, did, "r1", "Cube", &["m1", "m2"])), 1154 1232 ) 1155 1233 .await 1156 1234 .unwrap(); 1157 1235 assert_eq!(junction_count(&db, uri).await, 2); 1158 1236 1159 1237 // Update: same rkey, different models. Old junction rows must be replaced. 1160 - let mut upd = thing_event(2, did, "r1", "Cube v2", &["at://m3"]); 1238 + let mut upd = thing_event(2, did, "r1", "Cube v2", &["m3"]); 1161 1239 upd.action = "update".to_string(); 1162 1240 project_event(&db, &ProjectionInput::Record(upd)) 1163 1241 .await ··· 1166 1244 assert_eq!(junction_count(&db, uri).await, 1); 1167 1245 // m1/m2 gone, m3 present. 1168 1246 let has_m3 = sqlx::query_as::<_, CountRow>( 1169 - "SELECT COUNT(*) AS count FROM thing_models WHERE thing_uri = ? AND model_uri = 'at://m3'", 1247 + "SELECT COUNT(*) AS count FROM thing_models WHERE thing_uri = ? AND model_uri = 'at://did:plc:abc/space.polymodel.library.model/m3'", 1170 1248 ) 1171 1249 .bind(uri) 1172 1250 .fetch_one(&db) ··· 1183 1261 let uri = "at://did:plc:abc/space.polymodel.library.thing/r1"; 1184 1262 project_event( 1185 1263 &db, 1186 - &ProjectionInput::Record(thing_event(1, did, "r1", "Cube", &["at://m1"])), 1264 + &ProjectionInput::Record(thing_event(1, did, "r1", "Cube", &["m1"])), 1187 1265 ) 1188 1266 .await 1189 1267 .unwrap(); ··· 1204 1282 #[tokio::test] 1205 1283 async fn like_create_increments_count() { 1206 1284 let db = test_db().await; 1207 - let subject = "at://did:plc:author/.../thing"; 1285 + let subject = "at://did:plc:author/space.polymodel.library.thing/t1"; 1208 1286 project_event( 1209 1287 &db, 1210 1288 &ProjectionInput::Record(like_event(1, "did:plc:a", "k1", subject)), ··· 1225 1303 #[tokio::test] 1226 1304 async fn like_delete_decrements_count_two_step() { 1227 1305 let db = test_db().await; 1228 - let subject = "at://did:plc:author/.../thing"; 1306 + let subject = "at://did:plc:author/space.polymodel.library.thing/t1"; 1229 1307 project_event( 1230 1308 &db, 1231 1309 &ProjectionInput::Record(like_event(1, "did:plc:a", "k1", subject)), ··· 1252 1330 async fn like_delete_when_absent_is_idempotent() { 1253 1331 // Deleting a like that was never there must not panic nor go negative. 1254 1332 let db = test_db().await; 1255 - let subject = "at://did:plc:author/.../thing"; 1333 + let subject = "at://did:plc:author/space.polymodel.library.thing/t1"; 1256 1334 project_event( 1257 1335 &db, 1258 1336 &ProjectionInput::Record(delete_event( ··· 1271 1349 async fn duplicate_like_does_not_double_count() { 1272 1350 // Replay: same like event twice (e.g. eager write + firehose delivery). 1273 1351 let db = test_db().await; 1274 - let subject = "at://did:plc:author/.../thing"; 1352 + let subject = "at://did:plc:author/space.polymodel.library.thing/t1"; 1275 1353 project_event( 1276 1354 &db, 1277 1355 &ProjectionInput::Record(like_event(1, "did:plc:a", "k1", subject)), ··· 1301 1379 assert_eq!(cursor(&db).await, Some(7)); 1302 1380 project_event( 1303 1381 &db, 1304 - &ProjectionInput::Record(like_event(99, "did:plc:b", "k1", "at://x")), 1382 + &ProjectionInput::Record(like_event( 1383 + 99, 1384 + "did:plc:b", 1385 + "k1", 1386 + "at://did:plc:author/space.polymodel.library.thing/x", 1387 + )), 1305 1388 ) 1306 1389 .await 1307 1390 .unwrap(); ··· 1311 1394 #[tokio::test] 1312 1395 async fn save_create_and_delete_maintain_count() { 1313 1396 let db = test_db().await; 1314 - let subject = "at://did:plc:author/.../thing"; 1397 + let subject = "at://did:plc:author/space.polymodel.library.thing/t1"; 1315 1398 let save = ProjectionEvent { 1316 1399 seq: 1, 1317 1400 did: "did:plc:a".to_string(), 1318 1401 collection: "space.polymodel.graph.save".to_string(), 1319 1402 rkey: "s1".to_string(), 1320 1403 action: "create".to_string(), 1321 - record: Some(json!({ 1322 - "subject": {"uri": subject, "cid": "c"}, 1323 - "note": "useful", 1324 - "createdAt": "2024-03-01T00:00:00.000Z", 1325 - })), 1404 + record: Some( 1405 + to_data(&json!({ 1406 + "subject": {"uri": subject, "cid": "c"}, 1407 + "note": "useful", 1408 + "createdAt": "2024-03-01T00:00:00.000Z", 1409 + })) 1410 + .unwrap(), 1411 + ), 1326 1412 cid: Some("savecid".to_string()), 1327 1413 }; 1328 1414 project_event(&db, &ProjectionInput::Record(save)) ··· 1361 1447 collection: "space.polymodel.library.thing".to_string(), 1362 1448 rkey: "r1".to_string(), 1363 1449 action: "create".to_string(), 1364 - record: Some(json!({ 1365 - "name": "Calibration Cube", 1366 - "license": "MIT", "tags": ["cube"], 1367 - "createdAt": "2024-01-01T00:00:00.000Z", 1368 - "models": [], 1369 - })), 1450 + record: Some( 1451 + to_data(&json!({ 1452 + "name": "Calibration Cube", 1453 + "license": "MIT", "tags": ["cube"], 1454 + "createdAt": "2024-01-01T00:00:00.000Z", 1455 + "models": [], 1456 + })) 1457 + .unwrap(), 1458 + ), 1370 1459 cid: Some("cid1".to_string()), 1371 1460 }; 1372 1461 project_event(&db, &ProjectionInput::Record(create.clone())) ··· 1378 1467 update.seq = 2; 1379 1468 update.action = "update".to_string(); 1380 1469 update.cid = Some("cid2".to_string()); 1381 - update.record = Some(json!({ 1382 - "name": "Friendly Sphere", 1383 - "license": "MIT", "tags": ["sphere"], 1384 - "createdAt": "2024-01-02T00:00:00.000Z", 1385 - "models": [], 1386 - })); 1470 + update.record = Some( 1471 + to_data(&json!({ 1472 + "name": "Friendly Sphere", 1473 + "license": "MIT", "tags": ["sphere"], 1474 + "createdAt": "2024-01-02T00:00:00.000Z", 1475 + "models": [], 1476 + })) 1477 + .unwrap(), 1478 + ); 1387 1479 project_event(&db, &ProjectionInput::Record(update)) 1388 1480 .await 1389 1481 .unwrap(); ··· 1404 1496 .unwrap() 1405 1497 .count; 1406 1498 assert_eq!(live, 1); 1499 + } 1500 + 1501 + #[tokio::test] 1502 + async fn malformed_and_incomplete_records_skip_gracefully() { 1503 + // A non-object body and a lexicon-incomplete record (Thing missing the 1504 + // required `createdAt`) must both degrade to a skip: project_event 1505 + // returns Ok, the cursor advances, and no partial row is written. 1506 + let db = test_db().await; 1507 + let did = "did:plc:abc"; 1508 + 1509 + let poison = ProjectionEvent { 1510 + seq: 1, 1511 + did: did.to_string(), 1512 + collection: "space.polymodel.library.thing".to_string(), 1513 + rkey: "bad1".to_string(), 1514 + action: "create".to_string(), 1515 + record: Some(to_data(&json!("not an object")).unwrap()), 1516 + cid: Some("cid".to_string()), 1517 + }; 1518 + project_event(&db, &ProjectionInput::Record(poison)) 1519 + .await 1520 + .expect("malformed record must not poison the stream"); 1521 + assert_eq!(cursor(&db).await, Some(1)); 1522 + 1523 + let incomplete = ProjectionEvent { 1524 + seq: 2, 1525 + did: did.to_string(), 1526 + collection: "space.polymodel.library.thing".to_string(), 1527 + rkey: "bad2".to_string(), 1528 + action: "create".to_string(), 1529 + record: Some(to_data(&json!({"name": "Incomplete", "license": "MIT"})).unwrap()), 1530 + cid: Some("cid".to_string()), 1531 + }; 1532 + project_event(&db, &ProjectionInput::Record(incomplete)) 1533 + .await 1534 + .expect("incomplete record must not poison the stream"); 1535 + assert_eq!(cursor(&db).await, Some(2)); 1536 + 1537 + let rows: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM things") 1538 + .fetch_one(&db) 1539 + .await 1540 + .unwrap(); 1541 + assert_eq!(rows, 0, "no partial rows should exist for skipped records"); 1542 + } 1543 + 1544 + #[tokio::test] 1545 + async fn listitem_create_projects_with_direct_list_uri() { 1546 + // Regression: the old ref_uri("list") expected a nested {uri, cid} 1547 + // object, but the lexicon type is a direct AtUri. Conformant listitems 1548 + // were silently never projected. The typed path now projects them. 1549 + let db = test_db().await; 1550 + let list_uri = "at://did:plc:abc/space.polymodel.graph.list/mylist"; 1551 + let subject_uri = "at://did:plc:xyz/space.polymodel.library.thing/t1"; 1552 + let event = ProjectionEvent { 1553 + seq: 10, 1554 + did: "did:plc:a".to_string(), 1555 + collection: "space.polymodel.graph.listitem".to_string(), 1556 + rkey: "li1".to_string(), 1557 + action: "create".to_string(), 1558 + record: Some( 1559 + to_data(&json!({ 1560 + "list": list_uri, 1561 + "subject": {"uri": subject_uri, "cid": "c"}, 1562 + "createdAt": "2024-04-01T00:00:00.000Z", 1563 + })) 1564 + .unwrap(), 1565 + ), 1566 + cid: Some("licid".to_string()), 1567 + }; 1568 + project_event(&db, &ProjectionInput::Record(event)) 1569 + .await 1570 + .unwrap(); 1571 + 1572 + let row: Option<(String, String)> = sqlx::query_as( 1573 + "SELECT list_uri, subject_uri FROM listitems WHERE did = ? AND rkey = ?", 1574 + ) 1575 + .bind("did:plc:a") 1576 + .bind("li1") 1577 + .fetch_optional(&db) 1578 + .await 1579 + .unwrap(); 1580 + let (stored_list, stored_subject) = row.expect("listitem must be projected"); 1581 + assert_eq!(stored_list, list_uri); 1582 + assert_eq!(stored_subject, subject_uri); 1407 1583 } 1408 1584 }