Something went wrong. Try again.
A fork of @slices.network/slices forked from slices.network/slices
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310//! Actor management operations.//!//! This module handles database operations for ATProto actors (users/DIDs)//! tracked within slices, including batch insertion, querying, and filtering.
use super::client::Database;use super::types::{WhereClause, WhereCondition};use crate::errors::DatabaseError;use crate::models::Actor;
impl Database { /// Inserts multiple actors in batches with conflict resolution. /// /// Updates handle and indexed_at if an actor already exists for the /// (did, slice_uri) pair. pub async fn batch_insert_actors(&self, actors: &[Actor]) -> Result<(), DatabaseError> { if actors.is_empty() { return Ok(()); }
let mut tx = self.pool.begin().await?;
const CHUNK_SIZE: usize = 1000;
for chunk in actors.chunks(CHUNK_SIZE) { for actor in chunk { sqlx::query!( r#"INSERT INTO "actor" ("did", "handle", "slice_uri", "indexed_at") VALUES ($1, $2, $3, $4) ON CONFLICT ("did", "slice_uri") DO UPDATE SET "handle" = EXCLUDED."handle", "indexed_at" = EXCLUDED."indexed_at""#, actor.did, actor.handle, actor.slice_uri, actor.indexed_at ) .execute(&mut *tx) .await?; } }
tx.commit().await?; Ok(()) }
/// Queries actors for a slice with advanced filtering and cursor-based pagination. /// /// Supports: /// - Complex WHERE conditions (AND/OR, eq/in/contains operators) /// - Cursor-based pagination /// /// # Returns /// Tuple of (actors, next_cursor) where cursor is the last DID pub async fn get_slice_actors( &self, slice_uri: &str, limit: Option<i32>, cursor: Option<&str>, where_clause: Option<&WhereClause>, ) -> Result<(Vec<Actor>, Option<String>), DatabaseError> { let limit = limit.unwrap_or(50).min(100);
let mut where_clauses = vec![format!("slice_uri = $1")]; let mut param_count = 2;
// Build WHERE conditions for actors (handle table columns properly) let (and_conditions, or_conditions) = build_actor_where_conditions(where_clause, &mut param_count); where_clauses.extend(and_conditions);
if !or_conditions.is_empty() { let or_clause = format!("({})", or_conditions.join(" OR ")); where_clauses.push(or_clause); }
// Add cursor condition if let Some(_cursor_did) = cursor { where_clauses.push(format!("did > ${}", param_count)); param_count += 1; }
let where_sql = format!("WHERE {}", where_clauses.join(" AND "));
let query = format!( r#" SELECT did, handle, slice_uri, indexed_at FROM actor {} ORDER BY did ASC LIMIT ${} "#, where_sql, param_count );
let mut sqlx_query = sqlx::query_as::<_, Actor>(&query);
// Bind parameters in order sqlx_query = sqlx_query.bind(slice_uri);
// Bind WHERE clause parameters if let Some(clause) = where_clause { for condition in clause.conditions.values() { if let Some(eq_value) = &condition.eq { if let Some(str_val) = eq_value.as_str() { sqlx_query = sqlx_query.bind(str_val); } else { sqlx_query = sqlx_query.bind(eq_value); } } if let Some(in_values) = &condition.in_values { let str_values: Vec<String> = in_values .iter() .filter_map(|v| v.as_str().map(|s| s.to_string())) .collect(); sqlx_query = sqlx_query.bind(str_values); } if let Some(contains_value) = &condition.contains { sqlx_query = sqlx_query.bind(contains_value); } }
// Bind OR conditions if let Some(or_conditions) = &clause.or_conditions { for condition in or_conditions.values() { if let Some(eq_value) = &condition.eq { if let Some(str_val) = eq_value.as_str() { sqlx_query = sqlx_query.bind(str_val); } else { sqlx_query = sqlx_query.bind(eq_value); } } if let Some(in_values) = &condition.in_values { let str_values: Vec<String> = in_values .iter() .filter_map(|v| v.as_str().map(|s| s.to_string())) .collect(); sqlx_query = sqlx_query.bind(str_values); } if let Some(contains_value) = &condition.contains { sqlx_query = sqlx_query.bind(contains_value); } } } }
// Bind cursor parameter if let Some(cursor_did) = cursor { sqlx_query = sqlx_query.bind(cursor_did); }
// Bind limit sqlx_query = sqlx_query.bind(limit as i64);
let records = sqlx_query.fetch_all(&self.pool).await?;
let cursor = if records.len() < limit as usize { None // Last page - no more results } else { records.last().map(|actor| actor.did.clone()) };
Ok((records, cursor)) }
/// Gets all actors across all slices. /// /// # Returns /// Vector of (did, slice_uri) tuples pub async fn get_all_actors(&self) -> Result<Vec<(String, String)>, DatabaseError> { let rows = sqlx::query!( r#" SELECT did, slice_uri FROM actor "# ) .fetch_all(&self.pool) .await?;
Ok(rows .into_iter() .map(|row| (row.did, row.slice_uri)) .collect()) }
/// Checks if an actor has any records in a slice. /// /// Used before actor deletion to maintain referential integrity. pub async fn actor_has_records( &self, did: &str, slice_uri: &str, ) -> Result<bool, DatabaseError> { let count = sqlx::query!( r#" SELECT COUNT(*) as count FROM record WHERE did = $1 AND slice_uri = $2 "#, did, slice_uri ) .fetch_one(&self.pool) .await?; Ok(count.count.unwrap_or(0) > 0) }
/// Deletes an actor from a specific slice. /// /// # Returns /// Number of rows affected pub async fn delete_actor(&self, did: &str, slice_uri: &str) -> Result<u64, DatabaseError> { let result = sqlx::query!( r#" DELETE FROM actor WHERE did = $1 AND slice_uri = $2 "#, did, slice_uri ) .execute(&self.pool) .await?; Ok(result.rows_affected()) }
/// Deletes all actors for a specific slice. /// /// This is a destructive operation that removes all tracked actors /// from the specified slice. Actors will be recreated when records /// are re-indexed during sync. /// /// # Arguments /// * `slice_uri` - AT-URI of the slice to clear /// /// # Returns /// Number of actors deleted pub async fn delete_all_actors_for_slice( &self, slice_uri: &str, ) -> Result<u64, DatabaseError> { let result = sqlx::query!( r#" DELETE FROM actor WHERE slice_uri = $1 "#, slice_uri ) .execute(&self.pool) .await?; Ok(result.rows_affected()) }}
/// Builds WHERE conditions specifically for actor queries.////// Unlike the general query builder, this handles actor table columns directly/// rather than treating them as JSON paths.fn build_actor_where_conditions( where_clause: Option<&WhereClause>, param_count: &mut usize,) -> (Vec<String>, Vec<String>) { let mut where_clauses = Vec::new(); let mut or_clauses = Vec::new();
if let Some(clause) = where_clause { for (field, condition) in &clause.conditions { let field_clause = build_actor_single_condition(field, condition, param_count); if !field_clause.is_empty() { where_clauses.push(field_clause); } }
if let Some(or_conditions) = &clause.or_conditions { for (field, condition) in or_conditions { let field_clause = build_actor_single_condition(field, condition, param_count); if !field_clause.is_empty() { or_clauses.push(field_clause); } } } }
(where_clauses, or_clauses)}
/// Builds a single SQL condition clause for actor fields.fn build_actor_single_condition( field: &str, condition: &WhereCondition, param_count: &mut usize,) -> String { if let Some(_eq_value) = &condition.eq { let clause = format!("{} = ${}", field, param_count); *param_count += 1; clause } else if let Some(_in_values) = &condition.in_values { let clause = format!("{} = ANY(${})", field, param_count); *param_count += 1; clause } else if let Some(_contains_value) = &condition.contains { let clause = format!("{} ILIKE '%' || ${} || '%'", field, param_count); *param_count += 1; clause } else { String::new() }}