use async_trait::async_trait;
use std::sync::Arc;
use crate::ast::{Command, CommandType};
use crate::error::KipError;
use crate::parser::parse_kip;
use crate::request::{
ExecutionMode, OnError, Operation, OperationResult, OperationStatus, Request, Response,
ResponseExecution, TopLevelStatus,
};
#[async_trait]
pub trait Executor: Send + Sync {
async fn execute(&self, command: Command, request: &Request, operation: &Operation)
-> Response;
}
#[async_trait]
impl Executor for Box<dyn Executor> {
async fn execute(
&self,
command: Command,
request: &Request,
operation: &Operation,
) -> Response {
(**self).execute(command, request, operation).await
}
}
#[async_trait]
impl Executor for Arc<dyn Executor> {
async fn execute(
&self,
command: Command,
request: &Request,
operation: &Operation,
) -> Response {
(**self).execute(command, request, operation).await
}
}
#[async_trait]
impl Executor for &dyn Executor {
async fn execute(
&self,
command: Command,
request: &Request,
operation: &Operation,
) -> Response {
(**self).execute(command, request, operation).await
}
}
pub async fn execute_kip(
executor: &impl Executor,
command: &str,
dry_run: bool,
) -> (CommandType, Response) {
execute_one(executor, command, dry_run, |_| Ok(())).await
}
pub async fn execute_readonly(
executor: &impl Executor,
command: &str,
dry_run: bool,
) -> (CommandType, Response) {
execute_one(executor, command, dry_run, admits_readonly).await
}
fn admits_readonly(command: &Command) -> Result<(), KipError> {
if command.is_mutation() {
return Err(KipError::readonly_violation(
"this endpoint executes KQL and META only; KML mutations must go through the \
state-capable runtime",
));
}
Ok(())
}
async fn execute_one(
executor: &impl Executor,
text: &str,
dry_run: bool,
admits: impl Fn(&Command) -> Result<(), KipError>,
) -> (CommandType, Response) {
let command = match parse_kip(text) {
Ok(command) => command,
Err(err) => return (CommandType::Unknown, err.into()),
};
let language = CommandType::from(&command);
if let Err(err) = admits(&command) {
return (language, err.into());
}
let request = single_command_request(text, dry_run);
let response = executor
.execute(command, &request, &request.operations[0])
.await;
(language, response)
}
fn single_command_request(command: &str, dry_run: bool) -> Request {
let mut request = Request::single(command);
request.options = Some(crate::request::RequestOptions {
dry_run: Some(dry_run),
..Default::default()
});
request
}
pub async fn execute_request(executor: &impl Executor, request: &Request) -> Response {
run_request(executor, request, |_| Ok(())).await
}
pub async fn execute_request_readonly(executor: &impl Executor, request: &Request) -> Response {
run_request(executor, request, admits_readonly).await
}
async fn run_request(
executor: &impl Executor,
request: &Request,
admits: impl Fn(&Command) -> Result<(), KipError>,
) -> Response {
if let Err(err) = request.validate() {
return Response::from(err).with_request_id(request.request_id.clone());
}
let mode = request.execution_mode();
if mode == ExecutionMode::Atomic {
return Response::from(KipError::unsupported_capability(
"atomic execution needs one transaction, one snapshot and all-or-none commit; this \
helper runs operations one at a time and will not fake them",
))
.with_request_id(request.request_id.clone());
}
let parsed: Vec<Result<Command, KipError>> =
request.operations.iter().map(Operation::parse).collect();
for command in parsed.iter().filter_map(|parsed| parsed.as_ref().ok()) {
if let Err(err) = admits(command) {
return Response::from(err).with_request_id(request.request_id.clone());
}
}
let on_error = request
.execution
.as_ref()
.and_then(|e| e.on_error)
.unwrap_or(OnError::Stop);
let mut results = Vec::with_capacity(request.operations.len());
let mut stopped = false;
let mut receipt = None;
let mut snapshot = None;
let mut outcome_unknown_error = None;
for (operation, parsed) in request.operations.iter().zip(parsed) {
if stopped {
results.push(
OperationResult {
status: OperationStatus::Skipped,
..Default::default()
}
.with_op_id(operation.op_id.clone()),
);
continue;
}
let result = match parsed {
Ok(command) => {
let response = executor.execute(command, request, operation).await;
if response.status == TopLevelStatus::OutcomeUnknown {
if outcome_unknown_error.is_none() {
receipt = response.receipt.clone();
outcome_unknown_error = Some(response.error.clone().unwrap_or_else(|| {
KipError::outcome_unknown(
"the executor could not establish whether the operation committed",
)
.into()
}));
}
} else if outcome_unknown_error.is_none() && response.receipt.is_some() {
receipt = response.receipt.clone();
}
if response.snapshot.is_some() {
snapshot = response.snapshot.clone();
}
operation_result_from(response)
}
Err(err) => OperationResult::failed(err),
}
.with_op_id(operation.op_id.clone());
if result.status == OperationStatus::Failed
&& mode == ExecutionMode::Sequence
&& on_error == OnError::Stop
{
stopped = true;
}
results.push(result);
}
Response {
status: if outcome_unknown_error.is_some() {
TopLevelStatus::OutcomeUnknown
} else {
TopLevelStatus::derive(&results)
},
execution: Some(ResponseExecution {
mode,
on_error: Some(on_error),
isolation: request.execution.as_ref().and_then(|e| e.isolation.clone()),
idempotency_key: request
.execution
.as_ref()
.and_then(|e| e.idempotency_key.clone()),
extensions: None,
}),
results,
receipt,
snapshot,
error: outcome_unknown_error,
..Default::default()
}
.with_request_id(request.request_id.clone())
}
fn operation_result_from(response: Response) -> OperationResult {
let Response {
results,
warnings,
next_cursor,
error,
..
} = response;
if let Some(mut result) = results.into_iter().next() {
result.warnings.extend(warnings);
result.next_cursor = result.next_cursor.or(next_cursor);
return result;
}
match error {
Some(error) => OperationResult {
status: OperationStatus::Failed,
error: Some(error),
warnings,
next_cursor,
..Default::default()
},
None => OperationResult {
warnings,
next_cursor,
..OperationResult::no_effect()
},
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::ast::Json;
use crate::error::KipErrorCode;
use crate::request::{Execution, Operation, Receipt, ReceiptStatus};
fn committed_receipt(tx_id: &str) -> Receipt {
Receipt {
status: ReceiptStatus::Committed,
tx_id: Some(tx_id.into()),
space_seq: Some(42),
space_id: Some("space-1".into()),
snapshot_seq: Some(41),
committed_at: Some("2026-08-16T00:00:00Z".into()),
transaction_class: None,
request_digest: None,
semantic_plan_digest: None,
result_digest: None,
schema_environment_version: None,
change_summary: None,
proofs: vec![],
receipt_digest: None,
origin: None,
extensions: None,
}
}
struct EchoNexus;
#[async_trait]
impl Executor for EchoNexus {
async fn execute(
&self,
command: Command,
_request: &Request,
_operation: &Operation,
) -> Response {
Response::ok(Json::String(CommandType::from(&command).to_string()))
}
}
struct FailingNexus;
#[async_trait]
impl Executor for FailingNexus {
async fn execute(
&self,
_command: Command,
_request: &Request,
_operation: &Operation,
) -> Response {
Response::from(KipError::not_found_or_not_visible("nothing here"))
}
}
#[tokio::test]
async fn execute_kip_classifies_what_it_ran() {
let (language, response) =
execute_kip(&EchoNexus, r#"FIND(?x) WHERE { ?x {type: "T"} }"#, false).await;
assert_eq!(language, CommandType::Kql);
assert_eq!(response.first_result(), Some(&Json::String("KQL".into())));
let (language, _) = execute_kip(&EchoNexus, "not a command", false).await;
assert_eq!(language, CommandType::Unknown);
}
#[tokio::test]
async fn the_readonly_path_rejects_writes_by_semantics() {
let (language, response) =
execute_readonly(&EchoNexus, r#"TRANSITION :x TO "tombstoned""#, false).await;
assert_eq!(language, CommandType::Kml);
assert_eq!(
response.error.unwrap().parsed_code(),
Some(KipErrorCode::ReadonlyViolation)
);
for command in [
"DESCRIBE PRIMER",
r#"EXPORT CAPSULE :out WHERE { ?c {type: "T"} }"#,
r#"FIND(?x) WHERE { ?x {type: "T"} }"#,
"DESCRIBE SNAPSHOT",
] {
let (_, response) = execute_readonly(&EchoNexus, command, false).await;
assert_eq!(response.status, TopLevelStatus::Succeeded, "for {command}");
}
}
#[tokio::test]
async fn a_sequence_stops_but_keeps_what_already_ran() {
let request = Request {
execution: Some(Execution {
on_error: Some(OnError::Stop),
..Execution::new(ExecutionMode::Sequence)
}),
operations: vec![
Operation::new("DESCRIBE PROTOCOL").with_op_id("op-1"),
Operation::new("DESCRIBE PRIMER").with_op_id("op-2"),
],
..Default::default()
};
let response = execute_request(&FailingNexus, &request).await;
assert_eq!(response.results[0].status, OperationStatus::Failed);
assert_eq!(response.results[1].status, OperationStatus::Skipped);
assert_eq!(response.results[1].op_id.as_deref(), Some("op-2"));
assert_eq!(response.status, TopLevelStatus::Failed);
}
#[tokio::test]
async fn independent_operations_isolate_their_failures() {
let request = Request {
execution: Some(Execution::new(ExecutionMode::Independent)),
operations: vec![
Operation::new("DESCRIBE PROTOCOL"),
Operation::new("nonsense"),
],
..Default::default()
};
let response = execute_request(&EchoNexus, &request).await;
assert_eq!(response.results[0].status, OperationStatus::Succeeded);
assert_eq!(response.results[1].status, OperationStatus::Failed);
assert_eq!(response.status, TopLevelStatus::Partial);
}
#[tokio::test]
async fn a_commit_receipt_survives_the_batch_runner() {
use crate::request::Warning;
struct Committing;
#[async_trait]
impl Executor for Committing {
async fn execute(
&self,
_command: Command,
_request: &Request,
_operation: &Operation,
) -> Response {
Response {
receipt: Some(committed_receipt("tx-9")),
warnings: vec![Warning::Message("index lagged".into())],
..Response::ok(Json::Bool(true))
}
}
}
let response = execute_request(
&Committing,
&Request::single(r#"TRANSITION :x TO "archived""#),
)
.await;
assert_eq!(
response.receipt.as_ref().and_then(|r| r.tx_id.as_deref()),
Some("tx-9")
);
assert_eq!(response.results[0].warnings.len(), 1);
}
#[tokio::test]
async fn the_readonly_envelope_path_rejects_writes_by_semantics() {
let request = Request {
execution: Some(Execution::new(ExecutionMode::Sequence)),
operations: vec![
Operation::new("DESCRIBE PRIMER").with_op_id("op-1"),
Operation::new(r#"TRANSITION :x TO "tombstoned""#).with_op_id("op-2"),
],
..Default::default()
};
let response = execute_request_readonly(&EchoNexus, &request).await;
assert_eq!(response.status, TopLevelStatus::Failed);
assert_eq!(
response.error.as_ref().unwrap().parsed_code(),
Some(KipErrorCode::ReadonlyViolation)
);
assert_eq!(response.results.len(), 1);
assert_eq!(response.results[0].op_id, None);
}
#[tokio::test]
async fn the_readonly_envelope_path_serves_reads_and_meta() {
let request = Request {
execution: Some(Execution::new(ExecutionMode::Independent)),
operations: vec![
Operation::new("DESCRIBE PRIMER"),
Operation::new(r#"EXPORT CAPSULE :out WHERE { ?c {type: "T"} }"#),
Operation::new(r#"FIND(?x) WHERE { ?x {type: "T"} }"#),
],
..Default::default()
};
let response = execute_request_readonly(&EchoNexus, &request).await;
assert_eq!(response.status, TopLevelStatus::Succeeded, "{response:#?}");
assert_eq!(response.results.len(), 3);
}
#[tokio::test]
async fn an_unparseable_operation_fails_on_its_own_result_on_the_readonly_path() {
let response =
execute_request_readonly(&EchoNexus, &Request::single("not a command")).await;
assert_eq!(response.results.len(), 1);
assert_eq!(
response.results[0]
.error
.as_ref()
.and_then(|error| error.parsed_code()),
Some(KipErrorCode::InvalidSyntax)
);
}
#[tokio::test]
async fn atomic_execution_is_refused_rather_than_faked() {
let request = Request {
execution: Some(Execution::new(ExecutionMode::Atomic)),
operations: vec![
Operation::new(r#"TRANSITION :a TO "archived""#),
Operation::new(r#"TRANSITION :b TO "archived""#),
],
..Default::default()
};
let response = execute_request(&EchoNexus, &request).await;
assert_eq!(
response.error.unwrap().parsed_code(),
Some(KipErrorCode::UnsupportedCapability)
);
}
#[tokio::test]
async fn an_invalid_envelope_never_reaches_the_executor() {
let request = Request {
kip: "1.0".into(),
..Request::single("DESCRIBE PROTOCOL")
};
let response = execute_request(&EchoNexus, &request).await;
assert_eq!(
response.error.unwrap().parsed_code(),
Some(KipErrorCode::UnsupportedProtocolVersion)
);
}
#[tokio::test]
async fn the_executor_receives_the_complete_request_and_operation_context() {
struct ContextAware;
#[async_trait]
impl Executor for ContextAware {
async fn execute(
&self,
_command: Command,
request: &Request,
operation: &Operation,
) -> Response {
Response::ok(serde_json::json!({
"space": request.space.as_ref().and_then(|space| space.id.clone()),
"request_parameter": request
.parameters
.as_ref()
.and_then(|parameters| parameters.get("request_value"))
.cloned(),
"operation_parameter": operation
.parameters
.as_ref()
.and_then(|parameters| parameters.get("operation_value"))
.cloned(),
"dry_run": request.is_dry_run(),
}))
}
}
let request = serde_json::from_value::<Request>(serde_json::json!({
"kip": "2.0",
"space": {"id": "space-7"},
"operations": [{
"command": "DESCRIBE PROTOCOL",
"parameters": {"operation_value": 2}
}],
"parameters": {"request_value": 1},
"options": {"dry_run": true}
}))
.unwrap();
let response = execute_request(&ContextAware, &request).await;
assert_eq!(
response.results[0].result,
Some(serde_json::json!({
"space": "space-7",
"request_parameter": 1,
"operation_parameter": 2,
"dry_run": true,
}))
);
}
#[tokio::test]
async fn an_unknown_write_outcome_is_never_flattened_to_failed() {
struct UnknownOutcome;
#[async_trait]
impl Executor for UnknownOutcome {
async fn execute(
&self,
_command: Command,
_request: &Request,
_operation: &Operation,
) -> Response {
Response::outcome_unknown(KipError::outcome_unknown("connection dropped"))
}
}
let response = execute_request(
&UnknownOutcome,
&Request::single(r#"TRANSITION :x TO "archived""#),
)
.await;
assert_eq!(response.status, TopLevelStatus::OutcomeUnknown);
assert_eq!(
response
.error
.as_ref()
.and_then(|error| error.parsed_code()),
Some(KipErrorCode::OutcomeUnknown)
);
}
#[tokio::test]
async fn an_unknown_outcome_never_leaks_an_earlier_transactions_receipt() {
struct CommitThenUnknown;
#[async_trait]
impl Executor for CommitThenUnknown {
async fn execute(
&self,
_command: Command,
_request: &Request,
operation: &Operation,
) -> Response {
if operation.op_id.as_deref() == Some("known") {
Response {
receipt: Some(committed_receipt("tx-known")),
..Response::ok(Json::Bool(true))
}
} else {
Response::outcome_unknown(KipError::outcome_unknown("connection dropped"))
}
}
}
let request = Request {
execution: Some(Execution::new(ExecutionMode::Independent)),
operations: vec![
Operation::new(r#"TRANSITION :a TO "archived""#).with_op_id("known"),
Operation::new(r#"TRANSITION :b TO "archived""#).with_op_id("unknown"),
],
..Default::default()
};
let response = execute_request(&CommitThenUnknown, &request).await;
assert_eq!(response.status, TopLevelStatus::OutcomeUnknown);
assert_eq!(response.receipt, None);
}
}