use std::env; use std::sync::Arc; use anyhow::{anyhow, bail, Context, Result}; use argon2::{ password_hash::{PasswordHasher, SaltString}, Argon2, }; use diesel::{dsl::exists, prelude::*, select}; use once_cell::sync::Lazy; use reqwest::{Client, Method, StatusCode}; use serde_json::json; use uuid::Uuid; use backend::{ config::AppConfig, db::{self, PgPool}, jobs::{enqueue_job, JOB_ANALYZE_DOCUMENT}, models::{DocumentAsset, DocumentAssetObject, NewUser, NewUserMembership, Tenant, User}, s3, schema::{ document_asset_objects, document_assets, documents, tenants, user_memberships, users, }, storage::{ObjectStorage, S3Storage, TenantStorage}, utils::tracing::init_tracing, }; use rand::rngs::OsRng; static QUICKWIT_INDEX_TEMPLATE: Lazy = Lazy::new(|| { json!({ "version": "0.8", "index_id": "documents", "doc_mapping": { "tokenizers": [ { "name": "substring", "type": "ngram", "min_gram": 2, "max_gram": 20, "prefix_only": false } ], "field_mappings": [ { "name": "tenant_id", "type": "text", "stored": true }, { "name": "document_id", "type": "text", "stored": true }, { "name": "version_id", "type": "text", "stored": true }, { "name": "title", "type": "text", "tokenizer": "substring", "stored": true }, { "name": "text", "type": "text", "tokenizer": "substring", "record": "position" } ] }, "search_settings": { "default_search_fields": ["title", "text"] } }) }); #[derive(Debug)] enum Command { CreateUser { username: String, password: String, }, SetPassword { username: String, password: String, }, ListUsers, DeleteUser { username: String, }, CreateTenant { slug: String, storage_root: Option, quickwit_index: Option, }, DeleteTenant { slug: String, }, AddUserToTenant { username: String, slug: String, role: Option, }, RemoveUserFromTenant { username: String, slug: String, }, ReanalyzeDocuments { slug: String, }, ListTenants, DeleteAssets(String), QuickwitCreate(String), QuickwitDelete(String), } impl Command { fn usage() -> &'static str { "Usage: admin\n\ create-user \n\ set-password \n\ list-users\n\ delete-user \n\ create-tenant [storage_root] [quickwit_index]\n\ delete-tenant \n\ add-user-to-tenant [role]\n\ remove-user-from-tenant \n\ reanalyze-documents \n\ list-tenants\n\ delete-assets \n\ quickwit-create-index \n\ quickwit-delete-index " } 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"))?, password: args.next().ok_or_else(|| anyhow!("password required"))?, }), Some("set-password") => Ok(Self::SetPassword { username: args.next().ok_or_else(|| anyhow!("username required"))?, password: args.next().ok_or_else(|| anyhow!("password 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 { slug: args.next().ok_or_else(|| anyhow!("tenant slug required"))?, storage_root: args.next(), quickwit_index: args.next(), }), Some("delete-tenant") => Ok(Self::DeleteTenant { slug: args.next().ok_or_else(|| anyhow!("tenant slug required"))?, }), Some("add-user-to-tenant") => Ok(Self::AddUserToTenant { username: args.next().ok_or_else(|| anyhow!("username required"))?, slug: args.next().ok_or_else(|| anyhow!("tenant slug required"))?, role: args.next(), }), Some("remove-user-from-tenant") => Ok(Self::RemoveUserFromTenant { username: args.next().ok_or_else(|| anyhow!("username required"))?, slug: args.next().ok_or_else(|| anyhow!("tenant slug required"))?, }), Some("reanalyze-documents") => Ok(Self::ReanalyzeDocuments { slug: args.next().ok_or_else(|| anyhow!("tenant slug required"))?, }), Some("list-tenants") => Ok(Self::ListTenants), Some("delete-assets") => Ok(Self::DeleteAssets( args.next().ok_or_else(|| anyhow!("tenant slug required"))?, )), Some("quickwit-create-index") => Ok(Self::QuickwitCreate( args.next().ok_or_else(|| anyhow!("tenant slug required"))?, )), Some("quickwit-delete-index") => Ok(Self::QuickwitDelete( args.next().ok_or_else(|| anyhow!("tenant slug required"))?, )), _ => 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, password } => create_user(&pool, &username, &password)?, Command::SetPassword { username, password } => set_password(&pool, &username, &password)?, Command::ListUsers => list_users(&pool)?, Command::DeleteUser { username } => delete_user(&pool, &username)?, Command::CreateTenant { slug, storage_root, quickwit_index, } => create_tenant(&pool, &slug, storage_root, quickwit_index)?, Command::DeleteTenant { slug } => delete_tenant(&pool, &slug)?, Command::AddUserToTenant { username, slug, role, } => add_user_to_tenant(&pool, &username, &slug, role.as_deref())?, Command::RemoveUserFromTenant { username, slug } => { remove_user_from_tenant(&pool, &username, &slug)? } Command::ReanalyzeDocuments { slug } => reanalyze_documents(&pool, &slug)?, Command::ListTenants => list_tenants(&pool)?, Command::DeleteAssets(slug) => delete_assets_for_tenant(&config, &pool, &slug).await?, Command::QuickwitCreate(slug) => { quickwit_index(&config, &pool, &slug, Method::POST).await? } Command::QuickwitDelete(slug) => { quickwit_index(&config, &pool, &slug, Method::DELETE).await? } } Ok(()) } fn create_user(pool: &PgPool, username: &str, password: &str) -> Result<()> { if username.trim().is_empty() { bail!("username must not be empty"); } if password.is_empty() { bail!("password 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 password_hash = hash_password(password)?; let new_user = NewUser { id: Uuid::new_v4(), username: username.to_string(), password_hash, }; diesel::insert_into(users::table) .values(&new_user) .execute(&mut conn)?; println!("created user '{}' (id: {})", username, new_user.id); Ok(()) } fn set_password(pool: &PgPool, username: &str, password: &str) -> Result<()> { if password.is_empty() { bail!("password must not be empty"); } let mut conn = pool.get().context("failed to get database connection")?; let password_hash = hash_password(password)?; let updated = diesel::update(users::table.filter(users::username.eq(username))) .set(users::password_hash.eq(password_hash)) .execute(&mut conn)?; if updated == 0 { bail!("user '{}' not found", username); } println!("updated password for '{}'", username); Ok(()) } fn hash_password(password: &str) -> Result { let salt = SaltString::generate(&mut OsRng); let hash = Argon2::default() .hash_password(password.as_bytes(), &salt) .map_err(|err| anyhow!(err))?; Ok(hash.to_string()) } 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<(Uuid, String, String)> = user_memberships::table .inner_join(tenants::table) .filter(user_memberships::user_id.eq(user.id)) .select((tenants::id, tenants::slug, user_memberships::role)) .order((tenants::slug.asc(), user_memberships::role.asc())) .load(&mut conn)?; if memberships.is_empty() { println!("{} ({})", user.username, user.id); } else { let details: Vec = memberships .into_iter() .map(|(_, slug, role)| format!("{}: {}", slug, role)) .collect(); println!("{} ({}) -> {}", user.username, user.id, details.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, slug: &str, storage_root_arg: Option, quickwit_index_arg: Option, ) -> Result<()> { if slug.trim().is_empty() { bail!("tenant slug must not be empty"); } let mut conn = pool.get().context("failed to get database connection")?; let exists: bool = select(exists(tenants::table.filter(tenants::slug.eq(slug)))).get_result(&mut conn)?; if exists { bail!("tenant '{}' already exists", slug); } let id = Uuid::new_v4(); let storage_root = storage_root_arg .map(|mut s| { if s.is_empty() { format!("tenants/{}/", id) } else { if !s.ends_with('/') { s.push('/'); } s } }) .unwrap_or_else(|| format!("tenants/{}/", id)); let quickwit_index = quickwit_index_arg.unwrap_or_else(|| format!("documents-{}", id)); diesel::insert_into(tenants::table) .values(( tenants::id.eq(id), tenants::slug.eq(slug), tenants::storage_root.eq(Some(storage_root.clone())), tenants::quickwit_index.eq(Some(quickwit_index.clone())), tenants::status.eq("active"), tenants::config.eq(serde_json::json!({})), )) .execute(&mut conn)?; println!( "created tenant '{}' with id {}, storage_root '{}', quickwit_index '{}'", slug, id, storage_root, quickwit_index ); Ok(()) } fn delete_tenant(pool: &PgPool, slug: &str) -> Result<()> { let mut conn = pool.get().context("failed to get database connection")?; let tenant: Tenant = tenants::table .filter(tenants::slug.eq(slug)) .first(&mut conn) .optional()? .ok_or_else(|| anyhow!("tenant '{}' not found", slug))?; 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", slug); } diesel::delete(tenants::table.filter(tenants::id.eq(tenant.id))).execute(&mut conn)?; println!("deleted tenant '{}'", slug); Ok(()) } fn add_user_to_tenant(pool: &PgPool, username: &str, slug: &str, role: Option<&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))?; let tenant: Tenant = tenants::table .filter(tenants::slug.eq(slug)) .first(&mut conn) .optional()? .ok_or_else(|| anyhow!("tenant '{}' not found", slug))?; let membership = NewUserMembership { id: Uuid::new_v4(), user_id: user.id, tenant_id: tenant.id, role: role.unwrap_or("user").to_string(), }; diesel::insert_into(user_memberships::table) .values(&membership) .on_conflict((user_memberships::user_id, user_memberships::tenant_id)) .do_update() .set(user_memberships::role.eq(&membership.role)) .execute(&mut conn)?; println!( "added user '{}' to tenant '{}' with role '{}'", username, slug, membership.role ); Ok(()) } fn remove_user_from_tenant(pool: &PgPool, username: &str, slug: &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))?; let tenant: Tenant = tenants::table .filter(tenants::slug.eq(slug)) .first(&mut conn) .optional()? .ok_or_else(|| anyhow!("tenant '{}' not found", slug))?; 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, slug); } else { println!("removed user '{}' from tenant '{}'", username, slug); } Ok(()) } fn reanalyze_documents(pool: &PgPool, slug: &str) -> Result<()> { let mut conn = pool.get().context("failed to get database connection")?; let tenant: Tenant = tenants::table .filter(tenants::slug.eq(slug)) .first(&mut conn) .optional()? .ok_or_else(|| anyhow!("tenant '{}' not found", slug))?; 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", slug); 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, slug ); 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::slug.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.slug, tenant.id); } Ok(()) } async fn delete_assets_for_tenant( config: &AppConfig, pool: &PgPool, tenant_slug: &str, ) -> 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 .filter(tenants::slug.eq(tenant_slug)) .first(&mut conn) .optional() .context("failed to load tenant")? .ok_or_else(|| anyhow!("tenant '{}' not found", tenant_slug))?; let tenant_storage = TenantStorage::new(Arc::clone(&storage), &tenant) .with_context(|| format!("missing storage root for tenant {}", tenant.slug))?; 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.slug))?; if assets.is_empty() { println!("Tenant {}: no assets", tenant.slug); return Ok(()); } println!( "Tenant {} ({}): deleting {} assets…", tenant.slug, 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.slug))?; 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.slug ); } } 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.slug))?; 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.slug))?; println!("Tenant {}: asset records deleted.", tenant.slug); Ok(()) } async fn quickwit_index( config: &AppConfig, pool: &PgPool, slug: &str, 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 .filter(tenants::slug.eq(slug)) .first(&mut conn) .optional() .context("failed to query tenants")? .ok_or_else(|| anyhow!("tenant '{}' not found", slug))?; let client = Client::new(); let index_id = format!("documents-{}", tenant.id); let base_endpoint = endpoint.trim_end_matches('/'); match method { Method::POST => { let payload = render_index_template(&index_id); let response = client .post(format!("{}/api/v1/indexes", base_endpoint)) .header("content-type", "application/json") .body(payload) .send() .await .context("failed to send create index request")?; match response.status() { status if status.is_success() => { 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.slug, index_id ); } StatusCode::CONFLICT => { let lookup = client .get(format!("{}/api/v1/indexes/{}", base_endpoint, index_id)) .send() .await .context("failed to verify existing quickwit index")?; let lookup_status = lookup.status(); if !lookup_status.is_success() { let body = lookup.text().await.unwrap_or_default(); bail!( "quickwit reported conflict but index lookup failed with status {}: {}", lookup_status, body ); } 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.slug, index_id ); } status => { let body = response.text().await.unwrap_or_default(); bail!( "quickwit create index failed with status {}: {}", status, body ); } } } 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.slug); } status => { let body = response.text().await.unwrap_or_default(); bail!( "quickwit delete index failed with status {}: {}", status, body ); } } } _ => unreachable!(), } Ok(()) } fn render_index_template(index_id: &str) -> String { let mut template = QUICKWIT_INDEX_TEMPLATE.clone(); if let Some(obj) = template.as_object_mut() { obj.insert( "index_id".to_string(), serde_json::Value::String(index_id.to_string()), ); } template.to_string() }