//! GraphQL schema extensions for cross-collection record queries use async_graphql::dynamic::{Field, FieldFuture, FieldValue, InputObject, InputValue, Object, TypeRef}; use async_graphql::Value as GraphQLValue; use serde_json::Value; use crate::models::{Record, WhereClause}; use crate::database::Database; /// Container for record data #[derive(Clone)] pub struct SliceRecordContainer { pub uri: String, pub cid: String, pub did: String, pub collection: String, pub value: Value, pub indexed_at: String, } /// Container for slice records response #[derive(Clone)] pub struct SliceRecordsConnection { pub records: Vec, pub total_count: i64, pub has_next_page: bool, pub end_cursor: Option, } /// Edge data structure for Relay connections #[derive(Clone)] pub struct SliceRecordEdge { pub node: SliceRecordContainer, pub cursor: String, } /// Creates the SliceRecord GraphQL type pub fn create_slice_record_type() -> Object { let mut slice_record = Object::new("SliceRecord"); slice_record = slice_record.field(Field::new("uri", TypeRef::named_nn(TypeRef::STRING), |ctx| { FieldFuture::new(async move { let container = ctx.parent_value.try_downcast_ref::()?; Ok(Some(GraphQLValue::from(container.uri.clone()))) }) })); slice_record = slice_record.field(Field::new("cid", TypeRef::named_nn(TypeRef::STRING), |ctx| { FieldFuture::new(async move { let container = ctx.parent_value.try_downcast_ref::()?; Ok(Some(GraphQLValue::from(container.cid.clone()))) }) })); slice_record = slice_record.field(Field::new("did", TypeRef::named_nn(TypeRef::STRING), |ctx| { FieldFuture::new(async move { let container = ctx.parent_value.try_downcast_ref::()?; Ok(Some(GraphQLValue::from(container.did.clone()))) }) })); slice_record = slice_record.field(Field::new("collection", TypeRef::named_nn(TypeRef::STRING), |ctx| { FieldFuture::new(async move { let container = ctx.parent_value.try_downcast_ref::()?; Ok(Some(GraphQLValue::from(container.collection.clone()))) }) })); slice_record = slice_record.field(Field::new("value", TypeRef::named_nn(TypeRef::STRING), |ctx| { FieldFuture::new(async move { let container = ctx.parent_value.try_downcast_ref::()?; let json_str = serde_json::to_string(&container.value) .unwrap_or_else(|_| "{}".to_string()); Ok(Some(GraphQLValue::from(json_str))) }) })); slice_record = slice_record.field(Field::new("indexedAt", TypeRef::named_nn(TypeRef::STRING), |ctx| { FieldFuture::new(async move { let container = ctx.parent_value.try_downcast_ref::()?; Ok(Some(GraphQLValue::from(container.indexed_at.clone()))) }) })); slice_record } /// Creates the SliceRecordEdge GraphQL type pub fn create_slice_record_edge_type() -> Object { let mut edge = Object::new("SliceRecordEdge"); edge = edge.field(Field::new("node", TypeRef::named_nn("SliceRecord"), |ctx| { FieldFuture::new(async move { let edge_data = ctx.parent_value.try_downcast_ref::()?; Ok(Some(FieldValue::owned_any(edge_data.node.clone()))) }) })); edge = edge.field(Field::new("cursor", TypeRef::named_nn(TypeRef::STRING), |ctx| { FieldFuture::new(async move { let edge_data = ctx.parent_value.try_downcast_ref::()?; Ok(Some(GraphQLValue::from(edge_data.cursor.clone()))) }) })); edge } /// Creates the SliceRecordsConnection GraphQL type pub fn create_slice_records_connection_type() -> Object { let mut connection = Object::new("SliceRecordsConnection"); // Add totalCount field connection = connection.field(Field::new("totalCount", TypeRef::named_nn(TypeRef::INT), |ctx| { FieldFuture::new(async move { let container = ctx.parent_value.try_downcast_ref::()?; Ok(Some(GraphQLValue::from(container.total_count as i32))) }) })); // Add edges field (Relay standard) connection = connection.field(Field::new("edges", TypeRef::named_nn_list_nn("SliceRecordEdge"), |ctx| { FieldFuture::new(async move { let container = ctx.parent_value.try_downcast_ref::()?; let edges: Vec> = container.records .iter() .map(|record| { let record_container = SliceRecordContainer { uri: record.uri.clone(), cid: record.cid.clone(), did: record.did.clone(), collection: record.collection.clone(), value: record.json.clone(), indexed_at: record.indexed_at.to_rfc3339(), }; let edge = SliceRecordEdge { node: record_container, cursor: record.cid.clone(), // Use CID as cursor }; FieldValue::owned_any(edge) }) .collect(); Ok(Some(FieldValue::list(edges))) }) })); // Add pageInfo field connection = connection.field(Field::new("pageInfo", TypeRef::named_nn("PageInfo"), |ctx| { FieldFuture::new(async move { let container = ctx.parent_value.try_downcast_ref::()?; let mut page_info = async_graphql::indexmap::IndexMap::new(); page_info.insert( async_graphql::Name::new("hasNextPage"), GraphQLValue::from(container.has_next_page), ); page_info.insert( async_graphql::Name::new("hasPreviousPage"), GraphQLValue::from(false), ); // Add endCursor if let Some(ref cursor) = container.end_cursor { page_info.insert( async_graphql::Name::new("endCursor"), GraphQLValue::from(cursor.clone()), ); } // Add startCursor (first record's CID if available) if let Some(first_record) = container.records.first() { page_info.insert( async_graphql::Name::new("startCursor"), GraphQLValue::from(first_record.cid.clone()), ); } Ok(Some(FieldValue::owned_any(GraphQLValue::Object(page_info)))) }) })); connection } /// Parse a where input object into a WhereClause fn parse_where_input(value: async_graphql::dynamic::ValueAccessor) -> Option { if let Ok(obj) = value.object() { let mut conditions = std::collections::HashMap::new(); let mut or_clauses: Vec = Vec::new(); // Manually check for each filter field let fields = ["collection", "did", "uri", "cid", "indexedAt", "json"]; for field_name in fields { if let Some(filter_value) = obj.get(field_name) { if let Ok(filter_obj) = filter_value.object() { let mut condition = crate::database::WhereCondition { eq: None, in_values: None, contains: None, fuzzy: None, gt: None, gte: None, lt: None, lte: None, }; if let Some(eq_val) = filter_obj.get("eq") { if let Ok(s) = eq_val.string() { condition.eq = Some(serde_json::Value::String(s.to_string())); } } if let Some(contains_val) = filter_obj.get("contains") { if let Ok(s) = contains_val.string() { condition.contains = Some(s.to_string()); } } conditions.insert(field_name.to_string(), condition); } } } // Handle OR conditions if let Some(or_value) = obj.get("or") { if let Ok(or_list) = or_value.list() { let len = or_list.len(); for i in 0..len { if let Some(or_item) = or_list.get(i) { if let Some(or_clause) = parse_where_input(or_item) { or_clauses.push(or_clause); } } } } } if !conditions.is_empty() || !or_clauses.is_empty() { Some(WhereClause { conditions, or_conditions: None, and: None, or: if or_clauses.is_empty() { None } else { Some(or_clauses) }, }) } else { None } } else { None } } /// Creates the SliceRecordsWhereInput type for filtering pub fn create_slice_records_where_input() -> InputObject { InputObject::new("SliceRecordsWhereInput") .field(InputValue::new("collection", TypeRef::named("StringFilter"))) .field(InputValue::new("did", TypeRef::named("StringFilter"))) .field(InputValue::new("uri", TypeRef::named("StringFilter"))) .field(InputValue::new("cid", TypeRef::named("StringFilter"))) .field(InputValue::new("indexedAt", TypeRef::named("DateTimeFilter"))) .field(InputValue::new("json", TypeRef::named("StringFilter"))) .field(InputValue::new("or", TypeRef::named_list("SliceRecordsWhereInput"))) } /// Container for delete slice records response #[derive(Clone)] pub struct DeleteSliceRecordsOutput { pub message: String, pub records_deleted: i64, pub actors_deleted: i64, } /// Creates the DeleteSliceRecordsOutput GraphQL type pub fn create_delete_slice_records_output_type() -> Object { let mut output = Object::new("DeleteSliceRecordsOutput"); output = output.field(Field::new("message", TypeRef::named_nn(TypeRef::STRING), |ctx| { FieldFuture::new(async move { let container = ctx.parent_value.try_downcast_ref::()?; Ok(Some(GraphQLValue::from(container.message.clone()))) }) })); output = output.field(Field::new("recordsDeleted", TypeRef::named_nn(TypeRef::INT), |ctx| { FieldFuture::new(async move { let container = ctx.parent_value.try_downcast_ref::()?; Ok(Some(GraphQLValue::from(container.records_deleted as i32))) }) })); output = output.field(Field::new("actorsDeleted", TypeRef::named_nn(TypeRef::INT), |ctx| { FieldFuture::new(async move { let container = ctx.parent_value.try_downcast_ref::()?; Ok(Some(GraphQLValue::from(container.actors_deleted as i32))) }) })); output } /// Add sliceRecords query to the Query type pub fn add_slice_records_query( query: Object, database: Database, ) -> Object { let db_for_records = database.clone(); query.field( Field::new( "sliceRecords", TypeRef::named_nn("SliceRecordsConnection"), move |ctx| { let db = db_for_records.clone(); FieldFuture::new(async move { // Get slice URI - either directly or by looking up via actorHandle + rkey let slice_uri = if let Some(uri) = ctx.args.get("sliceUri").and_then(|v| v.string().ok()) { uri.to_string() } else if let (Some(actor_handle), Some(rkey)) = ( ctx.args.get("actorHandle").and_then(|v| v.string().ok()), ctx.args.get("rkey").and_then(|v| v.string().ok()), ) { // Look up the slice by actorHandle and rkey using the dedicated database method db.get_slice_uri_by_handle_and_rkey(actor_handle, rkey) .await .map_err(|e| async_graphql::Error::new(format!("Database error: {}", e)))? .ok_or_else(|| async_graphql::Error::new("Slice not found"))? } else { return Err(async_graphql::Error::new("Either sliceUri or both actorHandle and rkey are required")); }; // Relay-style pagination arguments let limit: i32 = ctx.args.get("first") .and_then(|v| v.i64().ok()) .map(|v| v as i32) .unwrap_or(50); let cursor: Option<&str> = ctx.args.get("after") .and_then(|v| v.string().ok()); // Parse where clause if provided let where_clause: Option = ctx.args.get("where") .and_then(|v| parse_where_input(v)); // Query the database for records with default sort order let default_sort = vec![ crate::database::SortField { field: "indexed_at".to_string(), direction: "desc".to_string(), } ]; let (records, next_cursor) = db.get_slice_collections_records( &slice_uri, Some(limit), cursor, Some(&default_sort), where_clause.as_ref(), ) .await .map_err(|e| async_graphql::Error::new(format!("Failed to fetch records: {}", e)))?; // Get total count for the query let total_count = db.count_slice_collections_records( &slice_uri, where_clause.as_ref(), ) .await .map_err(|e| async_graphql::Error::new(format!("Failed to count records: {}", e)))?; // Determine if there's a next page let has_next_page = next_cursor.is_some(); let connection = SliceRecordsConnection { records, total_count, has_next_page, end_cursor: next_cursor, }; Ok(Some(FieldValue::owned_any(connection))) }) }, ) .argument(InputValue::new("sliceUri", TypeRef::named(TypeRef::STRING))) .argument(InputValue::new("actorHandle", TypeRef::named(TypeRef::STRING))) .argument(InputValue::new("rkey", TypeRef::named(TypeRef::STRING))) .argument(InputValue::new("first", TypeRef::named(TypeRef::INT))) .argument(InputValue::new("after", TypeRef::named(TypeRef::STRING))) .argument(InputValue::new("where", TypeRef::named("SliceRecordsWhereInput"))) .description("Query records across all collections in a slice with filtering and pagination. Provide either sliceUri or both actorHandle and rkey.") ) } /// Add deleteSliceRecords mutation to the Mutation type pub fn add_delete_slice_records_mutation( mutation: Object, slice_uri: String, ) -> Object { mutation.field( Field::new( "deleteSliceRecords", TypeRef::named_nn("DeleteSliceRecordsOutput"), move |ctx| { let current_slice = slice_uri.clone(); FieldFuture::new(async move { // Get user_did from context (set by auth middleware) let user_did = ctx .data::() .map_err(|_| async_graphql::Error::new("Authentication required"))? .clone(); // Get slice parameter (defaults to current slice) let slice: String = match ctx.args.get("slice") { Some(val) => val.string()?.to_string(), None => current_slice, }; // Verify user owns this slice if !slice.starts_with(&format!("at://{}/", user_did)) { return Err(async_graphql::Error::new( "You do not have permission to clear this slice" )); } // Get pool from GraphQL context let pool = ctx.data::() .map_err(|_| async_graphql::Error::new("Database pool not found in context"))?; // Create Database instance from pool let db = Database::new(pool.clone()); // Delete all records for this slice let records_deleted = db .delete_all_records_for_slice(&slice) .await .map_err(|e| async_graphql::Error::new(format!("Failed to delete records: {}", e)))?; // Delete all actors for this slice let actors_deleted = db .delete_all_actors_for_slice(&slice) .await .map_err(|e| async_graphql::Error::new(format!("Failed to delete actors: {}", e)))?; let output = DeleteSliceRecordsOutput { message: format!( "Slice index cleared successfully. Deleted {} records and {} actors.", records_deleted, actors_deleted ), records_deleted: records_deleted as i64, actors_deleted: actors_deleted as i64, }; Ok(Some(FieldValue::owned_any(output))) }) }, ) .argument(InputValue::new("slice", TypeRef::named(TypeRef::STRING))) .description("Delete all records and actors from a slice index. Requires authentication and slice ownership.") ) }