use std::env; use std::sync::Arc; use anyhow::{anyhow, bail, Context, Result}; use diesel::{dsl::exists, prelude::*, select}; use reqwest::{Client, Method, StatusCode}; use uuid::Uuid; use backend::{ config::AppConfig, db::{self, PgPool}, documents::search::ensure_quickwit_index, jobs::{enqueue_job, JOB_ANALYZE_DOCUMENT}, models::{ DocumentAsset, DocumentAssetObject, NewUser, NewUserMembership, Tenant, TenantStatus, User, }, s3, schema::{ document_asset_objects, document_assets, documents, tenants, user_memberships, users, }, storage::{ObjectStorage, S3Storage, TenantStorage}, tenants::TenantService, utils::tracing::init_tracing, }; #[derive(Debug)] enum Command { CreateUser { username: String, }, ListUsers, DeleteUser { username: String, }, CreateTenant { name: String, storage_root: Option, quickwit_index: Option, }, DeleteTenant { tenant_id: Uuid, }, AddUserToTenant { username: String, tenant_id: Uuid, }, RemoveUserFromTenant { username: String, tenant_id: Uuid, }, ReanalyzeDocuments { tenant_id: Uuid, }, ListTenants, DeleteAssets(Uuid), QuickwitCreate(Uuid), QuickwitDelete(Uuid), } impl Command { fn usage() -> &'static str { "Usage: admin\n\ create-user \n\ list-users\n\ delete-user \n\ create-tenant [storage_root] [quickwit_index]\n\ delete-tenant \n\ add-user-to-tenant \n\ remove-user-from-tenant \n\ reanalyze-documents \n\ list-tenants\n\ delete-assets \n\ quickwit-create-index \n\ quickwit-delete-index " } fn parse_tenant_id(arg: Option) -> Result { let raw = arg.ok_or_else(|| anyhow!("tenant id required"))?; Uuid::parse_str(&raw).map_err(|_| anyhow!("invalid tenant id: {}", raw)) } fn parse() -> Result { let mut args = env::args().skip(1); match args.next().as_deref() { Some("create-user") => Ok(Self::CreateUser { username: args.next().ok_or_else(|| anyhow!("username required"))?, }), Some("list-users") => Ok(Self::ListUsers), Some("delete-user") => Ok(Self::DeleteUser { username: args.next().ok_or_else(|| anyhow!("username required"))?, }), Some("create-tenant") => Ok(Self::CreateTenant { name: args.next().ok_or_else(|| anyhow!("tenant name required"))?, storage_root: args.next(), quickwit_index: args.next(), }), Some("delete-tenant") => Ok(Self::DeleteTenant { tenant_id: Self::parse_tenant_id(args.next())?, }), Some("add-user-to-tenant") => Ok(Self::AddUserToTenant { username: args.next().ok_or_else(|| anyhow!("username required"))?, tenant_id: Self::parse_tenant_id(args.next())?, }), Some("remove-user-from-tenant") => Ok(Self::RemoveUserFromTenant { username: args.next().ok_or_else(|| anyhow!("username required"))?, tenant_id: Self::parse_tenant_id(args.next())?, }), Some("reanalyze-documents") => Ok(Self::ReanalyzeDocuments { tenant_id: Self::parse_tenant_id(args.next())?, }), Some("list-tenants") => Ok(Self::ListTenants), Some("delete-assets") => Ok(Self::DeleteAssets(Self::parse_tenant_id(args.next())?)), Some("quickwit-create-index") => { Ok(Self::QuickwitCreate(Self::parse_tenant_id(args.next())?)) } Some("quickwit-delete-index") => { Ok(Self::QuickwitDelete(Self::parse_tenant_id(args.next())?)) } _ => Err(anyhow!(Self::usage())), } } } #[tokio::main] async fn main() -> Result<()> { init_tracing("info"); let command = Command::parse()?; let config = AppConfig::load_and_log("admin")?; let pool = db::init_pool_with_size(&config.database_url, config.database_max_pool_size)?; match command { Command::CreateUser { username } => create_user(&pool, &username)?, Command::ListUsers => list_users(&pool)?, Command::DeleteUser { username } => delete_user(&pool, &username)?, Command::CreateTenant { name, storage_root, quickwit_index, } => create_tenant(&pool, &name, storage_root, quickwit_index)?, Command::DeleteTenant { tenant_id } => delete_tenant(&pool, tenant_id)?, Command::AddUserToTenant { username, tenant_id, } => add_user_to_tenant(&pool, &username, tenant_id)?, Command::RemoveUserFromTenant { username, tenant_id, } => remove_user_from_tenant(&pool, &username, tenant_id)?, Command::ReanalyzeDocuments { tenant_id } => reanalyze_documents(&pool, tenant_id)?, Command::ListTenants => list_tenants(&pool)?, Command::DeleteAssets(tenant_id) => { delete_assets_for_tenant(&config, &pool, tenant_id).await? } Command::QuickwitCreate(tenant_id) => { quickwit_index(&config, &pool, tenant_id, Method::POST).await? } Command::QuickwitDelete(tenant_id) => { quickwit_index(&config, &pool, tenant_id, Method::DELETE).await? } } Ok(()) } fn create_user(pool: &PgPool, username: &str) -> Result<()> { if username.trim().is_empty() { bail!("username must not be empty"); } let mut conn = pool.get().context("failed to get database connection")?; let exists: bool = select(exists(users::table.filter(users::username.eq(username)))).get_result(&mut conn)?; if exists { bail!("user '{}' already exists", username); } let new_user = NewUser { id: Uuid::new_v4(), username: username.to_string(), }; diesel::insert_into(users::table) .values(&new_user) .execute(&mut conn)?; println!("created user '{}' (id: {})", username, new_user.id); Ok(()) } fn list_users(pool: &PgPool) -> Result<()> { let mut conn = pool.get().context("failed to get database connection")?; let users_list: Vec = users::table.order(users::username.asc()).load(&mut conn)?; if users_list.is_empty() { println!("No users found."); return Ok(()); } for user in users_list { let memberships: Vec = user_memberships::table .inner_join(tenants::table) .filter(user_memberships::user_id.eq(user.id)) .select(tenants::name) .order(tenants::name.asc()) .load(&mut conn)?; if memberships.is_empty() { println!("{} ({})", user.username, user.id); } else { println!( "{} ({}) -> {}", user.username, user.id, memberships.join(", ") ); } } Ok(()) } fn delete_user(pool: &PgPool, username: &str) -> Result<()> { let mut conn = pool.get().context("failed to get database connection")?; let user: User = users::table .filter(users::username.eq(username)) .first(&mut conn) .optional()? .ok_or_else(|| anyhow!("user '{}' not found", username))?; diesel::delete(user_memberships::table.filter(user_memberships::user_id.eq(user.id))) .execute(&mut conn)?; diesel::delete(users::table.filter(users::id.eq(user.id))).execute(&mut conn)?; println!("deleted user '{}'", username); Ok(()) } fn create_tenant( pool: &PgPool, name: &str, storage_root_arg: Option, quickwit_index_arg: Option, ) -> Result<()> { let service = TenantService::new(pool.clone()); let tenant = service .create_tenant( name, storage_root_arg.as_deref(), quickwit_index_arg.as_deref(), TenantStatus::Creating, &[], None, ) .map_err(|err| anyhow!(format!("{err:?}")))?; let storage_root = tenant.storage_root.as_deref().unwrap_or(""); let quickwit_index = tenant.quickwit_index.as_deref().unwrap_or(""); println!( "created tenant '{}' with id {}, storage_root '{}', quickwit_index '{}', status '{}'", tenant.name, tenant.id, storage_root, quickwit_index, tenant.status.as_str() ); Ok(()) } fn delete_tenant(pool: &PgPool, tenant_id: Uuid) -> Result<()> { let mut conn = pool.get().context("failed to get database connection")?; let tenant: Tenant = tenants::table .find(tenant_id) .first(&mut conn) .optional()? .ok_or_else(|| anyhow!("tenant '{}' not found", tenant_id))?; let member_exists: bool = select(exists( user_memberships::table.filter(user_memberships::tenant_id.eq(tenant.id)), )) .get_result(&mut conn)?; if member_exists { bail!("tenant '{}' still has user memberships", tenant.name); } diesel::delete(tenants::table.filter(tenants::id.eq(tenant.id))).execute(&mut conn)?; println!("deleted tenant '{}'", tenant.name); Ok(()) } fn add_user_to_tenant(pool: &PgPool, username: &str, tenant_id: Uuid) -> Result<()> { let mut conn = pool.get().context("failed to get database connection")?; let user: User = users::table .filter(users::username.eq(username)) .first(&mut conn) .optional()? .ok_or_else(|| anyhow!("user '{}' not found", username))?; let tenant: Tenant = tenants::table .find(tenant_id) .first(&mut conn) .optional()? .ok_or_else(|| anyhow!("tenant '{}' not found", tenant_id))?; let membership = NewUserMembership { id: Uuid::new_v4(), user_id: user.id, tenant_id: tenant.id, }; diesel::insert_into(user_memberships::table) .values(&membership) .on_conflict((user_memberships::user_id, user_memberships::tenant_id)) .do_nothing() .execute(&mut conn)?; println!("added user '{}' to tenant '{}'", username, tenant.name); Ok(()) } fn remove_user_from_tenant(pool: &PgPool, username: &str, tenant_id: Uuid) -> Result<()> { let mut conn = pool.get().context("failed to get database connection")?; let user: User = users::table .filter(users::username.eq(username)) .first(&mut conn) .optional()? .ok_or_else(|| anyhow!("user '{}' not found", username))?; let tenant: Tenant = tenants::table .find(tenant_id) .first(&mut conn) .optional()? .ok_or_else(|| anyhow!("tenant '{}' not found", tenant_id))?; let removed = diesel::delete( user_memberships::table .filter(user_memberships::user_id.eq(user.id)) .filter(user_memberships::tenant_id.eq(tenant.id)), ) .execute(&mut conn)?; if removed == 0 { println!( "user '{}' was not a member of tenant '{}'", username, tenant.name ); } else { println!("removed user '{}' from tenant '{}'", username, tenant.name); } Ok(()) } fn reanalyze_documents(pool: &PgPool, tenant_id: Uuid) -> Result<()> { let mut conn = pool.get().context("failed to get database connection")?; let tenant: Tenant = tenants::table .find(tenant_id) .first(&mut conn) .optional()? .ok_or_else(|| anyhow!("tenant '{}' not found", tenant_id))?; let targets: Vec<(Uuid, Uuid)> = documents::table .filter(documents::tenant_id.eq(tenant.id)) .filter(documents::deleted_at.is_null()) .select((documents::id, documents::current_version_id)) .load(&mut conn)?; if targets.is_empty() { println!("tenant '{}' has no active documents", tenant.name); return Ok(()); } let mut queued = 0usize; for (document_id, version_id) in targets { enqueue_job( &mut conn, tenant.id, JOB_ANALYZE_DOCUMENT, serde_json::json!({ "document_id": document_id, "document_version_id": version_id, "force": true, }), None, ) .map_err(|err| anyhow!("failed to enqueue analyze job: {}", err))?; queued += 1; } println!( "queued {} documents for re-analysis in tenant '{}'", queued, tenant.name ); Ok(()) } fn list_tenants(pool: &PgPool) -> Result<()> { let mut conn = pool.get().context("failed to get database connection")?; let tenants: Vec = tenants::table .order(tenants::name.asc()) .load(&mut conn) .context("failed to load tenants")?; if tenants.is_empty() { println!("No tenants found."); return Ok(()); } for tenant in tenants { println!("{} {}", tenant.id, tenant.name); } Ok(()) } async fn delete_assets_for_tenant( config: &AppConfig, pool: &PgPool, tenant_id: Uuid, ) -> Result<()> { let s3_client = s3::build_client(config).await?; let storage: Arc = Arc::new(S3Storage::new(s3_client, config.s3_bucket.clone())); let mut conn = pool.get().context("failed to get database connection")?; let tenant: Tenant = tenants::table .find(tenant_id) .first(&mut conn) .optional()? .ok_or_else(|| anyhow!("tenant '{}' not found", tenant_id))?; let tenant_storage = TenantStorage::new(Arc::clone(&storage), &tenant) .with_context(|| format!("missing storage root for tenant {}", tenant.name))?; let assets: Vec = document_assets::table .filter(document_assets::tenant_id.eq(tenant.id)) .load(&mut conn) .with_context(|| format!("failed to load assets for tenant {}", tenant.name))?; if assets.is_empty() { println!("Tenant {}: no assets", tenant.name); return Ok(()); } println!( "Tenant {} ({}): deleting {} assets…", tenant.name, tenant.id, assets.len() ); let asset_ids: Vec = assets.iter().map(|asset| asset.id).collect(); let objects: Vec = document_asset_objects::table .filter(document_asset_objects::tenant_id.eq(tenant.id)) .filter(document_asset_objects::asset_id.eq_any(&asset_ids)) .load(&mut conn) .with_context(|| format!("failed to load asset objects for tenant {}", tenant.name))?; for object in &objects { if let Err(err) = tenant_storage.delete_object(&object.s3_key).await { eprintln!( "Failed to delete object {} (tenant {}): {err}", object.s3_key, tenant.name ); } } diesel::delete( document_asset_objects::table .filter(document_asset_objects::tenant_id.eq(tenant.id)) .filter(document_asset_objects::asset_id.eq_any(&asset_ids)), ) .execute(&mut conn) .with_context(|| format!("failed to remove asset objects for tenant {}", tenant.name))?; diesel::delete(document_assets::table.filter(document_assets::tenant_id.eq(tenant.id))) .execute(&mut conn) .with_context(|| format!("failed to remove asset records for tenant {}", tenant.name))?; println!("Tenant {}: asset records deleted.", tenant.name); Ok(()) } async fn quickwit_index( config: &AppConfig, pool: &PgPool, tenant_id: Uuid, method: Method, ) -> Result<()> { let endpoint = config .quickwit_endpoint .as_ref() .ok_or_else(|| anyhow!("quickwit endpoint not configured"))?; let mut conn = pool.get().context("failed to get database connection")?; let tenant: Tenant = tenants::table .find(tenant_id) .first(&mut conn) .optional()? .ok_or_else(|| anyhow!("tenant '{}' not found", tenant_id))?; let client = Client::new(); let index_id = format!("documents-{}", tenant.id); let base_endpoint = endpoint.trim_end_matches('/'); match method { Method::POST => { ensure_quickwit_index(&client, base_endpoint, &index_id) .await .context("failed to ensure quickwit index")?; diesel::update(tenants::table.filter(tenants::id.eq(tenant.id))) .set(tenants::quickwit_index.eq(Some(index_id.clone()))) .execute(&mut conn) .context("failed to update tenant quickwit_index")?; println!( "Tenant '{}' quickwit index set to '{}'.", tenant.name, index_id ); } Method::DELETE => { let response = client .delete(format!("{}/api/v1/indexes/{}", base_endpoint, index_id)) .send() .await .context("failed to send delete index request")?; match response.status() { status if status.is_success() || status == StatusCode::NOT_FOUND => { diesel::update(tenants::table.filter(tenants::id.eq(tenant.id))) .set(tenants::quickwit_index.eq::>(None)) .execute(&mut conn) .context("failed to clear tenant quickwit_index")?; println!("Tenant '{}' quickwit index cleared.", tenant.name); } status => { let body = response.text().await.unwrap_or_default(); bail!( "quickwit delete index failed with status {}: {}", status, body ); } } } _ => unreachable!(), } Ok(()) }