Something went wrong. Try again.
A lexicon-driven AppView for ATProto.
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844use serde::{Deserialize, Serialize};use serde_json::Value;use std::collections::HashMap;use std::sync::Arc;use tokio::sync::RwLock;use tracing::{info, warn};
/// The type of a lexicon's `main` definition.#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]#[serde(rename_all = "camelCase")]pub enum LexiconType { Record, Query, Procedure, /// Lexicons with no `main` def or a non-endpoint main type (token, object, string, etc.). Definitions,}
/// The action a procedure lexicon performs on its target collection.#[derive(Debug, Clone, PartialEq, Eq)]pub enum ProcedureAction { Create, Update, Delete, /// Backwards-compatible default: sniff for `uri` in input to decide create vs put. Upsert,}
impl ProcedureAction { /// Parse an optional action string into a `ProcedureAction`. /// Returns `Upsert` for `None`, or an error for unrecognized values. pub fn from_optional_str(s: Option<&str>) -> Result<Self, String> { match s { None => Ok(Self::Upsert), Some("create") => Ok(Self::Create), Some("update") => Ok(Self::Update), Some("delete") => Ok(Self::Delete), Some("upsert") => Ok(Self::Upsert), Some(other) => Err(format!( "invalid action '{other}': must be create, update, delete, or upsert" )), } }
/// Convert to an optional string for database storage. /// `Upsert` maps to `None` (the default). pub fn to_optional_str(&self) -> Option<&'static str> { match self { Self::Create => Some("create"), Self::Update => Some("update"), Self::Delete => Some("delete"), Self::Upsert => None, } }}
/// Metadata extracted from a raw lexicon JSON document.#[derive(Debug, Clone)]pub struct ParsedLexicon { /// The NSID, e.g. "games.gamesgamesgamesgames.game". pub id: String, /// What kind of endpoint this lexicon defines. pub lexicon_type: LexiconType, /// For records: the `key` field (e.g. "tid", "any", "literal:self"). pub record_key: Option<String>, /// For queries: the parameters schema from `defs.main.parameters`. pub parameters: Option<Value>, /// For procedures: the input schema from `defs.main.input`. pub input: Option<Value>, /// For queries and procedures: the output schema from `defs.main.output`. pub output: Option<Value>, /// For records: the record schema from `defs.main.record`. pub record_schema: Option<Value>, /// The entire raw lexicon JSON. pub raw: Value, /// Database revision number. pub revision: i32, /// For queries/procedures: the backing record collection NSID. pub target_collection: Option<String>, /// For procedures: the action this procedure performs (create, update, delete, upsert). pub action: ProcedureAction, /// Optional Lua script that replaces the built-in handler. pub script: Option<String>, /// Optional Lua script that runs when a record in this collection is indexed. pub index_hook: Option<String>, /// Optional per-NSID token cost for rate limiting. pub token_cost: Option<u32>, /// Optional space type NSID indicating this lexicon is designed for use within spaces of that type. pub space_type: Option<String>,}
impl ParsedLexicon { /// Parse a lexicon JSON document into a `ParsedLexicon`. pub fn parse( raw: Value, revision: i32, target_collection: Option<String>, action: ProcedureAction, script: Option<String>, index_hook: Option<String>, token_cost: Option<u32>, ) -> Result<Self, String> { let id = raw .get("id") .and_then(|v| v.as_str()) .ok_or("lexicon JSON missing 'id' field")? .to_string();
let main_def = raw.get("defs").and_then(|d| d.get("main"));
let main_type_str = main_def .and_then(|m| m.get("type")) .and_then(|t| t.as_str());
let lexicon_type = match main_type_str { Some("record") => LexiconType::Record, Some("query") => LexiconType::Query, Some("procedure") => LexiconType::Procedure, _ => LexiconType::Definitions, };
let record_key = main_def .and_then(|m| m.get("key")) .and_then(|k| k.as_str()) .map(|s| s.to_string());
let parameters = main_def.and_then(|m| m.get("parameters")).cloned(); let input = main_def.and_then(|m| m.get("input")).cloned(); let output = main_def.and_then(|m| m.get("output")).cloned(); let record_schema = main_def.and_then(|m| m.get("record")).cloned();
let space_type = raw .get("spaceType") .and_then(|v| v.as_str()) .map(|s| s.to_string());
Ok(Self { id, lexicon_type, record_key, parameters, input, output, record_schema, raw, revision, target_collection, action, script, index_hook, token_cost, space_type, }) }}
/// In-memory cache of all parsed lexicons, keyed by NSID.#[derive(Debug, Clone)]pub struct LexiconRegistry { inner: Arc<RwLock<HashMap<String, ParsedLexicon>>>,}
impl Default for LexiconRegistry { fn default() -> Self { Self::new() }}
impl LexiconRegistry { pub fn new() -> Self { Self { inner: Arc::new(RwLock::new(HashMap::new())), } }
/// Load all lexicons from the database, replacing any existing entries. pub async fn load_from_db(&self, db: &sqlx::AnyPool) -> Result<(), String> { #[allow(clippy::type_complexity)] let rows: Vec<( String, String, i32, Option<String>, Option<String>, Option<String>, Option<String>, Option<i32>, )> = sqlx::query_as( "SELECT id, lexicon_json, revision, target_collection, action, script, index_hook, token_cost FROM lexicons", ) .fetch_all(db) .await .map_err(|e| format!("failed to load lexicons: {e}"))?;
let mut inner = self.inner.write().await; inner.clear();
let mut loaded = 0u32; for ( id, json_str, revision, target_collection, action_str, script, index_hook, token_cost, ) in rows { let json: Value = match serde_json::from_str(&json_str) { Ok(v) => v, Err(e) => { warn!(%id, "failed to parse lexicon_json: {e}"); continue; } }; let action = match ProcedureAction::from_optional_str(action_str.as_deref()) { Ok(a) => a, Err(e) => { warn!(%id, "invalid action value: {e}"); ProcedureAction::Upsert } }; match ParsedLexicon::parse( json, revision, target_collection, action, script, index_hook, token_cost.map(|c| c as u32), ) { Ok(parsed) => { inner.insert(id, parsed); loaded += 1; } Err(e) => { warn!(%id, "failed to parse lexicon: {e}"); } } }
info!(count = loaded, "loaded lexicons into registry"); Ok(()) }
/// Insert or update a single lexicon in the registry. pub async fn upsert(&self, parsed: ParsedLexicon) { let mut inner = self.inner.write().await; inner.insert(parsed.id.clone(), parsed); }
/// Remove a lexicon from the registry by NSID. pub async fn remove(&self, id: &str) -> bool { let mut inner = self.inner.write().await; inner.remove(id).is_some() }
/// Get a single lexicon by NSID. pub async fn get(&self, id: &str) -> Option<ParsedLexicon> { let inner = self.inner.read().await; inner.get(id).cloned() }
/// Return NSIDs of all record-type lexicons. pub async fn get_record_collections(&self) -> Vec<String> { let inner = self.inner.read().await; inner .values() .filter(|lex| lex.lexicon_type == LexiconType::Record) .map(|lex| lex.id.clone()) .collect() }
/// Return NSIDs of all query-type lexicons. pub async fn get_queries(&self) -> Vec<String> { let inner = self.inner.read().await; inner .values() .filter(|lex| lex.lexicon_type == LexiconType::Query) .map(|lex| lex.id.clone()) .collect() }
/// Return NSIDs of all procedure-type lexicons. pub async fn get_procedures(&self) -> Vec<String> { let inner = self.inner.read().await; inner .values() .filter(|lex| lex.lexicon_type == LexiconType::Procedure) .map(|lex| lex.id.clone()) .collect() }
/// Get the index_hook for a record-type lexicon by its collection NSID. pub async fn get_index_hook(&self, collection: &str) -> Option<String> { let inner = self.inner.read().await; inner.get(collection).and_then(|lex| lex.index_hook.clone()) }
/// Return the total count of registered lexicons. pub async fn count(&self) -> usize { let inner = self.inner.read().await; inner.len() }}
#[cfg(test)]mod tests { use super::*; use serde_json::json;
// ----------------------------------------------------------------------- // ParsedLexicon::parse // -----------------------------------------------------------------------
fn record_lexicon_json() -> Value { json!({ "lexicon": 1, "id": "games.gamesgamesgamesgames.game", "defs": { "main": { "type": "record", "key": "tid", "record": { "type": "object", "properties": { "title": { "type": "string" } } } } } }) }
fn query_lexicon_json() -> Value { json!({ "lexicon": 1, "id": "games.gamesgamesgamesgames.listGames", "defs": { "main": { "type": "query", "parameters": { "type": "params", "properties": { "limit": { "type": "integer" } } }, "output": { "encoding": "application/json" } } } }) }
fn procedure_lexicon_json() -> Value { json!({ "lexicon": 1, "id": "games.gamesgamesgamesgames.createGame", "defs": { "main": { "type": "procedure", "input": { "encoding": "application/json" }, "output": { "encoding": "application/json" } } } }) }
fn definitions_lexicon_json() -> Value { json!({ "lexicon": 1, "id": "games.gamesgamesgamesgames.defs", "defs": { "genre": { "type": "string", "knownValues": ["action", "rpg"] } } }) }
#[test] fn parse_record_lexicon() { let parsed = ParsedLexicon::parse( record_lexicon_json(), 1, None, ProcedureAction::Upsert, None, None, None, ) .unwrap(); assert_eq!(parsed.id, "games.gamesgamesgamesgames.game"); assert_eq!(parsed.lexicon_type, LexiconType::Record); assert_eq!(parsed.record_key, Some("tid".into())); assert!(parsed.record_schema.is_some()); assert!(parsed.parameters.is_none()); assert!(parsed.input.is_none()); }
#[test] fn parse_query_lexicon() { let parsed = ParsedLexicon::parse( query_lexicon_json(), 2, Some("games.gamesgamesgamesgames.game".into()), ProcedureAction::Upsert, None, None, None, ) .unwrap(); assert_eq!(parsed.lexicon_type, LexiconType::Query); assert!(parsed.parameters.is_some()); assert!(parsed.output.is_some()); assert_eq!( parsed.target_collection, Some("games.gamesgamesgamesgames.game".into()) ); assert_eq!(parsed.revision, 2); }
#[test] fn parse_procedure_lexicon() { let parsed = ParsedLexicon::parse( procedure_lexicon_json(), 1, None, ProcedureAction::Upsert, None, None, None, ) .unwrap(); assert_eq!(parsed.lexicon_type, LexiconType::Procedure); assert!(parsed.input.is_some()); assert!(parsed.output.is_some()); }
#[test] fn parse_procedure_with_action() { let parsed = ParsedLexicon::parse( procedure_lexicon_json(), 1, None, ProcedureAction::Delete, None, None, None, ) .unwrap(); assert_eq!(parsed.action, ProcedureAction::Delete); }
#[test] fn parse_definitions_lexicon() { let parsed = ParsedLexicon::parse( definitions_lexicon_json(), 1, None, ProcedureAction::Upsert, None, None, None, ) .unwrap(); assert_eq!(parsed.lexicon_type, LexiconType::Definitions); }
#[test] fn parse_missing_id_returns_error() { let raw = json!({"lexicon": 1, "defs": {}}); let result = ParsedLexicon::parse(raw, 1, None, ProcedureAction::Upsert, None, None, None); assert!(result.is_err()); assert!(result.unwrap_err().contains("id")); }
#[test] fn parse_preserves_raw_json() { let raw = record_lexicon_json(); let parsed = ParsedLexicon::parse( raw.clone(), 1, None, ProcedureAction::Upsert, None, None, None, ) .unwrap(); assert_eq!(parsed.raw, raw); }
#[test] fn parse_target_collection_passthrough() { let parsed = ParsedLexicon::parse( query_lexicon_json(), 1, Some("custom.collection".into()), ProcedureAction::Upsert, None, None, None, ) .unwrap(); assert_eq!(parsed.target_collection, Some("custom.collection".into())); }
// ----------------------------------------------------------------------- // LexiconRegistry // -----------------------------------------------------------------------
#[tokio::test] async fn registry_new_is_empty() { let reg = LexiconRegistry::new(); assert_eq!(reg.count().await, 0); }
#[tokio::test] async fn registry_upsert_and_get() { let reg = LexiconRegistry::new(); let parsed = ParsedLexicon::parse( record_lexicon_json(), 1, None, ProcedureAction::Upsert, None, None, None, ) .unwrap(); reg.upsert(parsed).await;
let got = reg.get("games.gamesgamesgamesgames.game").await; assert!(got.is_some()); assert_eq!(got.unwrap().lexicon_type, LexiconType::Record); }
#[tokio::test] async fn registry_upsert_replaces() { let reg = LexiconRegistry::new(); let v1 = ParsedLexicon::parse( record_lexicon_json(), 1, None, ProcedureAction::Upsert, None, None, None, ) .unwrap(); reg.upsert(v1).await;
let v2 = ParsedLexicon::parse( record_lexicon_json(), 5, None, ProcedureAction::Upsert, None, None, None, ) .unwrap(); reg.upsert(v2).await;
assert_eq!(reg.count().await, 1); assert_eq!( reg.get("games.gamesgamesgamesgames.game") .await .unwrap() .revision, 5 ); }
#[tokio::test] async fn registry_remove_existing() { let reg = LexiconRegistry::new(); let parsed = ParsedLexicon::parse( record_lexicon_json(), 1, None, ProcedureAction::Upsert, None, None, None, ) .unwrap(); reg.upsert(parsed).await;
assert!(reg.remove("games.gamesgamesgamesgames.game").await); assert_eq!(reg.count().await, 0); }
#[tokio::test] async fn registry_remove_nonexistent() { let reg = LexiconRegistry::new(); assert!(!reg.remove("nonexistent").await); }
#[tokio::test] async fn registry_get_nonexistent() { let reg = LexiconRegistry::new(); assert!(reg.get("nonexistent").await.is_none()); }
#[tokio::test] async fn registry_type_filtered_collections() { let reg = LexiconRegistry::new();
let record = ParsedLexicon::parse( record_lexicon_json(), 1, None, ProcedureAction::Upsert, None, None, None, ) .unwrap(); let query = ParsedLexicon::parse( query_lexicon_json(), 1, None, ProcedureAction::Upsert, None, None, None, ) .unwrap(); let procedure = ParsedLexicon::parse( procedure_lexicon_json(), 1, None, ProcedureAction::Upsert, None, None, None, ) .unwrap(); let defs = ParsedLexicon::parse( definitions_lexicon_json(), 1, None, ProcedureAction::Upsert, None, None, None, ) .unwrap();
reg.upsert(record).await; reg.upsert(query).await; reg.upsert(procedure).await; reg.upsert(defs).await;
assert_eq!(reg.count().await, 4);
let records = reg.get_record_collections().await; assert_eq!(records.len(), 1); assert!(records.contains(&"games.gamesgamesgamesgames.game".to_string()));
let queries = reg.get_queries().await; assert_eq!(queries.len(), 1); assert!(queries.contains(&"games.gamesgamesgamesgames.listGames".to_string()));
let procedures = reg.get_procedures().await; assert_eq!(procedures.len(), 1); assert!(procedures.contains(&"games.gamesgamesgamesgames.createGame".to_string())); }
// ----------------------------------------------------------------------- // ProcedureAction // -----------------------------------------------------------------------
#[test] fn procedure_action_from_none_is_upsert() { assert_eq!( ProcedureAction::from_optional_str(None).unwrap(), ProcedureAction::Upsert ); }
#[test] fn procedure_action_from_known_values() { assert_eq!( ProcedureAction::from_optional_str(Some("create")).unwrap(), ProcedureAction::Create ); assert_eq!( ProcedureAction::from_optional_str(Some("update")).unwrap(), ProcedureAction::Update ); assert_eq!( ProcedureAction::from_optional_str(Some("delete")).unwrap(), ProcedureAction::Delete ); assert_eq!( ProcedureAction::from_optional_str(Some("upsert")).unwrap(), ProcedureAction::Upsert ); }
#[test] fn procedure_action_from_invalid_returns_error() { let result = ProcedureAction::from_optional_str(Some("invalid")); assert!(result.is_err()); assert!(result.unwrap_err().contains("invalid")); }
#[test] fn procedure_action_to_optional_str_roundtrip() { assert_eq!(ProcedureAction::Create.to_optional_str(), Some("create")); assert_eq!(ProcedureAction::Update.to_optional_str(), Some("update")); assert_eq!(ProcedureAction::Delete.to_optional_str(), Some("delete")); assert_eq!(ProcedureAction::Upsert.to_optional_str(), None); }
// ----------------------------------------------------------------------- // index_hook // -----------------------------------------------------------------------
#[test] fn parse_preserves_index_hook() { let parsed = ParsedLexicon::parse( record_lexicon_json(), 1, None, ProcedureAction::Upsert, None, Some("function handle() end".into()), None, ) .unwrap(); assert_eq!(parsed.index_hook, Some("function handle() end".into())); }
#[test] fn parse_index_hook_none_by_default() { let parsed = ParsedLexicon::parse( record_lexicon_json(), 1, None, ProcedureAction::Upsert, None, None, None, ) .unwrap(); assert!(parsed.index_hook.is_none()); }
#[tokio::test] async fn registry_get_index_hook_returns_script() { let reg = LexiconRegistry::new(); let parsed = ParsedLexicon::parse( record_lexicon_json(), 1, None, ProcedureAction::Upsert, None, Some("function handle() log('hook') end".into()), None, ) .unwrap(); reg.upsert(parsed).await;
let script = reg.get_index_hook("games.gamesgamesgamesgames.game").await; assert_eq!(script, Some("function handle() log('hook') end".into())); }
#[tokio::test] async fn registry_get_index_hook_returns_none_when_absent() { let reg = LexiconRegistry::new(); let parsed = ParsedLexicon::parse( record_lexicon_json(), 1, None, ProcedureAction::Upsert, None, None, None, ) .unwrap(); reg.upsert(parsed).await;
let script = reg.get_index_hook("games.gamesgamesgamesgames.game").await; assert!(script.is_none()); }
#[tokio::test] async fn registry_get_index_hook_returns_none_for_unknown() { let reg = LexiconRegistry::new(); let script = reg.get_index_hook("nonexistent").await; assert!(script.is_none()); }
#[test] fn parse_space_type_from_lexicon() { let raw = json!({ "lexicon": 1, "id": "com.example.forum.post", "spaceType": "com.example.forum", "defs": { "main": { "type": "record", "key": "tid", "record": { "type": "object", "properties": { "text": { "type": "string" } } } } } }); let parsed = ParsedLexicon::parse(raw, 1, None, ProcedureAction::Upsert, None, None, None).unwrap(); assert_eq!(parsed.space_type.as_deref(), Some("com.example.forum")); }
#[test] fn parse_space_type_none_by_default() { let parsed = ParsedLexicon::parse( record_lexicon_json(), 1, None, ProcedureAction::Upsert, None, None, None, ) .unwrap(); assert!(parsed.space_type.is_none()); }}