use serde::{Deserialize, Serialize};
use serde_json::json;
use std::borrow::Cow;
use crate::{
CommandType, Json, Map,
error::KipError,
executor::{Executor, execute_kip, execute_readonly},
};
#[derive(Clone, Debug, Deserialize, Serialize)]
#[serde(untagged)]
pub enum CommandItem {
Simple(String),
WithParams {
command: String,
#[serde(default)]
parameters: Map<String, Json>,
},
}
#[derive(Clone, Debug, Default, Deserialize, Serialize)]
pub struct Request {
#[serde(default)]
pub command: String,
#[serde(default)]
pub commands: Vec<CommandItem>,
#[serde(default)]
pub parameters: Map<String, Json>,
#[serde(default)]
pub dry_run: bool,
#[serde(default)]
pub readonly: bool,
}
impl Request {
pub fn is_batch(&self) -> bool {
!self.commands.is_empty()
}
pub fn iter_commands(
&self,
) -> impl Iterator<Item = (Cow<'_, str>, Cow<'_, Map<String, Json>>)> {
let shared_params = &self.parameters;
let single_command = &self.command;
let commands = &self.commands;
let is_batch = !commands.is_empty();
let single_iter = if !is_batch {
Some(std::iter::once((
Cow::Borrowed(single_command.as_str()),
Cow::Borrowed(shared_params),
)))
} else {
None
};
let batch_iter = if is_batch {
Some(commands.iter().map(move |item| match item {
CommandItem::Simple(cmd) => {
(Cow::Borrowed(cmd.as_str()), Cow::Borrowed(shared_params))
}
CommandItem::WithParams {
command,
parameters,
} => {
if parameters.is_empty() {
(
Cow::Borrowed(command.as_str()),
Cow::Borrowed(shared_params),
)
} else if shared_params.is_empty() {
(Cow::Borrowed(command.as_str()), Cow::Borrowed(parameters))
} else {
let mut merged = shared_params.clone();
for (k, v) in parameters {
merged.insert(k.clone(), v.clone());
}
(Cow::Borrowed(command.as_str()), Cow::Owned(merged))
}
}
}))
} else {
None
};
single_iter
.into_iter()
.flatten()
.chain(batch_iter.into_iter().flatten())
}
fn substitute_params(command: &str, parameters: &Map<String, Json>) -> String {
if parameters.is_empty() {
return command.to_string();
}
let mut result = command.to_string();
for (key, value) in parameters {
let placeholder = format!(":{}", key);
let replacement = match value {
Json::Number(n) => n.to_string(),
Json::Bool(b) => b.to_string(),
Json::Null => "null".to_string(),
_ => serde_json::to_string(value).unwrap_or_else(|_| "null".to_string()),
};
result = result.replace(&placeholder, &replacement);
}
result
}
fn find_placeholders_in_strings(
command: &str,
parameters: &Map<String, Json>,
) -> Result<(), String> {
let mut warnings = Vec::new();
for key in parameters.keys() {
let placeholder = format!(":{}", key);
let mut search_start = 0;
while let Some(pos) = command[search_start..].find(&placeholder) {
let abs_pos = search_start + pos;
let before = &command[..abs_pos];
let mut in_string = false;
let mut chars = before.chars().peekable();
while let Some(ch) = chars.next() {
if ch == '\\' {
chars.next();
} else if ch == '"' {
in_string = !in_string;
}
}
if in_string {
warnings.push((key.clone(), abs_pos));
}
search_start = abs_pos + placeholder.len();
}
}
if warnings.is_empty() {
return Ok(());
}
let param_names: Vec<_> = warnings.iter().map(|(n, _)| n.as_str()).collect();
Err(format!(
"Possible cause: placeholder(s) {:?} appear to be inside quoted strings. \
Placeholders must occupy a full JSON value position, not be embedded \
inside strings.",
param_names
))
}
pub fn to_command(&self) -> Cow<'_, str> {
if self.parameters.is_empty() {
Cow::Borrowed(&self.command)
} else {
Cow::Owned(Self::substitute_params(&self.command, &self.parameters))
}
}
pub fn readonly(&mut self) -> &mut Self {
self.readonly = true;
self
}
pub async fn execute(&self, nexus: &impl Executor) -> (CommandType, Response) {
if self.is_batch() {
self.execute_batch(nexus).await
} else {
let command = self.to_command();
let (cmd_type, response) = if self.readonly {
execute_readonly(nexus, &command, self.dry_run).await
} else {
execute_kip(nexus, &command, self.dry_run).await
};
if let Response::Err { ref error, .. } = response
&& error.code.starts_with("KIP_1")
&& !self.parameters.is_empty()
{
let warnings = Self::find_placeholders_in_strings(&self.command, &self.parameters);
if let Err(extra_hint) = warnings {
let mut new_error = error.clone();
new_error.hint = Some(match &error.hint {
Some(existing) => format!("{} {}", existing, extra_hint),
None => extra_hint,
});
return (
cmd_type,
Response::Err {
error: new_error,
result: None,
},
);
}
}
(cmd_type, response)
}
}
async fn execute_batch(&self, nexus: &impl Executor) -> (CommandType, Response) {
let mut results = Vec::with_capacity(self.commands.len());
let mut command_type = CommandType::Unknown;
for (cmd, params) in self.iter_commands() {
let substituted = Self::substitute_params(&cmd, ¶ms);
let (cmd_type, response) = if self.readonly {
execute_readonly(nexus, &substituted, self.dry_run).await
} else {
execute_kip(nexus, &substituted, self.dry_run).await
};
if command_type != CommandType::Kml && cmd_type != CommandType::Unknown {
command_type = cmd_type.clone();
}
match response {
Response::Ok { .. } => {
results.push(response);
}
Response::Err { mut error, .. } => {
if error.code.starts_with("KIP_1") && !params.is_empty() {
let warnings = Self::find_placeholders_in_strings(&cmd, ¶ms);
if let Err(extra_hint) = warnings {
error.hint = Some(match &error.hint {
Some(existing) => format!("{} {}", existing, extra_hint),
None => extra_hint,
});
}
}
results.push(Response::err(error));
if cmd_type == CommandType::Kml {
return (cmd_type, Response::ok(json!(results)));
}
}
}
}
(command_type, Response::ok(json!(results)))
}
}
#[derive(Clone, Debug, Deserialize, Serialize, PartialEq, Eq)]
#[serde(untagged)]
pub enum Response {
Ok {
result: Json,
#[serde(skip_serializing_if = "Option::is_none")]
next_cursor: Option<String>,
},
Err {
error: ErrorObject,
#[serde(skip_serializing_if = "Option::is_none")]
result: Option<Json>,
},
}
impl Response {
pub fn ok(result: Json) -> Self {
Self::Ok {
result,
next_cursor: None,
}
}
pub fn err(error: impl Into<ErrorObject>) -> Self {
Self::Err {
error: error.into(),
result: None,
}
}
pub fn into_result(self) -> Result<Json, ErrorObject> {
match self {
Self::Ok { result, .. } => Ok(result),
Self::Err { error, .. } => Err(error),
}
}
}
#[derive(Clone, Debug, Default, Deserialize, Serialize, PartialEq, Eq)]
pub struct ErrorObject {
pub code: String,
pub message: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub hint: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub data: Option<Json>,
}
impl From<String> for ErrorObject {
fn from(message: String) -> Self {
ErrorObject {
code: "KIP_4000".to_string(),
message,
hint: None,
data: None,
}
}
}
impl From<serde_json::Error> for ErrorObject {
fn from(error: serde_json::Error) -> Self {
ErrorObject {
code: "KIP_1001".to_string(),
message: error.to_string(),
hint: Some("Check JSON data format is valid.".to_string()),
data: None,
}
}
}
impl From<KipError> for ErrorObject {
fn from(error: KipError) -> Self {
let hint = error.hint().to_string();
ErrorObject {
code: error.code_str().to_string(),
message: error.message,
hint: Some(hint),
data: None,
}
}
}
impl std::fmt::Display for ErrorObject {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
if let Some(hint) = &self.hint {
write!(f, "[{}] {}, Hint: {}", self.code, self.message, hint)
} else {
write!(f, "[{}] {}", self.code, self.message)
}
}
}
impl From<KipError> for Response {
fn from(error: KipError) -> Self {
Response::Err {
error: error.into(),
result: None,
}
}
}
impl<E> From<Result<(Json, Option<String>), E>> for Response
where
E: Into<ErrorObject>,
{
fn from(result: Result<(Json, Option<String>), E>) -> Self {
match result {
Ok((result, next_cursor)) => Response::Ok {
result,
next_cursor,
},
Err(err) => Response::Err {
error: err.into(),
result: None,
},
}
}
}
impl<E> From<Result<Json, E>> for Response
where
E: Into<ErrorObject>,
{
fn from(result: Result<Json, E>) -> Self {
match result {
Ok(result) => Response::Ok {
result,
next_cursor: None,
},
Err(err) => Response::Err {
error: err.into(),
result: None,
},
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{Command, parse_kml, parse_kql};
use async_trait::async_trait;
use serde_json::json;
#[test]
fn test_to_command_empty_parameters() {
let request = Request {
command: "FIND(?drug) WHERE { ?drug {type: \"Drug\"} }".to_string(),
..Default::default()
};
let result = request.to_command();
assert_eq!(result, "FIND(?drug) WHERE { ?drug {type: \"Drug\"} }");
assert!(matches!(result, Cow::Borrowed(_)));
assert!(parse_kql(&result).is_ok());
}
#[test]
fn test_to_command_string_parameter() {
let mut parameters = Map::new();
parameters.insert(
"symptom_name".to_string(),
Json::String("Headache".to_string()),
);
let request = Request {
command: "FIND(?symptom) WHERE { ?symptom {name: :symptom_name} }".to_string(),
parameters,
..Default::default()
};
let result = request.to_command();
assert_eq!(
result,
"FIND(?symptom) WHERE { ?symptom {name: \"Headache\"} }"
);
assert!(matches!(result, Cow::Owned(_)));
assert!(parse_kql(&result).is_ok());
}
#[test]
fn test_to_command_number_parameter() {
let mut parameters = Map::new();
parameters.insert(
"limit".to_string(),
Json::Number(serde_json::Number::from(10)),
);
parameters.insert(
"risk_level".to_string(),
Json::Number(serde_json::Number::from_f64(3.5).unwrap()),
);
let request = Request {
command: "FIND(?drug) WHERE { ?drug {type: \"Drug\"} FILTER(?drug.attributes.risk_level < :risk_level) } LIMIT :limit".to_string(),
parameters,
..Default::default()
};
let result = request.to_command();
assert_eq!(
result,
"FIND(?drug) WHERE { ?drug {type: \"Drug\"} FILTER(?drug.attributes.risk_level < 3.5) } LIMIT 10"
);
assert!(parse_kql(&result).is_ok());
}
#[test]
fn test_to_command_object_parameter() {
let mut parameters = Map::new();
parameters.insert(
"metadata".to_string(),
json!({"confidence": 0.95, "source": "clinical_trial"}),
);
let request = Request {
command:
"UPSERT { CONCEPT ?drug { {type: \"Drug\", name: \"TestDrug\"} } WITH METADATA :metadata }"
.to_string(),
parameters,
..Default::default()
};
let result = request.to_command();
assert_eq!(
result,
"UPSERT { CONCEPT ?drug { {type: \"Drug\", name: \"TestDrug\"} } WITH METADATA {\"confidence\":0.95,\"source\":\"clinical_trial\"} }"
);
assert!(parse_kml(&result).is_ok());
}
#[test]
fn test_to_command_multiple_parameters() {
let mut parameters = Map::new();
parameters.insert(
"symptom_name".to_string(),
Json::String("Headache".to_string()),
);
parameters.insert(
"limit".to_string(),
Json::Number(serde_json::Number::from(5)),
);
parameters.insert("include_experimental".to_string(), Json::Bool(false));
let request = Request {
command: r#"
FIND(?drug.name)
WHERE {
?symptom{name: :symptom_name}
(?drug, "treats", ?symptom)
FILTER(?drug.attributes.experimental == :include_experimental)
}
LIMIT :limit
"#
.to_string(),
parameters,
..Default::default()
};
let result = request.to_command();
let expected = r#"
FIND(?drug.name)
WHERE {
?symptom{name: "Headache"}
(?drug, "treats", ?symptom)
FILTER(?drug.attributes.experimental == false)
}
LIMIT 5
"#;
assert_eq!(result, expected);
assert!(parse_kql(&result).is_ok());
}
#[test]
fn test_to_command_same_parameter_multiple_times() {
let mut parameters = Map::new();
parameters.insert(
"drug_type".to_string(),
Json::String("Analgesic".to_string()),
);
let request = Request {
command:
"FIND(?drug1, ?drug2) WHERE { ?drug1 {type: :drug_type} ?drug2 {type: :drug_type} }"
.to_string(),
parameters,
..Default::default()
};
let result = request.to_command();
assert_eq!(
result,
"FIND(?drug1, ?drug2) WHERE { ?drug1 {type: \"Analgesic\"} ?drug2 {type: \"Analgesic\"} }"
);
assert!(parse_kql(&result).is_ok());
}
#[test]
fn test_to_command_parameter_not_found() {
let mut parameters = Map::new();
parameters.insert(
"existing_param".to_string(),
Json::String("value".to_string()),
);
let request = Request {
command: "FIND(?item) WHERE { ?item{name: :missing_param} }".to_string(),
parameters,
..Default::default()
};
let result = request.to_command();
assert_eq!(result, "FIND(?item) WHERE { ?item{name: :missing_param} }");
assert!(parse_kql(&result).is_err());
}
#[test]
fn test_to_command_special_characters_in_string() {
let mut parameters = Map::new();
parameters.insert(
"special_name".to_string(),
Json::String("Drug with \"quotes\" and :symbols".to_string()),
);
let request = Request {
command: "FIND(?drug) WHERE { ?drug{name: :special_name} }".to_string(),
parameters,
..Default::default()
};
let result = request.to_command();
assert_eq!(
result,
"FIND(?drug) WHERE { ?drug{name: \"Drug with \\\"quotes\\\" and :symbols\"} }"
);
assert!(parse_kql(&result).is_ok());
}
#[test]
fn test_to_command_complex_kip_example() {
let mut parameters = Map::new();
parameters.insert(
"symptom_name".to_string(),
Json::String("Brain Fog".to_string()),
);
parameters.insert(
"confidence_threshold".to_string(),
Json::Number(serde_json::Number::from_f64(0.8).unwrap()),
);
parameters.insert(
"max_results".to_string(),
Json::Number(serde_json::Number::from(20)),
);
let request = Request {
command: r#"
FIND(?drug.name, ?drug.metadata.confidence)
WHERE {
(?drug, "treats", {name: :symptom_name})
FILTER(?drug.metadata.confidence > :confidence_threshold)
}
ORDER BY ?drug.metadata.confidence DESC
LIMIT :max_results
"#
.to_string(),
parameters,
..Default::default()
};
let result = request.to_command();
let expected = r#"
FIND(?drug.name, ?drug.metadata.confidence)
WHERE {
(?drug, "treats", {name: "Brain Fog"})
FILTER(?drug.metadata.confidence > 0.8)
}
ORDER BY ?drug.metadata.confidence DESC
LIMIT 20
"#;
assert_eq!(result, expected);
assert!(parse_kql(expected).is_ok());
}
#[test]
fn test_response() {
let res = Response::ok(json!("Success"));
assert_eq!(
serde_json::to_string(&res).unwrap(),
r#"{"result":"Success"}"#
);
assert_eq!(
res,
serde_json::from_str(r#"{"result":"Success"}"#).unwrap()
);
let res = Response::Ok {
result: json!("Success"),
next_cursor: Some("abcdef".to_string()),
};
assert_eq!(
serde_json::to_string(&res).unwrap(),
r#"{"result":"Success","next_cursor":"abcdef"}"#
);
let res = Response::err(ErrorObject {
code: "KIP_4003".to_string(),
message: "An error occurred".to_string(),
hint: Some("Contact system administrator.".to_string()),
data: Some(json!("Additional info")),
});
assert_eq!(
serde_json::to_string(&res).unwrap(),
r#"{"error":{"code":"KIP_4003","message":"An error occurred","hint":"Contact system administrator.","data":"Additional info"}}"#
);
}
#[test]
fn test_command_item_deserialization() {
let simple: CommandItem = serde_json::from_str(r#""DESCRIBE PRIMER""#).unwrap();
assert!(matches!(simple, CommandItem::Simple(s) if s == "DESCRIBE PRIMER"));
let with_params: CommandItem = serde_json::from_str(
r#"{"command": "FIND(?x) WHERE { ?x {name: :name} }", "parameters": {"name": "test"}}"#,
)
.unwrap();
match with_params {
CommandItem::WithParams {
command,
parameters,
} => {
assert_eq!(command, "FIND(?x) WHERE { ?x {name: :name} }");
assert_eq!(
parameters.get("name"),
Some(&Json::String("test".to_string()))
);
}
_ => panic!("Expected WithParams variant"),
}
let with_empty_params: CommandItem =
serde_json::from_str(r#"{"command": "DESCRIBE PRIMER", "parameters": {}}"#).unwrap();
match with_empty_params {
CommandItem::WithParams {
command,
parameters,
} => {
assert_eq!(command, "DESCRIBE PRIMER");
assert!(parameters.is_empty());
}
_ => panic!("Expected WithParams variant"),
}
}
#[test]
fn test_batch_request_deserialization() {
let json_str = r#"{
"commands": [
"DESCRIBE PRIMER",
"FIND(?t.name) WHERE { ?t {type: \"$ConceptType\"} } LIMIT 50",
{
"command": "UPSERT { CONCEPT ?e { {type:\"Event\", name: :name} } }",
"parameters": { "name": "MyEvent" }
}
],
"parameters": { "limit": 10 }
}"#;
let request: Request = serde_json::from_str(json_str).unwrap();
assert!(request.is_batch());
assert_eq!(request.commands.len(), 3);
assert_eq!(
request.parameters.get("limit"),
Some(&Json::Number(10.into()))
);
assert!(matches!(&request.commands[0], CommandItem::Simple(s) if s == "DESCRIBE PRIMER"));
assert!(matches!(&request.commands[1], CommandItem::Simple(s) if s.contains("FIND")));
match &request.commands[2] {
CommandItem::WithParams {
command,
parameters,
} => {
assert!(command.contains("UPSERT"));
assert_eq!(
parameters.get("name"),
Some(&Json::String("MyEvent".to_string()))
);
}
_ => panic!("Expected WithParams variant"),
}
}
#[test]
fn test_iter_commands_single_mode() {
let mut parameters = Map::new();
parameters.insert("name".to_string(), Json::String("test".to_string()));
let request = Request {
command: "FIND(?x) WHERE { ?x {name: :name} }".to_string(),
commands: vec![],
parameters: parameters.clone(),
..Default::default()
};
assert!(!request.is_batch());
let items: Vec<_> = request.iter_commands().collect();
assert_eq!(items.len(), 1);
assert_eq!(items[0].0.as_ref(), "FIND(?x) WHERE { ?x {name: :name} }");
assert_eq!(
items[0].1.get("name"),
Some(&Json::String("test".to_string()))
);
}
#[test]
fn test_iter_commands_batch_mode() {
let mut shared_params = Map::new();
shared_params.insert("limit".to_string(), Json::Number(10.into()));
shared_params.insert(
"shared".to_string(),
Json::String("shared_value".to_string()),
);
let mut cmd_params = Map::new();
cmd_params.insert("name".to_string(), Json::String("MyEvent".to_string()));
cmd_params.insert("shared".to_string(), Json::String("overridden".to_string()));
let request = Request {
command: String::new(),
commands: vec![
CommandItem::Simple("DESCRIBE PRIMER".to_string()),
CommandItem::WithParams {
command: "FIND(?x) WHERE { ?x {type: :type} }".to_string(),
parameters: Map::new(),
},
CommandItem::WithParams {
command: "UPSERT { CONCEPT ?e { {name: :name} } }".to_string(),
parameters: cmd_params,
},
],
parameters: shared_params,
..Default::default()
};
assert!(request.is_batch());
let items: Vec<_> = request.iter_commands().collect();
assert_eq!(items.len(), 3);
assert_eq!(items[0].0.as_ref(), "DESCRIBE PRIMER");
assert_eq!(items[0].1.get("limit"), Some(&Json::Number(10.into())));
assert_eq!(items[1].0.as_ref(), "FIND(?x) WHERE { ?x {type: :type} }");
assert_eq!(items[1].1.get("limit"), Some(&Json::Number(10.into())));
assert_eq!(
items[2].0.as_ref(),
"UPSERT { CONCEPT ?e { {name: :name} } }"
);
assert_eq!(
items[2].1.get("name"),
Some(&Json::String("MyEvent".to_string()))
);
assert_eq!(items[2].1.get("limit"), Some(&Json::Number(10.into())));
assert_eq!(
items[2].1.get("shared"),
Some(&Json::String("overridden".to_string()))
);
}
#[test]
fn test_batch_request_serialization() {
let mut cmd_params = Map::new();
cmd_params.insert("name".to_string(), Json::String("MyEvent".to_string()));
let mut shared_params = Map::new();
shared_params.insert("limit".to_string(), Json::Number(10.into()));
let request = Request {
command: String::new(),
commands: vec![
CommandItem::Simple("DESCRIBE PRIMER".to_string()),
CommandItem::WithParams {
command: "UPSERT { CONCEPT ?e { {name: :name} } }".to_string(),
parameters: cmd_params,
},
],
parameters: shared_params,
dry_run: true,
readonly: true,
};
let json_str = serde_json::to_string(&request).unwrap();
let parsed: Request = serde_json::from_str(&json_str).unwrap();
assert!(parsed.is_batch());
assert_eq!(parsed.commands.len(), 2);
assert!(parsed.dry_run);
assert!(parsed.readonly);
}
#[test]
fn test_single_command_mode_default() {
let request = Request {
command: "DESCRIBE PRIMER".to_string(),
..Default::default()
};
assert!(!request.is_batch());
assert!(request.commands.is_empty());
}
#[test]
fn test_validate_placeholder_usage_valid() {
let mut parameters = Map::new();
parameters.insert("name".to_string(), Json::String("John".to_string()));
parameters.insert("age".to_string(), Json::Number(25.into()));
let request = Request {
command: r#"SET ATTRIBUTES { name: :name, age: :age }"#.to_string(),
parameters,
..Default::default()
};
let warnings = Request::find_placeholders_in_strings(&request.command, &request.parameters);
assert!(warnings.is_ok(), "Should have no warnings for valid usage");
}
#[test]
fn test_validate_placeholder_usage_invalid() {
let mut parameters = Map::new();
parameters.insert("user_id".to_string(), Json::String("user123".to_string()));
let request = Request {
command: r#"SET ATTRIBUTES { summary: "Hello :user_id, welcome!" }"#.to_string(),
parameters,
..Default::default()
};
let warnings = Request::find_placeholders_in_strings(&request.command, &request.parameters);
assert!(
warnings.unwrap_err().contains("user_id"),
"Should warn about 'user_id' placeholder"
);
}
#[test]
fn test_validate_placeholder_usage_mixed() {
let mut parameters = Map::new();
parameters.insert("name".to_string(), Json::String("John".to_string()));
parameters.insert("user_id".to_string(), Json::String("user123".to_string()));
let request = Request {
command: r#"SET ATTRIBUTES { name: :name, summary: "User :user_id joined" }"#
.to_string(),
parameters,
..Default::default()
};
let warnings = Request::find_placeholders_in_strings(&request.command, &request.parameters);
assert!(
warnings.unwrap_err().contains("user_id"),
"Should warn about 'user_id' placeholder"
);
}
#[test]
fn test_validate_placeholder_usage_escaped_quotes() {
let mut parameters = Map::new();
parameters.insert("name".to_string(), Json::String("John".to_string()));
let request = Request {
command: r#"SET ATTRIBUTES { desc: "Say \"hello\" to :name" }"#.to_string(),
parameters,
..Default::default()
};
let warnings = Request::find_placeholders_in_strings(&request.command, &request.parameters);
assert!(
warnings.unwrap_err().contains("name"),
"Should warn about 'name' placeholder"
);
}
#[derive(Debug)]
struct MockExecutor;
#[async_trait]
impl Executor for MockExecutor {
async fn execute(&self, command: Command, _dry_run: bool) -> Response {
match command {
Command::Kql(query) => Response::Ok {
result: json!({
"type": "kql",
"find_count": query.find_clause.expressions.len()
}),
next_cursor: None,
},
Command::Kml(_) => Response::ok(json!({"type": "kml", "upserted": 1})),
Command::Meta(_) => Response::ok(json!({"type": "meta"})),
}
}
}
#[derive(Debug)]
struct FailingExecutor;
#[async_trait]
impl Executor for FailingExecutor {
async fn execute(&self, _command: Command, _dry_run: bool) -> Response {
Response::err(ErrorObject {
code: "KIP_3001".to_string(),
message: "Not found".to_string(),
hint: Some("Check that the referenced concept exists.".to_string()),
data: None,
})
}
}
#[derive(Debug)]
struct KmlFailingExecutor;
#[async_trait]
impl Executor for KmlFailingExecutor {
async fn execute(&self, command: Command, _dry_run: bool) -> Response {
match command {
Command::Kql(_) => Response::ok(json!({"type": "kql"})),
Command::Meta(_) => Response::ok(json!({"type": "meta"})),
Command::Kml(_) => Response::err(ErrorObject {
code: "KIP_4002".to_string(),
message: "Write failed".to_string(),
hint: Some("Check write constraints before retrying.".to_string()),
data: None,
}),
}
}
}
#[tokio::test]
async fn test_execute_single_kql_command() {
let executor = MockExecutor;
let request = Request {
command: r#"FIND(?drug.name) WHERE { ?drug {type: "Drug"} } LIMIT 10"#.to_string(),
..Default::default()
};
let (cmd_type, response) = request.execute(&executor).await;
assert_eq!(cmd_type, CommandType::Kql);
match response {
Response::Ok {
result,
next_cursor,
} => {
assert_eq!(result["type"], "kql");
assert_eq!(result["find_count"], 1);
assert!(next_cursor.is_none());
}
_ => panic!("Expected Ok response"),
}
}
#[tokio::test]
async fn test_execute_single_kml_command() {
let executor = MockExecutor;
let request = Request {
command: r#"UPSERT { CONCEPT ?d { {type: "Drug", name: "Aspirin"} } }"#.to_string(),
..Default::default()
};
let (cmd_type, response) = request.execute(&executor).await;
assert_eq!(cmd_type, CommandType::Kml);
match response {
Response::Ok { result, .. } => {
assert_eq!(result["type"], "kml");
assert_eq!(result["upserted"], 1);
}
_ => panic!("Expected Ok response"),
}
}
#[tokio::test]
async fn test_execute_single_meta_command() {
let executor = MockExecutor;
let request = Request {
command: "DESCRIBE PRIMER".to_string(),
..Default::default()
};
let (cmd_type, response) = request.execute(&executor).await;
assert_eq!(cmd_type, CommandType::Meta);
match response {
Response::Ok { result, .. } => {
assert_eq!(result["type"], "meta");
}
_ => panic!("Expected Ok response"),
}
}
#[tokio::test]
async fn test_execute_single_command_with_params() {
let executor = MockExecutor;
let mut parameters = Map::new();
parameters.insert("name".to_string(), Json::String("Aspirin".to_string()));
parameters.insert("limit".to_string(), Json::Number(5.into()));
let request = Request {
command: r#"FIND(?drug) WHERE { ?drug {name: :name} } LIMIT :limit"#.to_string(),
parameters,
..Default::default()
};
let (cmd_type, response) = request.execute(&executor).await;
assert_eq!(cmd_type, CommandType::Kql);
assert!(matches!(response, Response::Ok { .. }));
}
#[tokio::test]
async fn test_execute_single_command_syntax_error() {
let executor = MockExecutor;
let request = Request {
command: "INVALID COMMAND SYNTAX".to_string(),
..Default::default()
};
let (cmd_type, response) = request.execute(&executor).await;
assert_eq!(cmd_type, CommandType::Unknown);
match response {
Response::Err { error, .. } => {
assert!(error.code.starts_with("KIP_1"));
}
_ => panic!("Expected Err response for syntax error"),
}
}
#[tokio::test]
async fn test_execute_single_command_syntax_error_with_placeholder_hint() {
let executor = MockExecutor;
let mut parameters = Map::new();
parameters.insert("name".to_string(), Json::String("test".to_string()));
let request = Request {
command: r#"FIND(?x) WHERE { ?x {name: "Hello :name"} }"#.to_string(),
parameters,
..Default::default()
};
let (cmd_type, response) = request.execute(&executor).await;
assert_eq!(cmd_type, CommandType::Unknown);
match response {
Response::Err { error, .. } => {
assert!(error.code.starts_with("KIP_1"));
assert!(error.hint.as_ref().unwrap().contains("name"));
}
_ => panic!("Expected Err response"),
}
}
#[tokio::test]
async fn test_execute_single_command_executor_error() {
let executor = FailingExecutor;
let request = Request {
command: r#"FIND(?drug) WHERE { ?drug {type: "Drug"} }"#.to_string(),
..Default::default()
};
let (cmd_type, response) = request.execute(&executor).await;
assert_eq!(cmd_type, CommandType::Kql);
match response {
Response::Err { error, .. } => {
assert_eq!(error.code, "KIP_3001");
assert_eq!(error.message, "Not found");
}
_ => panic!("Expected Err response from failing executor"),
}
}
#[tokio::test]
async fn test_execute_batch_all_success() {
let executor = MockExecutor;
let request = Request {
commands: vec![
CommandItem::Simple("DESCRIBE PRIMER".to_string()),
CommandItem::Simple(
r#"FIND(?t.name) WHERE { ?t {type: "$ConceptType"} } LIMIT 50"#.to_string(),
),
],
..Default::default()
};
let (cmd_type, response) = request.execute(&executor).await;
assert_eq!(cmd_type, CommandType::Kql);
match response {
Response::Ok { result, .. } => {
let arr = result.as_array().unwrap();
assert_eq!(arr.len(), 2);
assert_eq!(arr[0]["result"]["type"], "meta");
assert_eq!(arr[1]["result"]["type"], "kql");
}
_ => panic!("Expected Ok response for batch"),
}
}
#[tokio::test]
async fn test_execute_batch_with_params() {
let executor = MockExecutor;
let mut shared_params = Map::new();
shared_params.insert("limit".to_string(), Json::Number(10.into()));
let mut cmd_params = Map::new();
cmd_params.insert("name".to_string(), Json::String("TestDrug".to_string()));
let request = Request {
commands: vec![
CommandItem::Simple("DESCRIBE PRIMER".to_string()),
CommandItem::WithParams {
command: r#"UPSERT { CONCEPT ?d { {type: "Drug", name: :name} } }"#.to_string(),
parameters: cmd_params,
},
],
parameters: shared_params,
..Default::default()
};
let (cmd_type, response) = request.execute(&executor).await;
assert_eq!(cmd_type, CommandType::Kml);
match response {
Response::Ok { result, .. } => {
let arr = result.as_array().unwrap();
assert_eq!(arr.len(), 2);
assert_eq!(arr[0]["result"]["type"], "meta");
assert_eq!(arr[1]["result"]["type"], "kml");
}
_ => panic!("Expected Ok response for batch with params"),
}
}
#[tokio::test]
async fn test_execute_batch_continues_on_syntax_error() {
let executor = MockExecutor;
let request = Request {
commands: vec![
CommandItem::Simple("DESCRIBE PRIMER".to_string()),
CommandItem::Simple("INVALID SYNTAX HERE".to_string()),
CommandItem::Simple(r#"FIND(?t) WHERE { ?t {type: "$ConceptType"} }"#.to_string()),
],
..Default::default()
};
let (_cmd_type, response) = request.execute(&executor).await;
match response {
Response::Ok { result, .. } => {
let arr = result.as_array().unwrap();
assert_eq!(arr.len(), 3);
assert!(arr[0]["result"].is_object());
assert!(arr[1]["error"].is_object());
assert!(
arr[1]["error"]["code"]
.as_str()
.unwrap()
.starts_with("KIP_1")
);
assert_eq!(arr[2]["result"]["type"], "kql");
}
_ => panic!("Expected Ok response wrapping batch results"),
}
}
#[tokio::test]
async fn test_execute_batch_continues_on_non_kml_executor_error() {
let executor = FailingExecutor;
let request = Request {
commands: vec![
CommandItem::Simple(r#"FIND(?t) WHERE { ?t {type: "$ConceptType"} }"#.to_string()),
CommandItem::Simple("DESCRIBE PRIMER".to_string()),
],
..Default::default()
};
let (_cmd_type, response) = request.execute(&executor).await;
match response {
Response::Ok { result, .. } => {
let arr = result.as_array().unwrap();
assert_eq!(arr.len(), 2);
assert!(arr[0]["error"].is_object());
assert_eq!(arr[0]["error"]["code"], "KIP_3001");
assert!(arr[1]["error"].is_object());
assert_eq!(arr[1]["error"]["code"], "KIP_3001");
}
_ => panic!("Expected Ok response wrapping batch results"),
}
}
#[tokio::test]
async fn test_execute_batch_stops_on_kml_error() {
let executor = KmlFailingExecutor;
let request = Request {
commands: vec![
CommandItem::Simple("DESCRIBE PRIMER".to_string()),
CommandItem::Simple(
r#"UPSERT { CONCEPT ?d { {type: "Drug", name: "NewDrug"} } }"#.to_string(),
),
CommandItem::Simple(r#"FIND(?d) WHERE { ?d {type: "Drug"} } LIMIT 5"#.to_string()),
],
..Default::default()
};
let (cmd_type, response) = request.execute(&executor).await;
assert_eq!(cmd_type, CommandType::Kml);
match response {
Response::Ok { result, .. } => {
let arr = result.as_array().unwrap();
assert_eq!(arr.len(), 2);
assert_eq!(arr[0]["result"]["type"], "meta");
assert!(arr[1]["error"].is_object());
assert_eq!(arr[1]["error"]["code"], "KIP_4002");
}
_ => panic!("Expected Ok response wrapping batch results"),
}
}
#[tokio::test]
async fn test_execute_batch_mixed_command_types() {
let executor = MockExecutor;
let request = Request {
commands: vec![
CommandItem::Simple("DESCRIBE PRIMER".to_string()),
CommandItem::Simple(r#"FIND(?d) WHERE { ?d {type: "Drug"} } LIMIT 5"#.to_string()),
CommandItem::Simple(
r#"UPSERT { CONCEPT ?d { {type: "Drug", name: "NewDrug"} } }"#.to_string(),
),
],
..Default::default()
};
let (cmd_type, response) = request.execute(&executor).await;
assert_eq!(cmd_type, CommandType::Kml);
match response {
Response::Ok { result, .. } => {
let arr = result.as_array().unwrap();
assert_eq!(arr.len(), 3);
assert_eq!(arr[0]["result"]["type"], "meta");
assert_eq!(arr[1]["result"]["type"], "kql");
assert_eq!(arr[2]["result"]["type"], "kml");
}
_ => panic!("Expected Ok response"),
}
}
#[tokio::test]
async fn test_execute_batch_empty_commands() {
let executor = MockExecutor;
let request = Request {
commands: vec![],
..Default::default()
};
assert!(!request.is_batch());
let (cmd_type, response) = request.execute(&executor).await;
assert_eq!(cmd_type, CommandType::Unknown);
assert!(matches!(response, Response::Err { .. }));
}
#[tokio::test]
async fn test_execute_batch_syntax_error_with_placeholder_hint() {
let executor = MockExecutor;
let mut params = Map::new();
params.insert("name".to_string(), Json::String("test".to_string()));
let request = Request {
commands: vec![
CommandItem::Simple("DESCRIBE PRIMER".to_string()),
CommandItem::WithParams {
command: r#"FIND(?x) WHERE { ?x {name: "Hello :name"} }"#.to_string(),
parameters: params,
},
],
..Default::default()
};
let (_cmd_type, response) = request.execute(&executor).await;
match response {
Response::Ok { result, .. } => {
let arr = result.as_array().unwrap();
assert_eq!(arr.len(), 2);
assert!(arr[0]["result"].is_object());
let error = &arr[1]["error"];
assert!(error["code"].as_str().unwrap().starts_with("KIP_1"));
assert!(error["hint"].as_str().unwrap().contains("name"));
}
_ => panic!("Expected Ok response wrapping batch results"),
}
}
}