586 lines
19 KiB
Rust
586 lines
19 KiB
Rust
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::{
|
|
auth::password,
|
|
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,
|
|
password: String,
|
|
},
|
|
SetPassword {
|
|
username: String,
|
|
password: String,
|
|
},
|
|
ListUsers,
|
|
DeleteUser {
|
|
username: String,
|
|
},
|
|
CreateTenant {
|
|
name: String,
|
|
storage_root: Option<String>,
|
|
quickwit_index: Option<String>,
|
|
},
|
|
DeleteTenant {
|
|
name: String,
|
|
},
|
|
AddUserToTenant {
|
|
username: String,
|
|
name: String,
|
|
},
|
|
RemoveUserFromTenant {
|
|
username: String,
|
|
name: String,
|
|
},
|
|
ReanalyzeDocuments {
|
|
name: String,
|
|
},
|
|
ListTenants,
|
|
DeleteAssets(String),
|
|
QuickwitCreate(String),
|
|
QuickwitDelete(String),
|
|
}
|
|
|
|
impl Command {
|
|
fn usage() -> &'static str {
|
|
"Usage: admin\n\
|
|
create-user <username> <password>\n\
|
|
set-password <username> <password>\n\
|
|
list-users\n\
|
|
delete-user <username>\n\
|
|
create-tenant <name> [storage_root] [quickwit_index]\n\
|
|
delete-tenant <name>\n\
|
|
add-user-to-tenant <username> <name>\n\
|
|
remove-user-from-tenant <username> <name>\n\
|
|
reanalyze-documents <name>\n\
|
|
list-tenants\n\
|
|
delete-assets <name>\n\
|
|
quickwit-create-index <name>\n\
|
|
quickwit-delete-index <name>"
|
|
}
|
|
|
|
fn parse() -> Result<Self> {
|
|
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 {
|
|
name: args.next().ok_or_else(|| anyhow!("tenant name required"))?,
|
|
storage_root: args.next(),
|
|
quickwit_index: args.next(),
|
|
}),
|
|
Some("delete-tenant") => Ok(Self::DeleteTenant {
|
|
name: args.next().ok_or_else(|| anyhow!("tenant name required"))?,
|
|
}),
|
|
Some("add-user-to-tenant") => Ok(Self::AddUserToTenant {
|
|
username: args.next().ok_or_else(|| anyhow!("username required"))?,
|
|
name: args.next().ok_or_else(|| anyhow!("tenant name required"))?,
|
|
}),
|
|
Some("remove-user-from-tenant") => Ok(Self::RemoveUserFromTenant {
|
|
username: args.next().ok_or_else(|| anyhow!("username required"))?,
|
|
name: args.next().ok_or_else(|| anyhow!("tenant name required"))?,
|
|
}),
|
|
Some("reanalyze-documents") => Ok(Self::ReanalyzeDocuments {
|
|
name: args.next().ok_or_else(|| anyhow!("tenant name required"))?,
|
|
}),
|
|
Some("list-tenants") => Ok(Self::ListTenants),
|
|
Some("delete-assets") => Ok(Self::DeleteAssets(
|
|
args.next().ok_or_else(|| anyhow!("tenant name required"))?,
|
|
)),
|
|
Some("quickwit-create-index") => Ok(Self::QuickwitCreate(
|
|
args.next().ok_or_else(|| anyhow!("tenant name required"))?,
|
|
)),
|
|
Some("quickwit-delete-index") => Ok(Self::QuickwitDelete(
|
|
args.next().ok_or_else(|| anyhow!("tenant name 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 {
|
|
name,
|
|
storage_root,
|
|
quickwit_index,
|
|
} => create_tenant(&pool, &name, storage_root, quickwit_index)?,
|
|
Command::DeleteTenant { name } => delete_tenant(&pool, &name)?,
|
|
Command::AddUserToTenant { username, name } => add_user_to_tenant(&pool, &username, &name)?,
|
|
Command::RemoveUserFromTenant { username, name } => {
|
|
remove_user_from_tenant(&pool, &username, &name)?
|
|
}
|
|
Command::ReanalyzeDocuments { name } => reanalyze_documents(&pool, &name)?,
|
|
Command::ListTenants => list_tenants(&pool)?,
|
|
Command::DeleteAssets(name) => delete_assets_for_tenant(&config, &pool, &name).await?,
|
|
Command::QuickwitCreate(name) => {
|
|
quickwit_index(&config, &pool, &name, Method::POST).await?
|
|
}
|
|
Command::QuickwitDelete(name) => {
|
|
quickwit_index(&config, &pool, &name, 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 = password::hash_password(password).map_err(|err| anyhow!(err))?;
|
|
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 = password::hash_password(password).map_err(|err| anyhow!(err))?;
|
|
|
|
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 list_users(pool: &PgPool) -> Result<()> {
|
|
let mut conn = pool.get().context("failed to get database connection")?;
|
|
|
|
let users_list: Vec<User> = 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<String> = 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<String>,
|
|
quickwit_index_arg: Option<String>,
|
|
) -> 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("<none>");
|
|
let quickwit_index = tenant.quickwit_index.as_deref().unwrap_or("<none>");
|
|
|
|
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, name: &str) -> Result<()> {
|
|
let mut conn = pool.get().context("failed to get database connection")?;
|
|
|
|
let tenant: Tenant = tenants::table
|
|
.filter(tenants::name.eq(name))
|
|
.first(&mut conn)
|
|
.optional()?
|
|
.ok_or_else(|| anyhow!("tenant '{}' not found", name))?;
|
|
|
|
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", name);
|
|
}
|
|
|
|
diesel::delete(tenants::table.filter(tenants::id.eq(tenant.id))).execute(&mut conn)?;
|
|
println!("deleted tenant '{}'", name);
|
|
Ok(())
|
|
}
|
|
|
|
fn add_user_to_tenant(pool: &PgPool, username: &str, name: &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::name.eq(name))
|
|
.first(&mut conn)
|
|
.optional()?
|
|
.ok_or_else(|| anyhow!("tenant '{}' not found", name))?;
|
|
|
|
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, name);
|
|
Ok(())
|
|
}
|
|
|
|
fn remove_user_from_tenant(pool: &PgPool, username: &str, name: &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::name.eq(name))
|
|
.first(&mut conn)
|
|
.optional()?
|
|
.ok_or_else(|| anyhow!("tenant '{}' not found", name))?;
|
|
|
|
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, name);
|
|
} else {
|
|
println!("removed user '{}' from tenant '{}'", username, name);
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
fn reanalyze_documents(pool: &PgPool, name: &str) -> Result<()> {
|
|
let mut conn = pool.get().context("failed to get database connection")?;
|
|
|
|
let tenant: Tenant = tenants::table
|
|
.filter(tenants::name.eq(name))
|
|
.first(&mut conn)
|
|
.optional()?
|
|
.ok_or_else(|| anyhow!("tenant '{}' not found", name))?;
|
|
|
|
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", 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, name
|
|
);
|
|
Ok(())
|
|
}
|
|
|
|
fn list_tenants(pool: &PgPool) -> Result<()> {
|
|
let mut conn = pool.get().context("failed to get database connection")?;
|
|
let tenants: Vec<Tenant> = 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.name, tenant.id);
|
|
}
|
|
|
|
Ok(())
|
|
}
|
|
|
|
async fn delete_assets_for_tenant(
|
|
config: &AppConfig,
|
|
pool: &PgPool,
|
|
tenant_name: &str,
|
|
) -> Result<()> {
|
|
let s3_client = s3::build_client(config).await?;
|
|
let storage: Arc<dyn ObjectStorage> =
|
|
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::name.eq(tenant_name))
|
|
.first(&mut conn)
|
|
.optional()
|
|
.context("failed to load tenant")?
|
|
.ok_or_else(|| anyhow!("tenant '{}' not found", tenant_name))?;
|
|
|
|
let tenant_storage = TenantStorage::new(Arc::clone(&storage), &tenant)
|
|
.with_context(|| format!("missing storage root for tenant {}", tenant.name))?;
|
|
|
|
let assets: Vec<DocumentAsset> = 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<Uuid> = assets.iter().map(|asset| asset.id).collect();
|
|
|
|
let objects: Vec<DocumentAssetObject> = 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,
|
|
name: &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::name.eq(name))
|
|
.first(&mut conn)
|
|
.optional()
|
|
.context("failed to query tenants")?
|
|
.ok_or_else(|| anyhow!("tenant '{}' not found", name))?;
|
|
|
|
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::<Option<String>>(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(())
|
|
}
|