//! PM-43 server-side publish drafts + app-internal image upload. //! //! Drafts are app-internal (not ATProto records, not federated): owner-scoped //! rows in the `drafts` table holding a [`DraftThingInput`] JSON payload. They //! never reach the actor's PDS until publish, and they live in their own table //! so feeds (which read only `things`) never surface them. //! //! Every handler authenticates with the strict OAuth session (same extractor as //! the PM-28 writes) and scopes every query by the authenticated DID, so no //! actor can read, write, publish, or delete another actor's drafts. //! //! Routes (same-origin, cookie-authenticated; app-internal, not `/xrpc`): //! - `POST /app/drafts` create (server assigns an opaque id) //! - `PUT /app/drafts/{draft_id}` upsert an existing draft //! - `GET /app/drafts` list draft summaries //! - `GET /app/drafts/{draft_id}` fetch one draft payload //! - `DELETE /app/drafts/{draft_id}` delete a draft //! - `POST /app/drafts/{draft_id}/publish` assemble + publish + delete-on-success //! - `POST /app/images` upload an image blob, return an `#image` use axum::Json; use axum::body::Body; use axum::extract::{Path, State}; use axum::http::{HeaderMap, HeaderValue, Response}; use jacquard::client::Agent; use jacquard_axum::oauth::ExtractOAuthSession; use jacquard_common::types::blob::BlobRef; use jacquard_common::types::string::Did; use polymodel_api::space_polymodel::library::publish_thing::{PublishThingOutput, ThingInput}; use polymodel_api::space_polymodel::library::{AspectRatio, Image}; use super::content_type::content_type; use super::error::{AppResult, db, internal, invalid_request, not_found}; use super::state::AppState; use super::writes::{ ExtractSession, agent_upload_blob, authenticated_did, ensure_polymodel_profile, publish_composition, }; use crate::publish::draft::{ DeleteDraftResponse, DraftIdResponse, DraftSummary, DraftThingInput, ImageUploadResponse, }; /// Image blobs are buffered fully (cover/preview art is small); cap to keep the /// app from buffering an unbounded body. const MAX_IMAGE_SIZE: usize = 20 * 1024 * 1024; // ---------------------------------------------------------------------------- // Owner-scoped DB helpers (pure; unit-tested directly) // ---------------------------------------------------------------------------- fn now_nanos() -> i64 { chrono::Utc::now().timestamp_nanos_opt().unwrap_or_default() } /// Insert or update a draft row, refreshing the denormalized `name`/`model_count` /// and `updated_at`. `created_at` is preserved across updates. pub(super) async fn upsert_draft( state: &AppState, owner: &Did, draft_id: &str, draft: &DraftThingInput, ) -> AppResult<()> { let payload_json = serde_json::to_string(draft).map_err(|e| internal(format!("draft serialize: {e}")))?; let name = draft.display_name().map(ToOwned::to_owned); let model_count = draft.model_count(); let now = now_nanos(); let owner_did = owner.as_ref(); db(sqlx::query!( "INSERT INTO drafts (owner_did, draft_id, payload_json, name, model_count, created_at, updated_at) \ VALUES (?, ?, ?, ?, ?, ?, ?) \ ON CONFLICT(owner_did, draft_id) DO UPDATE SET \ payload_json = excluded.payload_json, \ name = excluded.name, \ model_count = excluded.model_count, \ updated_at = excluded.updated_at", owner_did, draft_id, payload_json, name, model_count, now, now, ) .execute(&state.pool) .await)?; Ok(()) } /// Load a single owner-scoped draft payload, if present. pub(super) async fn get_draft( state: &AppState, owner: &Did, draft_id: &str, ) -> AppResult> { let owner_did = owner.as_ref(); let row = db(sqlx::query!( "SELECT payload_json FROM drafts WHERE owner_did = ? AND draft_id = ?", owner_did, draft_id, ) .fetch_optional(&state.pool) .await)?; match row { Some(row) => { let draft: DraftThingInput = serde_json::from_str(&row.payload_json) .map_err(|e| internal(format!("draft deserialize: {e}")))?; Ok(Some(draft)) } None => Ok(None), } } /// List owner-scoped draft summaries, most-recently-updated first. pub(super) async fn list_drafts(state: &AppState, owner: &Did) -> AppResult> { let owner_did = owner.as_ref(); let rows = db(sqlx::query!( "SELECT draft_id, name, model_count, created_at, updated_at \ FROM drafts WHERE owner_did = ? ORDER BY updated_at DESC", owner_did, ) .fetch_all(&state.pool) .await)?; Ok(rows .into_iter() .map(|row| DraftSummary { draft_id: row.draft_id, name: row.name, model_count: row.model_count, created_at: row.created_at, updated_at: row.updated_at, }) .collect()) } /// Delete an owner-scoped draft; returns whether a row existed. pub(super) async fn delete_draft(state: &AppState, owner: &Did, draft_id: &str) -> AppResult { let owner_did = owner.as_ref(); let result = db(sqlx::query!( "DELETE FROM drafts WHERE owner_did = ? AND draft_id = ?", owner_did, draft_id, ) .execute(&state.pool) .await)?; Ok(result.rows_affected() > 0) } /// Promote a draft to a published thing: load → assemble → publish → delete. /// /// The publish step is injected so the ordering invariant (the draft row is /// deleted **only after** a successful publish) is unit-testable without a live /// PDS agent. If `publish` returns `Err`, the draft is left untouched. pub(super) async fn promote_draft( state: &AppState, owner: &Did, draft_id: &str, publish: F, ) -> AppResult where F: FnOnce(ThingInput) -> Fut, Fut: std::future::Future>, { let draft = get_draft(state, owner, draft_id) .await? .ok_or_else(not_found)?; let thing = draft .assemble() .map_err(|errors| invalid_request(errors.to_string()))?; let output = publish(thing).await?; // Reached only on a successful publish; a failure short-circuits above and // leaves the draft recoverable. delete_draft(state, owner, draft_id).await?; Ok(output) } // ---------------------------------------------------------------------------- // Handlers // ---------------------------------------------------------------------------- pub(super) async fn create_draft( State(state): State, ExtractOAuthSession(session): ExtractSession, Json(draft): Json, ) -> AppResult> { let agent = Agent::from(session); let actor = authenticated_did(&agent).await?; let draft_id = ulid::Ulid::new().to_string(); upsert_draft(&state, &actor, &draft_id, &draft).await?; Ok(Json(DraftIdResponse { draft_id })) } pub(super) async fn put_draft( State(state): State, ExtractOAuthSession(session): ExtractSession, Path(draft_id): Path, Json(draft): Json, ) -> AppResult> { let agent = Agent::from(session); let actor = authenticated_did(&agent).await?; upsert_draft(&state, &actor, &draft_id, &draft).await?; Ok(Json(DraftIdResponse { draft_id })) } pub(super) async fn get_draft_handler( State(state): State, ExtractOAuthSession(session): ExtractSession, Path(draft_id): Path, ) -> AppResult> { let agent = Agent::from(session); let actor = authenticated_did(&agent).await?; let draft = get_draft(&state, &actor, &draft_id) .await? .ok_or_else(not_found)?; Ok(Json(draft)) } pub(super) async fn list_drafts_handler( State(state): State, ExtractOAuthSession(session): ExtractSession, ) -> AppResult>> { let agent = Agent::from(session); let actor = authenticated_did(&agent).await?; let drafts = list_drafts(&state, &actor).await?; Ok(Json(drafts)) } pub(super) async fn delete_draft_handler( State(state): State, ExtractOAuthSession(session): ExtractSession, Path(draft_id): Path, ) -> AppResult> { let agent = Agent::from(session); let actor = authenticated_did(&agent).await?; let deleted = delete_draft(&state, &actor, &draft_id).await?; Ok(Json(DeleteDraftResponse { deleted })) } pub(super) async fn publish_draft( State(state): State, ExtractOAuthSession(session): ExtractSession, Path(draft_id): Path, ) -> AppResult> { let agent = Agent::from(session); let actor = authenticated_did(&agent).await?; let _ = ensure_polymodel_profile(&state, &agent, &actor, false).await?; let _write_guard = state.write_lock.lock().await; let output = promote_draft(&state, &actor, &draft_id, |thing| { publish_composition(&state, &agent, &actor, thing) }) .await?; Ok(Json(output)) } pub(super) async fn get_staged_resource( State(state): State, ExtractOAuthSession(session): ExtractSession, Path(upload_id): Path, ) -> AppResult> { let agent = Agent::from(session); let actor = authenticated_did(&agent).await?; let file = super::writes::load_upload_file(&state, &actor, &upload_id).await?; let size = u64::try_from(file.size) .map_err(|_| invalid_request("staged resource has negative length"))?; let digest = file .digest .as_deref() .ok_or_else(|| invalid_request("staged resource has no digest"))?; if size == 0 || size > crate::publish::draft::MAX_PREVIEW_RESOURCE_BYTES || digest.len() != 32 { return Err(invalid_request("staged resource exceeds preview bounds")); } super::downloads::validate_file_manifest(&file)?; let cids = super::downloads::ordered_chunks(&file) .iter() .map(|chunk| chunk.blob.blob().r#ref.as_str().to_owned()) .collect::>(); let bytes = super::ldraw::fetch_staged_bytes( &state, &actor, &cids, size, digest, crate::publish::draft::MAX_PREVIEW_RESOURCE_BYTES, ) .await?; let mut response = Response::new(Body::from(bytes)); response.headers_mut().insert( axum::http::header::CONTENT_TYPE, HeaderValue::from_str(file.mime_type.as_str()) .map_err(|_| invalid_request("invalid staged resource MIME"))?, ); Ok(response) } pub(super) async fn upload_image( State(state): State, ExtractOAuthSession(session): ExtractSession, headers: HeaderMap, body: Body, ) -> AppResult> { let agent = Agent::from(session); let actor = authenticated_did(&agent).await?; let _ = ensure_polymodel_profile(&state, &agent, &actor, false).await?; let mime_type = content_type(&headers); let alt = headers .get("x-polymodel-alt") .and_then(|v| v.to_str().ok()) .unwrap_or("") .to_string(); let bytes = axum::body::to_bytes(body, MAX_IMAGE_SIZE) .await .map_err(|_| invalid_request("image body exceeds the size limit or could not be read"))?; if bytes.is_empty() { return Err(invalid_request( "image upload requires a non-empty raw byte body", )); } // `aspectRatio` is required by `#image` but is not returned by uploadBlob, so // decode width/height from the bytes (header-only probe, no full decode). let dims = imagesize::blob_size(&bytes) .map_err(|e| invalid_request(format!("unsupported or corrupt image: {e}")))?; let blob = agent_upload_blob(&agent, bytes.to_vec(), &mime_type).await?; let image = Image { alt: alt.into(), aspect_ratio: AspectRatio { height: dims.height as i64, width: dims.width as i64, extra_data: None, }, image: BlobRef::from(blob), extra_data: None, }; Ok(Json(ImageUploadResponse { image })) } #[cfg(test)] mod tests { use std::str::FromStr; use jacquard_common::types::string::Did; use sqlx::sqlite::{SqliteConnectOptions, SqlitePoolOptions}; use super::*; use crate::publish::draft::{DraftModelInput, DraftPartInput}; const DID_A: &str = "did:plc:aaaaaaaaaaaaaaaaaaaaaaaa"; const DID_B: &str = "did:plc:bbbbbbbbbbbbbbbbbbbbbbbb"; async fn state() -> AppState { let options = SqliteConnectOptions::from_str("sqlite::memory:").unwrap(); let pool = SqlitePoolOptions::new() .max_connections(1) .connect_with(options) .await .unwrap(); sqlx::migrate!("./migrations").run(&pool).await.unwrap(); let bootstrap = crate::oauth::bootstrap_oauth(pool.clone(), Some("http://localhost")) .expect("ephemeral OAuth bootstrap for tests"); AppState::new(pool, bootstrap) } fn did(s: &str) -> Did { Did::new_owned(s).unwrap() } fn draft_named(name: &str) -> DraftThingInput { DraftThingInput { name: Some(name.to_string()), license: Some("CC-BY-4.0".to_string()), models: vec![DraftModelInput { name: Some("Model".to_string()), parts: vec![DraftPartInput { name: Some("part.stl".to_string()), upload_id: Some("abcdef0123456789-42".to_string()), ..Default::default() }], ..Default::default() }], ..Default::default() } } #[tokio::test] async fn upsert_and_get_round_trips_owner_scoped() { let state = state().await; let owner = did(DID_A); let draft = draft_named("Widget"); upsert_draft(&state, &owner, "d1", &draft).await.unwrap(); let loaded = get_draft(&state, &owner, "d1").await.unwrap(); assert_eq!(loaded, Some(draft)); // A different owner cannot see it. let other = get_draft(&state, &did(DID_B), "d1").await.unwrap(); assert_eq!(other, None); } #[tokio::test] async fn upsert_updates_denormalized_summary_and_preserves_created_at() { let state = state().await; let owner = did(DID_A); upsert_draft(&state, &owner, "d1", &draft_named("First")) .await .unwrap(); let created_at: i64 = sqlx::query_scalar("SELECT created_at FROM drafts WHERE draft_id = 'd1'") .fetch_one(&state.pool) .await .unwrap(); upsert_draft(&state, &owner, "d1", &draft_named("Second")) .await .unwrap(); let summaries = list_drafts(&state, &owner).await.unwrap(); assert_eq!(summaries.len(), 1); assert_eq!(summaries[0].name.as_deref(), Some("Second")); assert_eq!(summaries[0].model_count, 1); let created_after: i64 = sqlx::query_scalar("SELECT created_at FROM drafts WHERE draft_id = 'd1'") .fetch_one(&state.pool) .await .unwrap(); assert_eq!( created_at, created_after, "created_at must be preserved on update" ); } #[tokio::test] async fn list_is_owner_scoped_and_ordered() { let state = state().await; let a = did(DID_A); upsert_draft(&state, &a, "d1", &draft_named("One")) .await .unwrap(); upsert_draft(&state, &a, "d2", &draft_named("Two")) .await .unwrap(); upsert_draft(&state, &did(DID_B), "d3", &draft_named("Other")) .await .unwrap(); let summaries = list_drafts(&state, &a).await.unwrap(); assert_eq!(summaries.len(), 2); // Most recently updated first. assert_eq!(summaries[0].draft_id, "d2"); assert_eq!(summaries[1].draft_id, "d1"); } #[tokio::test] async fn delete_is_owner_scoped() { let state = state().await; let a = did(DID_A); upsert_draft(&state, &a, "d1", &draft_named("One")) .await .unwrap(); // Another owner cannot delete it. assert!(!delete_draft(&state, &did(DID_B), "d1").await.unwrap()); assert!(get_draft(&state, &a, "d1").await.unwrap().is_some()); // The owner can. assert!(delete_draft(&state, &a, "d1").await.unwrap()); assert!(get_draft(&state, &a, "d1").await.unwrap().is_none()); // Deleting again reports no row. assert!(!delete_draft(&state, &a, "d1").await.unwrap()); } #[tokio::test] async fn promote_failure_keeps_draft() { let state = state().await; let owner = did(DID_A); upsert_draft(&state, &owner, "d1", &draft_named("Keepme")) .await .unwrap(); let result = promote_draft(&state, &owner, "d1", |_thing| async { Err(internal("simulated publish failure")) }) .await; assert!(result.is_err(), "publish failure must propagate"); // The invariant: a failed publish must not delete the draft. assert!( get_draft(&state, &owner, "d1").await.unwrap().is_some(), "draft must survive a failed publish" ); } #[tokio::test] async fn promote_invalid_draft_is_rejected_and_kept() { let state = state().await; let owner = did(DID_A); // Missing license/models → assemble fails before any publish attempt. let partial = DraftThingInput { name: Some("WIP".to_string()), ..Default::default() }; upsert_draft(&state, &owner, "d1", &partial).await.unwrap(); let mut published = false; let result = promote_draft(&state, &owner, "d1", |_thing| { published = true; async { unreachable!("publish must not run for an invalid draft") } }) .await; assert!(result.is_err()); assert!(!published, "an invalid draft must never reach publish"); assert!(get_draft(&state, &owner, "d1").await.unwrap().is_some()); } }