//! Bluesky profile import functionality. //! //! This module imports a user's Bluesky profile (`app.bsky.actor.profile`) to create //! a Smokesignal profile (`events.smokesignal.profile`) on first login. use std::sync::Arc; use atproto_client::{ client::{Auth, DPoPAuth, post_dpop_bytes_with_headers}, com::atproto::repo::{PutRecordRequest, PutRecordResponse, get_blob, get_record, put_record}, }; use atproto_identity::resolve::IdentityResolver; use atproto_record::lexicon::TypedBlob; use bytes::Bytes; use http::header::CONTENT_TYPE; use serde::Deserialize; use crate::{ atproto::lexicon::{ bluesky_profile::{BlueskyProfile, NSID as BLUESKY_PROFILE_NSID}, profile::{NSID as SMOKESIGNAL_PROFILE_NSID, Profile as SmokesignalProfile}, }, facets::{FacetLimits, parse_facets_from_text}, storage::{ StoragePool, content::ContentStorage, profile::{profile_get_by_did, profile_insert}, }, }; use super::errors::ProfileImportError; /// Maximum length for display name (Smokesignal limit) const MAX_DISPLAY_NAME_LENGTH: usize = 200; /// Maximum length for description (Smokesignal limit) const MAX_DESCRIPTION_LENGTH: usize = 5000; /// Import a user's Bluesky profile to create a Smokesignal profile. /// /// This function is called on first login to provide seamless onboarding. /// If the user already has a Smokesignal profile, this is a no-op. /// /// # Arguments /// * `http_client` - HTTP client for making requests /// * `content_storage` - Storage for avatar images /// * `pool` - Database pool /// * `identity_resolver` - Resolver for parsing facets (mentions) /// * `dpop_auth` - DPoP authentication for PDS writes /// * `pds_endpoint` - User's PDS endpoint /// * `did` - User's DID /// * `handle` - User's handle (for display_name fallback) /// * `facet_limits` - Limits for facet parsing /// /// # Returns /// * `Ok(true)` - Profile was successfully imported /// * `Ok(false)` - Import was skipped (profile exists or no Bluesky profile) /// * `Err(...)` - Import failed #[allow(clippy::too_many_arguments)] pub(crate) async fn import_bluesky_profile( http_client: &reqwest::Client, content_storage: &Arc, pool: &StoragePool, identity_resolver: &Arc, dpop_auth: &DPoPAuth, pds_endpoint: &str, did: &str, handle: &str, facet_limits: &FacetLimits, ) -> Result { // Check if user already has a Smokesignal profile if let Ok(Some(_)) = profile_get_by_did(pool, did).await { tracing::debug!(did = %did, "User already has Smokesignal profile, skipping import"); return Ok(false); } // Fetch Bluesky profile from user's PDS (public read, no auth needed) let record_response = get_record( http_client, &Auth::None, pds_endpoint, did, BLUESKY_PROFILE_NSID, "self", None, ) .await .map_err(|e| { tracing::warn!(did = %did, error = %e, "Failed to fetch Bluesky profile"); ProfileImportError::FetchFailed(e.to_string()) })?; let bluesky_profile: BlueskyProfile = match record_response { atproto_client::com::atproto::repo::GetRecordResponse::Record { value, .. } => { serde_json::from_value(value).map_err(|e| { tracing::warn!(did = %did, error = %e, "Failed to parse Bluesky profile"); ProfileImportError::ParseFailed(e.to_string()) })? } atproto_client::com::atproto::repo::GetRecordResponse::Error(e) => { if e.error.as_deref() == Some("RecordNotFound") { tracing::debug!(did = %did, "No Bluesky profile found, skipping import"); return Ok(false); } tracing::warn!(did = %did, error = ?e.error, "PDS error fetching Bluesky profile"); return Err(ProfileImportError::FetchFailed(e.error_message())); } }; // Build Smokesignal profile from Bluesky profile let mut smokesignal_profile = SmokesignalProfile { display_name: None, description: None, profile_host: Some("bsky.app".to_string()), // Default to Bluesky facets: None, avatar: None, banner: None, // Don't import banner (different aspect ratios) extra: std::collections::HashMap::new(), }; // Copy display_name (truncate if needed) if let Some(display_name) = &bluesky_profile.display_name { let truncated = if display_name.len() > MAX_DISPLAY_NAME_LENGTH { display_name[..MAX_DISPLAY_NAME_LENGTH].to_string() } else { display_name.clone() }; smokesignal_profile.display_name = Some(truncated); } // Copy description and parse facets (truncate if needed) if let Some(description) = &bluesky_profile.description { let truncated = if description.len() > MAX_DESCRIPTION_LENGTH { description[..MAX_DESCRIPTION_LENGTH].to_string() } else { description.clone() }; // Parse facets from description if let Some(facets) = parse_facets_from_text(&truncated, identity_resolver.as_ref(), facet_limits).await { smokesignal_profile.facets = Some(facets); } smokesignal_profile.description = Some(truncated); } // Import avatar if present if let Some(ref avatar_blob) = bluesky_profile.avatar { match import_avatar( http_client, content_storage, dpop_auth, pds_endpoint, did, avatar_blob, ) .await { Ok(Some(new_blob)) => { smokesignal_profile.avatar = Some(new_blob); tracing::debug!(did = %did, "Successfully imported avatar"); } Ok(None) => { tracing::debug!(did = %did, "Avatar import returned None, skipping avatar"); } Err(e) => { // Log but don't fail - profile can still be created without avatar tracing::warn!(did = %did, error = %e, "Failed to import avatar, continuing without it"); } } } // Build typed profile for PDS let typed_profile = TypedSmokesignalProfile { r#type: SMOKESIGNAL_PROFILE_NSID.to_string(), profile: smokesignal_profile.clone(), }; // Write profile to user's PDS let put_request = PutRecordRequest { repo: did.to_string(), collection: SMOKESIGNAL_PROFILE_NSID.to_string(), record_key: "self".to_string(), record: typed_profile.clone(), validate: false, swap_record: None, swap_commit: None, }; let put_response = put_record( http_client, &Auth::DPoP(dpop_auth.clone()), pds_endpoint, put_request, ) .await .map_err(|e| ProfileImportError::PdsWriteFailed(e.to_string()))?; match put_response { PutRecordResponse::StrongRef { uri, cid, .. } => { // Determine display name for local storage let display_name_for_db = smokesignal_profile .display_name .as_ref() .filter(|s| !s.trim().is_empty()) .map(|s| s.as_str()) .unwrap_or(handle); // Store profile locally if let Err(e) = profile_insert( pool, &uri, &cid, did, display_name_for_db, &smokesignal_profile, ) .await { tracing::error!(did = %did, error = %e, "Failed to store imported profile locally"); return Err(ProfileImportError::StorageFailed(e.to_string())); } tracing::info!(did = %did, uri = %uri, "Successfully imported Bluesky profile to Smokesignal"); Ok(true) } PutRecordResponse::Error(e) => { tracing::error!(did = %did, error = ?e.error, "PDS returned error for profile import"); Err(ProfileImportError::PdsWriteFailed(e.error_message())) } } } /// Import avatar from Bluesky profile. /// /// Downloads the avatar from the user's PDS, processes it through the image pipeline, /// re-uploads it to the user's PDS, and stores it locally. async fn import_avatar( http_client: &reqwest::Client, content_storage: &Arc, dpop_auth: &DPoPAuth, pds_endpoint: &str, did: &str, avatar_blob: &TypedBlob, ) -> Result, ProfileImportError> { let blob_cid = &avatar_blob.inner.ref_.link; // Validate size (max 3MB) if avatar_blob.inner.size > 3_000_000 { tracing::debug!(did = %did, size = avatar_blob.inner.size, "Avatar exceeds max size, skipping"); return Ok(None); } // Download blob from PDS (public, no auth needed) let image_bytes = get_blob(http_client, pds_endpoint, did, blob_cid) .await .map_err(|e| ProfileImportError::AvatarDownloadFailed(e.to_string()))?; // Process avatar through image pipeline (validates and converts to 400x400 PNG) let processed = crate::image::process_avatar(&image_bytes) .map_err(|e| ProfileImportError::AvatarProcessFailed(e.to_string()))?; // Upload processed avatar to user's PDS let new_blob = upload_blob_to_pds( http_client, dpop_auth, pds_endpoint, &processed, "image/png", ) .await .map_err(|e| ProfileImportError::AvatarUploadFailed(e.to_string()))?; // Store avatar locally in content storage let image_path = format!("{}.png", new_blob.inner.ref_.link); if let Err(e) = content_storage.write_content(&image_path, &processed).await { tracing::warn!( did = %did, cid = %new_blob.inner.ref_.link, error = %e, "Failed to store avatar in content storage, continuing anyway" ); // Don't fail - the PDS upload succeeded } Ok(Some(new_blob)) } /// Upload a blob to the user's PDS using DPoP authentication. async fn upload_blob_to_pds( http_client: &reqwest::Client, dpop_auth: &DPoPAuth, pds_endpoint: &str, data: &[u8], mime_type: &str, ) -> Result { let upload_url = format!("{}/xrpc/com.atproto.repo.uploadBlob", pds_endpoint); let mut headers = http::HeaderMap::default(); headers.insert(CONTENT_TYPE, mime_type.parse().unwrap()); let blob_response = post_dpop_bytes_with_headers( http_client, dpop_auth, &upload_url, Bytes::copy_from_slice(data), &headers, ) .await .map_err(|e| e.to_string())?; serde_json::from_value::(blob_response) .map(|r| r.blob) .map_err(|e| e.to_string()) } /// Response from com.atproto.repo.uploadBlob #[derive(Deserialize)] struct CreateBlobResponse { blob: TypedBlob, } /// Typed Smokesignal profile with $type field for PDS storage #[derive(Clone, serde::Serialize, serde::Deserialize)] struct TypedSmokesignalProfile { #[serde(rename = "$type")] r#type: String, #[serde(flatten)] profile: SmokesignalProfile, }