Something went wrong. Try again.
A lexicon-driven AppView for ATProto.
Something went wrong. Try again.
Rust
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396use axum::Json;use axum::extract::{Path, State};use axum::http::StatusCode;use serde_json::Value;
use crate::AppState;use crate::db::{adapt_sql, now_rfc3339};use crate::error::AppError;use crate::event_log::{EventLog, Severity, log_event};use crate::lexicon::{LexiconType, ParsedLexicon, ProcedureAction};
use super::auth::UserAuth;use super::permissions::Permission;use super::types::{LexiconSummary, UploadLexiconBody};
/// Broadcast the current record collection list so the Jetstream/// subscription syncs its filter.async fn notify_collections(state: &AppState) { let collections = state.lexicons.get_record_collections().await; let _ = state.collections_tx.send(collections);}
/// POST /admin/lexicons — upload (upsert) a lexicon.pub(super) async fn upload_lexicon( State(state): State<AppState>, auth: UserAuth, Json(body): Json<UploadLexiconBody>,) -> Result<(StatusCode, Json<Value>), AppError> { auth.require(Permission::LexiconsCreate).await?; let backend = state.db_backend; // Validate basic structure let lexicon_version = body .lexicon_json .get("lexicon") .and_then(|v| v.as_i64()) .ok_or_else(|| { AppError::BadRequest("lexicon JSON must have a numeric 'lexicon' field".into()) })?;
if lexicon_version != 1 { return Err(AppError::BadRequest(format!( "unsupported lexicon version: {lexicon_version}" ))); }
let id = body .lexicon_json .get("id") .and_then(|v| v.as_str()) .ok_or_else(|| AppError::BadRequest("lexicon JSON must have a string 'id' field".into()))? .to_string();
// Validate action let action = ProcedureAction::from_optional_str(body.action.as_deref()).map_err(AppError::BadRequest)?;
// Validate it parses correctly ParsedLexicon::parse( body.lexicon_json.clone(), 1, body.target_collection.clone(), action.clone(), body.script.clone(), body.index_hook.clone(), body.token_cost.map(|c| c as u32), ) .map_err(|e| AppError::BadRequest(format!("failed to parse lexicon: {e}")))?;
// Validate script if provided if let Some(ref script) = body.script { crate::lua::validate_script(script).map_err(AppError::BadRequest)?; }
// Validate index_hook if provided if let Some(ref script) = body.index_hook { crate::lua::validate_script(script).map_err(AppError::BadRequest)?; }
let action_str = action.to_optional_str(); let has_script = body.script.is_some(); let lexicon_json_str = serde_json::to_string(&body.lexicon_json).unwrap_or_default(); let now = now_rfc3339();
// Upsert into database let sql = adapt_sql( r#" INSERT INTO lexicons (id, lexicon_json, backfill, target_collection, action, script, index_hook, token_cost, source, created_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, 'manual', ?) ON CONFLICT (id) DO UPDATE SET lexicon_json = EXCLUDED.lexicon_json, backfill = EXCLUDED.backfill, target_collection = EXCLUDED.target_collection, action = EXCLUDED.action, script = EXCLUDED.script, index_hook = EXCLUDED.index_hook, token_cost = EXCLUDED.token_cost, source = 'manual', revision = lexicons.revision + 1, updated_at = ? RETURNING revision "#, backend, ); let row: (i32,) = sqlx::query_as(&sql) .bind(&id) .bind(&lexicon_json_str) .bind(if body.backfill { 1_i32 } else { 0_i32 }) .bind(&body.target_collection) .bind(action_str) .bind(&body.script) .bind(&body.index_hook) .bind(body.token_cost) .bind(&now) .bind(&now) .fetch_one(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to upsert lexicon: {e}")))?;
let revision = row.0;
// Update in-memory registry with correct revision let parsed = ParsedLexicon::parse( body.lexicon_json, revision, body.target_collection, action, body.script, body.index_hook.clone(), body.token_cost.map(|c| c as u32), ) .map_err(|e| AppError::Internal(format!("failed to re-parse lexicon: {e}")))?; let is_record = parsed.lexicon_type == LexiconType::Record; state.lexicons.upsert(parsed).await;
if is_record { notify_collections(&state).await; }
let status = if revision == 1 { StatusCode::CREATED } else { StatusCode::OK };
let event_type = if status == StatusCode::CREATED { "lexicon.created" } else { "lexicon.updated" }; log_event( &state.db, EventLog { event_type: event_type.to_string(), severity: Severity::Info, actor_did: Some(auth.did.clone()), subject: Some(id.clone()), detail: serde_json::json!({ "revision": revision, "has_script": has_script, "has_index_hook": body.index_hook.is_some(), "source": "manual", }), }, backend, ) .await;
Ok(( status, Json(serde_json::json!({ "id": id, "revision": revision, })), ))}
/// GET /admin/lexicons — list all lexicons.pub(super) async fn list_lexicons( State(state): State<AppState>, auth: UserAuth,) -> Result<Json<Vec<LexiconSummary>>, AppError> { auth.require(Permission::LexiconsRead).await?; let backend = state.db_backend; let sql = adapt_sql( "SELECT id, revision, lexicon_json, backfill, action, target_collection, script, index_hook, source, authority_did, last_fetched_at, created_at, updated_at, token_cost FROM lexicons ORDER BY id", backend, ); #[allow(clippy::type_complexity)] let rows: Vec<( String, i32, String, i32, Option<String>, Option<String>, Option<String>, Option<String>, String, Option<String>, Option<String>, String, String, Option<i32>, )> = sqlx::query_as(&sql) .fetch_all(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to list lexicons: {e}")))?;
let summaries: Vec<LexiconSummary> = rows .into_iter() .map( |( id, revision, json_str, backfill, action, target_collection, script, index_hook, source, authority_did, last_fetched_at, created_at, updated_at, token_cost, )| { let json: Value = serde_json::from_str(&json_str).unwrap_or_default(); let parsed = ParsedLexicon::parse( json, revision, None, ProcedureAction::Upsert, None, None, None, ); let lexicon_type = parsed .as_ref() .map(|p| format!("{:?}", p.lexicon_type).to_lowercase()) .unwrap_or_else(|_| "unknown".into()); let record_schema = parsed .ok() .filter(|p| p.lexicon_type == LexiconType::Record) .and_then(|p| p.record_schema);
LexiconSummary { id, revision, lexicon_type, backfill: backfill != 0, action, target_collection, has_script: script.is_some(), has_index_hook: index_hook.is_some(), source, authority_did, last_fetched_at, created_at, updated_at, record_schema, token_cost, } }, ) .collect();
Ok(Json(summaries))}
/// GET /admin/lexicons/:id — get a single lexicon.pub(super) async fn get_lexicon( State(state): State<AppState>, auth: UserAuth, Path(id): Path<String>,) -> Result<Json<Value>, AppError> { auth.require(Permission::LexiconsRead).await?; let backend = state.db_backend; let sql = adapt_sql( "SELECT id, revision, lexicon_json, backfill, action, target_collection, script, index_hook, source, authority_did, last_fetched_at, created_at, updated_at, token_cost FROM lexicons WHERE id = ?", backend, ); #[allow(clippy::type_complexity)] let row: Option<( String, i32, String, i32, Option<String>, Option<String>, Option<String>, Option<String>, String, Option<String>, Option<String>, String, String, Option<i32>, )> = sqlx::query_as(&sql) .bind(&id) .fetch_optional(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to get lexicon: {e}")))?;
let ( id, revision, lexicon_json_str, backfill, action, target_collection, script, index_hook, source, authority_did, last_fetched_at, created_at, updated_at, token_cost, ) = row.ok_or_else(|| AppError::NotFound(format!("lexicon '{id}' not found")))?;
let lexicon_json: Value = serde_json::from_str(&lexicon_json_str).unwrap_or_default();
let lexicon_type = ParsedLexicon::parse( lexicon_json.clone(), revision, None, ProcedureAction::Upsert, None, None, None, ) .map(|p| format!("{:?}", p.lexicon_type).to_lowercase()) .unwrap_or_else(|_| "unknown".into());
let has_script = script.is_some();
Ok(Json(serde_json::json!({ "id": id, "revision": revision, "lexicon_json": lexicon_json, "lexicon_type": lexicon_type, "backfill": backfill != 0, "action": action, "target_collection": target_collection, "has_script": has_script, "script": script, "has_index_hook": index_hook.is_some(), "index_hook": index_hook, "source": source, "authority_did": authority_did, "last_fetched_at": last_fetched_at, "created_at": created_at, "updated_at": updated_at, "token_cost": token_cost, })))}
/// DELETE /admin/lexicons/:id — remove a lexicon.pub(super) async fn delete_lexicon( State(state): State<AppState>, auth: UserAuth, Path(id): Path<String>,) -> Result<StatusCode, AppError> { auth.require(Permission::LexiconsDelete).await?; let backend = state.db_backend; let sql = adapt_sql("DELETE FROM lexicons WHERE id = ?", backend); let result = sqlx::query(&sql) .bind(&id) .execute(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to delete lexicon: {e}")))?;
if result.rows_affected() == 0 { return Err(AppError::NotFound(format!("lexicon '{id}' not found"))); }
state.lexicons.remove(&id).await; notify_collections(&state).await;
log_event( &state.db, EventLog { event_type: "lexicon.deleted".to_string(), severity: Severity::Info, actor_did: Some(auth.did.clone()), subject: Some(id.clone()), detail: serde_json::json!({}), }, backend, ) .await;
Ok(StatusCode::NO_CONTENT)}