use serde::{Deserialize, Serialize};
use std::collections::BTreeMap;
use crate::ast::{Command, CommandType, Json, Map};
use crate::error::{ErrorObject, KipError, KipErrorCode};
use crate::parser::{MAX_KIP_BATCH_COMMANDS, parse_kip, validate_command};
pub const KIP_VERSION: &str = "2.0";
pub mod limits {
pub const REQUEST_ID: usize = 256;
pub const IDEMPOTENCY_KEY: usize = 1024;
pub const SPACE_ID: usize = 512;
pub const SPACE_URI: usize = 2048;
pub const COMPATIBILITY_PROFILE: usize = 128;
pub const OPAQUE_TOKEN: usize = 8192;
pub const ISOLATION: usize = 64;
pub const SHORT_LABEL: usize = 256;
pub const RISK: usize = 128;
pub const LOCALE: usize = 64;
pub const SOURCE_ACTOR: usize = 512;
pub const ELEMENT_REFERENCE_KEY: usize = 1024;
pub const FACET_NAME: usize = 512;
pub const CLIENT_KEY: usize = 1024;
}
#[derive(Clone, Debug, Deserialize, Serialize, PartialEq)]
#[serde(deny_unknown_fields)]
pub struct Request {
pub kip: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub request_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub space: Option<SpaceSelector>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub compatibility_profile: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub execution: Option<Execution>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub read: Option<ReadBinding>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub ingest: Option<IngestContext>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub preconditions: Option<Preconditions>,
pub operations: Vec<Operation>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub parameters: Option<Map<String, Json>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub context: Option<RequestContext>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub requires: Option<Map<String, Json>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub options: Option<RequestOptions>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<Map<String, Json>>,
}
impl Default for Request {
fn default() -> Self {
Self {
kip: KIP_VERSION.to_string(),
request_id: None,
space: None,
compatibility_profile: None,
execution: None,
read: None,
ingest: None,
preconditions: None,
operations: Vec::new(),
parameters: None,
context: None,
requires: None,
options: None,
extensions: None,
}
}
}
impl Request {
pub fn single(command: impl Into<String>) -> Self {
Self {
operations: vec![Operation::new(command)],
..Default::default()
}
}
pub fn is_dry_run(&self) -> bool {
self.options
.as_ref()
.and_then(|o| o.dry_run)
.unwrap_or(false)
}
pub fn execution_mode(&self) -> ExecutionMode {
self.execution
.as_ref()
.map(|e| e.mode)
.unwrap_or(ExecutionMode::Independent)
}
pub fn from_json(source: &str) -> Result<Self, KipError> {
Self::from_value(crate::parse_canonical_json(source)?)
}
pub fn from_value(value: Json) -> Result<Self, KipError> {
if let Some(entries) = value.pointer("/ingest/evidence").and_then(Json::as_array) {
for entry in entries {
if let Some(at) = entry.get("observed_at") {
crate::timestamp::validate_value(at, "ingest.observed_at")?;
}
}
}
let request: Self = serde_json::from_value(value)
.map_err(|e| KipError::invalid_request_envelope(e.to_string()))?;
request.validate()?;
Ok(request)
}
pub fn validate(&self) -> Result<(), KipError> {
crate::validate_json(
&serde_json::to_value(self)
.map_err(|e| KipError::invalid_request_envelope(e.to_string()))?,
)?;
if self.kip != KIP_VERSION {
return Err(KipError::unsupported_protocol_version(format!(
"this runtime speaks KIP {KIP_VERSION}, the request declares {:?}",
self.kip
)));
}
if self.operations.is_empty() {
return Err(KipError::invalid_request_envelope(
"a request must carry at least one operation",
));
}
if self.operations.len() > MAX_KIP_BATCH_COMMANDS {
return Err(KipError::resource_exhausted(format!(
"batch of {} operations exceeds maximum {MAX_KIP_BATCH_COMMANDS}",
self.operations.len()
)));
}
validate_optional_non_empty(&self.request_id, "request_id", limits::REQUEST_ID)?;
validate_optional_non_empty(
&self.compatibility_profile,
"compatibility_profile",
limits::COMPATIBILITY_PROFILE,
)?;
validate_extensions(&self.extensions, "extensions")?;
validate_block(self.space.as_ref())?;
validate_block(self.execution.as_ref())?;
validate_block(self.read.as_ref())?;
validate_block(self.preconditions.as_ref())?;
validate_block(self.context.as_ref())?;
validate_block(self.options.as_ref())?;
if self.operations.len() > 1 && self.execution.is_none() {
return Err(KipError::invalid_request_envelope(
"a multi-operation request must declare execution.mode: independent, sequence \
or atomic โ operations[] is a batch, not a transaction",
));
}
let mut seen_ops: Vec<&str> = Vec::new();
for operation in &self.operations {
operation.validate()?;
if let Some(op_id) = &operation.op_id {
if seen_ops.contains(&op_id.as_str()) {
return Err(KipError::invalid_request_envelope(format!(
"op_id {op_id:?} is used by two operations in one request"
)));
}
seen_ops.push(op_id);
}
}
if let Some(parameters) = &self.parameters {
for name in parameters.keys() {
validate_binding_name(name, "parameter")?;
}
}
if let Some(requires) = &self.requires {
for name in requires.keys() {
validate_capability_name(name)?;
}
}
if let Some(ingest) = &self.ingest {
ingest.validate()?;
let mut every_operation_parsed = true;
let mut opens_a_transaction = false;
for operation in &self.operations {
match operation.parse() {
Ok(command) if command.is_mutation() => {
opens_a_transaction = true;
break;
}
Ok(_) => {}
Err(_) => every_operation_parsed = false,
}
}
if every_operation_parsed && !opens_a_transaction {
return Err(KipError::invalid_request_envelope(
"an `ingest` block mints Evidence inside the request's transaction, so the \
request must carry at least one KML operation; a read-only request would \
drop the observation while reporting success",
));
}
}
Ok(())
}
pub fn critical_extensions(&self) -> Vec<&str> {
let blocks = [
self.extensions.as_ref(),
self.execution.as_ref().and_then(|e| e.extensions.as_ref()),
self.read.as_ref().and_then(|r| r.extensions.as_ref()),
self.preconditions
.as_ref()
.and_then(|p| p.extensions.as_ref()),
self.context.as_ref().and_then(|c| c.extensions.as_ref()),
self.options.as_ref().and_then(|o| o.extensions.as_ref()),
self.ingest.as_ref().and_then(|i| i.extensions.as_ref()),
];
let operation_blocks = self.operations.iter().flat_map(|operation| {
[
operation.extensions.as_ref(),
operation
.options
.as_ref()
.and_then(|o| o.extensions.as_ref()),
]
});
let ingest_blocks = self
.ingest
.iter()
.flat_map(|ingest| ingest.evidence.iter().map(|e| e.extensions.as_ref()));
let mut critical = Vec::new();
for block in blocks
.into_iter()
.chain(operation_blocks)
.chain(ingest_blocks)
.flatten()
{
for (name, value) in block {
if value.get("critical") == Some(&Json::Bool(true)) {
critical.push(name.as_str());
}
}
}
critical.sort_unstable();
critical.dedup();
critical
}
pub fn parse_operations(&self) -> Result<Vec<Command>, KipError> {
self.validate()?;
self.operations
.iter()
.map(|operation| operation.parse())
.collect()
}
}
#[derive(Clone, Debug, Default, Deserialize, Serialize, PartialEq)]
#[serde(deny_unknown_fields)]
pub struct SpaceSelector {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub uri: Option<String>,
}
impl EnvelopeBlock for SpaceSelector {
fn validate(&self) -> Result<(), KipError> {
if self.id.is_none() && self.uri.is_none() {
return Err(KipError::invalid_request_envelope(
"space must identify a MemorySpace by `id`, `uri`, or both",
));
}
validate_optional_non_empty(&self.id, "space.id", limits::SPACE_ID)?;
validate_optional_non_empty(&self.uri, "space.uri", limits::SPACE_URI)
}
}
#[derive(Clone, Debug, Deserialize, Serialize, PartialEq)]
#[serde(deny_unknown_fields)]
pub struct Execution {
pub mode: ExecutionMode,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub on_error: Option<OnError>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub isolation: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub idempotency_key: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<Map<String, Json>>,
}
impl Execution {
pub fn new(mode: ExecutionMode) -> Self {
Self {
mode,
on_error: None,
isolation: None,
idempotency_key: None,
extensions: None,
}
}
pub fn effective_on_error(&self) -> OnError {
self.on_error.unwrap_or_default()
}
}
impl EnvelopeBlock for Execution {
fn validate(&self) -> Result<(), KipError> {
if self.mode == ExecutionMode::Atomic && self.on_error == Some(OnError::Continue) {
return Err(KipError::invalid_request_envelope(
"an atomic transaction cannot continue past an error: it commits all or none",
));
}
validate_optional_non_empty(&self.isolation, "execution.isolation", limits::ISOLATION)?;
validate_optional_non_empty(
&self.idempotency_key,
"execution.idempotency_key",
limits::IDEMPOTENCY_KEY,
)?;
validate_extensions(&self.extensions, "execution.extensions")
}
}
#[derive(Clone, Copy, Debug, Default, Deserialize, Serialize, PartialEq, Eq, Hash)]
#[serde(rename_all = "lowercase")]
pub enum ExecutionMode {
#[default]
Independent,
Sequence,
Atomic,
}
impl ExecutionMode {
pub fn is_transactional(&self) -> bool {
matches!(self, ExecutionMode::Atomic)
}
}
#[derive(Clone, Copy, Debug, Default, Deserialize, Serialize, PartialEq, Eq, Hash)]
#[serde(rename_all = "lowercase")]
pub enum OnError {
#[default]
Stop,
Continue,
}
#[derive(Clone, Debug, Default, Deserialize, Serialize, PartialEq)]
#[serde(deny_unknown_fields)]
pub struct ReadBinding {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub snapshot_token: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<Map<String, Json>>,
}
impl EnvelopeBlock for ReadBinding {
fn validate(&self) -> Result<(), KipError> {
validate_optional_non_empty(
&self.snapshot_token,
"read.snapshot_token",
limits::OPAQUE_TOKEN,
)?;
validate_extensions(&self.extensions, "read.extensions")
}
}
#[derive(Clone, Debug, Default, Deserialize, Serialize, PartialEq)]
#[serde(deny_unknown_fields)]
pub struct Preconditions {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub space_seq: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub schema_environment_version: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<Map<String, Json>>,
}
impl EnvelopeBlock for Preconditions {
fn validate(&self) -> Result<(), KipError> {
validate_extensions(&self.extensions, "preconditions.extensions")
}
}
#[derive(Clone, Debug, Default, Deserialize, Serialize, PartialEq)]
#[serde(deny_unknown_fields)]
pub struct Operation {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub op_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub language: Option<CommandType>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub command: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub ast: Option<Command>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub parameters: Option<Map<String, Json>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub idempotency_key: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub options: Option<OperationOptions>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<Map<String, Json>>,
}
#[derive(Clone, Debug, Default, Deserialize, Serialize, PartialEq)]
#[serde(deny_unknown_fields)]
pub struct OperationOptions {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<Map<String, Json>>,
}
impl Operation {
pub fn new(command: impl Into<String>) -> Self {
Self {
command: Some(command.into()),
..Default::default()
}
}
pub fn with_op_id(mut self, op_id: impl Into<String>) -> Self {
self.op_id = Some(op_id.into());
self
}
pub fn with_parameters(mut self, parameters: Map<String, Json>) -> Self {
self.parameters = Some(parameters);
self
}
pub fn validate(&self) -> Result<(), KipError> {
match (&self.command, &self.ast) {
(Some(command), None) if !command.trim().is_empty() => {}
(Some(_), None) => {
return Err(KipError::invalid_request_envelope(
"an operation's command must not be empty",
));
}
(None, Some(_)) => {}
(Some(_), Some(_)) => {
return Err(KipError::invalid_request_envelope(
"an operation carries either `command` text or a pre-parsed `ast`, never both",
));
}
(None, None) => {
return Err(KipError::invalid_request_envelope(
"an operation must carry either `command` text or a pre-parsed `ast`",
));
}
}
validate_optional_non_empty(&self.op_id, "op_id", limits::REQUEST_ID)?;
validate_optional_non_empty(
&self.idempotency_key,
"operation.idempotency_key",
limits::IDEMPOTENCY_KEY,
)?;
validate_extensions(&self.extensions, "operation.extensions")?;
if let Some(options) = &self.options {
validate_extensions(&options.extensions, "operation.options.extensions")?;
}
if let Some(parameters) = &self.parameters {
for name in parameters.keys() {
validate_binding_name(name, "parameter")?;
}
}
Ok(())
}
pub fn parse(&self) -> Result<Command, KipError> {
self.validate()?;
let command = match (&self.command, &self.ast) {
(Some(text), _) => parse_kip(text)?,
(None, Some(ast)) => {
validate_command(ast)?;
ast.clone()
}
(None, None) => unreachable!("validate rejected the empty operation"),
};
if let Some(declared) = self.language {
let actual = CommandType::from(&command);
if declared != actual {
return Err(KipError::language_mismatch(format!(
"the operation declares {declared} but the command is {actual}"
)));
}
}
Ok(command)
}
}
#[derive(Clone, Debug, Default, Deserialize, Serialize, PartialEq)]
#[serde(deny_unknown_fields)]
pub struct RequestContext {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub purpose: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub risk: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub locale: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub client: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<Map<String, Json>>,
}
impl EnvelopeBlock for RequestContext {
fn validate(&self) -> Result<(), KipError> {
validate_optional_non_empty(&self.purpose, "context.purpose", limits::SHORT_LABEL)?;
validate_optional_non_empty(&self.risk, "context.risk", limits::RISK)?;
validate_optional_non_empty(&self.locale, "context.locale", limits::LOCALE)?;
validate_optional_non_empty(&self.client, "context.client", limits::SHORT_LABEL)?;
validate_extensions(&self.extensions, "context.extensions")
}
}
#[derive(Clone, Debug, Default, Deserialize, Serialize, PartialEq)]
#[serde(deny_unknown_fields)]
pub struct RequestOptions {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub dry_run: Option<bool>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub deadline_ms: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<Map<String, Json>>,
}
impl EnvelopeBlock for RequestOptions {
fn validate(&self) -> Result<(), KipError> {
validate_extensions(&self.extensions, "options.extensions")?;
if self.deadline_ms == Some(0) {
return Err(KipError::invalid_request_envelope(
"options.deadline_ms must be greater than zero",
));
}
Ok(())
}
}
#[derive(Clone, Debug, Default, Deserialize, Serialize, PartialEq)]
#[serde(deny_unknown_fields)]
pub struct IngestContext {
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub evidence: Vec<IngestEvidence>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<Map<String, Json>>,
}
impl IngestContext {
pub fn validate(&self) -> Result<(), KipError> {
if self.evidence.is_empty() {
return Err(KipError::invalid_request_envelope(
"an ingest context must carry at least one Evidence entry",
));
}
validate_extensions(&self.extensions, "ingest.extensions")?;
let mut seen: Vec<&str> = Vec::new();
for entry in &self.evidence {
entry.validate()?;
if seen.contains(&entry.key.as_str()) {
return Err(KipError::invalid_request_envelope(format!(
"ingest key {:?} is claimed by two Evidence entries",
entry.key
)));
}
seen.push(&entry.key);
}
Ok(())
}
}
#[derive(Clone, Debug, Default, Deserialize, Serialize, PartialEq)]
#[serde(deny_unknown_fields)]
pub struct ElementReference {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub id: Option<String>,
#[serde(default, rename = "type", skip_serializing_if = "Option::is_none")]
pub r#type: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub key: Option<String>,
}
impl ElementReference {
pub fn by_id(id: impl Into<String>) -> Self {
Self {
id: Some(id.into()),
r#type: None,
key: None,
}
}
pub fn by_key(r#type: impl Into<String>, key: impl Into<String>) -> Self {
Self {
id: None,
r#type: Some(r#type.into()),
key: Some(key.into()),
}
}
pub fn validate(&self, what: &str) -> Result<(), KipError> {
match (&self.id, &self.r#type, &self.key) {
(Some(id), None, None) => validate_bounded(id, what, limits::SOURCE_ACTOR),
(None, Some(r#type), Some(key)) => {
validate_bounded(r#type, what, limits::ELEMENT_REFERENCE_KEY)?;
validate_bounded(key, what, limits::ELEMENT_REFERENCE_KEY)
}
_ => Err(KipError::invalid_request_envelope(format!(
"{what} is an element reference: {{id}} or {{type, key}}, never a name"
))),
}
}
}
#[derive(Clone, Debug, Default, Deserialize, Serialize, PartialEq)]
#[serde(deny_unknown_fields)]
pub struct IngestEvidence {
pub key: String,
pub evidence_class: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub payload: Option<Json>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub payload_artifact: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub media_type: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub observed_at: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub source_actor: Option<ElementReference>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub client_key: Option<String>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub facets: BTreeMap<String, Map<String, Json>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<Map<String, Json>>,
}
impl IngestEvidence {
pub fn validate(&self) -> Result<(), KipError> {
validate_binding_name(&self.key, "ingest key")?;
if self.evidence_class.trim().is_empty() {
return Err(KipError::invalid_request_envelope(
"an ingest Evidence entry must declare an evidence_class",
));
}
validate_bounded(
&self.evidence_class,
"ingest evidence_class",
limits::SHORT_LABEL,
)?;
validate_optional_non_empty(
&self.payload_artifact,
"ingest payload_artifact",
limits::OPAQUE_TOKEN,
)?;
validate_optional_non_empty(&self.media_type, "ingest media_type", limits::SHORT_LABEL)?;
if let Some(at) = &self.observed_at {
crate::timestamp::parse(at, "ingest.observed_at")?;
}
if let Some(source_actor) = &self.source_actor {
source_actor.validate("ingest source_actor")?;
}
validate_optional_non_empty(&self.client_key, "ingest client_key", limits::CLIENT_KEY)?;
for facet in self.facets.keys() {
validate_bounded(facet, "ingest facets name", limits::FACET_NAME)?;
}
validate_extensions(&self.extensions, "ingest evidence extensions")?;
match (&self.payload, &self.payload_artifact) {
(Some(_), None) | (None, Some(_)) => Ok(()),
_ => Err(KipError::invalid_request_envelope(format!(
"ingest entry {:?} must declare exactly one of payload / payload_artifact",
self.key
))),
}
}
}
trait EnvelopeBlock {
fn validate(&self) -> Result<(), KipError>;
}
fn validate_block<T: EnvelopeBlock>(block: Option<&T>) -> Result<(), KipError> {
match block {
Some(block) => block.validate(),
None => Ok(()),
}
}
fn validate_optional_non_empty(
value: &Option<String>,
what: &str,
max: usize,
) -> Result<(), KipError> {
let Some(value) = value else { return Ok(()) };
validate_bounded(value, what, max)
}
fn validate_bounded(value: &str, what: &str, max: usize) -> Result<(), KipError> {
if value.trim().is_empty() {
return Err(KipError::invalid_request_envelope(format!(
"{what} must not be empty"
)));
}
let length = value.chars().count();
if length > max {
return Err(KipError::invalid_request_envelope(format!(
"{what} is {length} characters, and the wire schema allows at most {max}"
)));
}
Ok(())
}
fn validate_extensions(extensions: &Option<Map<String, Json>>, what: &str) -> Result<(), KipError> {
let Some(extensions) = extensions else {
return Ok(());
};
for name in extensions.keys() {
if !is_namespaced_extension(name) {
return Err(KipError::invalid_identifier(format!(
"{what} key {name:?} must be namespaced as <vendor>/<feature>"
)));
}
}
Ok(())
}
fn is_namespaced_extension(name: &str) -> bool {
let Some((namespace, rest)) = name.split_once('/') else {
return false;
};
let head_ok = |segment: &str| {
segment
.chars()
.next()
.is_some_and(|c| c.is_ascii_alphanumeric())
};
head_ok(namespace)
&& namespace
.chars()
.all(|c| c.is_ascii_alphanumeric() || matches!(c, '.' | '_' | '-'))
&& head_ok(rest)
&& rest
.chars()
.all(|c| c.is_ascii_alphanumeric() || matches!(c, '.' | '_' | '-' | '/'))
}
fn validate_capability_name(name: &str) -> Result<(), KipError> {
let valid = name.chars().next().is_some_and(|c| c.is_ascii_alphabetic())
&& name
.chars()
.all(|c| c.is_ascii_alphanumeric() || matches!(c, '_' | '.' | '-'));
if valid {
Ok(())
} else {
Err(KipError::invalid_identifier(format!(
"capability requirement {name:?} must match [A-Za-z][A-Za-z0-9_.-]*"
)))
}
}
fn validate_binding_name(name: &str, what: &str) -> Result<(), KipError> {
let mut chars = name.chars();
let valid = match chars.next() {
Some(c) if c.is_ascii_alphabetic() || c == '_' => {
chars.all(|c| c.is_ascii_alphanumeric() || c == '_')
}
_ => false,
};
if valid {
Ok(())
} else {
Err(KipError::invalid_identifier(format!(
"{what} {name:?} must match [A-Za-z_][A-Za-z0-9_]*"
)))
}
}
#[derive(Clone, Debug, Deserialize, Serialize, PartialEq)]
pub struct Response {
pub kip: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub request_id: Option<String>,
pub status: TopLevelStatus,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub execution: Option<ResponseExecution>,
pub results: Vec<OperationResult>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub context: Option<ResponseContext>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub snapshot: Option<SnapshotContext>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub receipt: Option<Receipt>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub warnings: Vec<Warning>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub next_cursor: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error: Option<ErrorObject>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<Map<String, Json>>,
}
impl Default for Response {
fn default() -> Self {
Self {
kip: KIP_VERSION.to_string(),
request_id: None,
status: TopLevelStatus::Succeeded,
execution: None,
results: Vec::new(),
context: None,
snapshot: None,
receipt: None,
warnings: Vec::new(),
next_cursor: None,
error: None,
extensions: None,
}
}
}
impl Response {
pub fn ok(result: Json) -> Self {
Self {
results: vec![OperationResult::ok(result)],
..Default::default()
}
}
pub fn failed(error: impl Into<ErrorObject>) -> Self {
let error = error.into();
Self {
status: TopLevelStatus::Failed,
results: vec![OperationResult::failed(error.clone())],
error: Some(error),
..Default::default()
}
}
pub fn from_results(results: Vec<OperationResult>) -> Self {
let status = TopLevelStatus::derive(&results);
Self {
status,
results,
..Default::default()
}
}
pub fn outcome_unknown(error: impl Into<ErrorObject>) -> Self {
Self {
status: TopLevelStatus::OutcomeUnknown,
error: Some(error.into()),
..Default::default()
}
}
pub fn with_request_id(mut self, request_id: Option<String>) -> Self {
self.request_id = request_id;
self
}
pub fn first_result(&self) -> Option<&Json> {
self.results.first().and_then(|r| r.result.as_ref())
}
}
impl From<KipError> for Response {
fn from(err: KipError) -> Self {
Response::failed(err)
}
}
#[derive(Clone, Copy, Debug, Default, Deserialize, Serialize, PartialEq, Eq, Hash)]
#[serde(rename_all = "snake_case")]
pub enum TopLevelStatus {
#[default]
Succeeded,
Failed,
Partial,
OutcomeUnknown,
}
impl TopLevelStatus {
pub fn derive(results: &[OperationResult]) -> Self {
if results.is_empty() {
return TopLevelStatus::Succeeded;
}
let succeeded = results
.iter()
.filter(|r| {
matches!(
r.status,
OperationStatus::Succeeded | OperationStatus::NoEffect
)
})
.count();
if succeeded == results.len() {
TopLevelStatus::Succeeded
} else if succeeded == 0 {
TopLevelStatus::Failed
} else {
TopLevelStatus::Partial
}
}
}
#[derive(Clone, Copy, Debug, Default, Deserialize, Serialize, PartialEq, Eq, Hash)]
#[serde(rename_all = "snake_case")]
pub enum OperationStatus {
#[default]
Succeeded,
Failed,
Skipped,
RolledBack,
NoEffect,
}
#[derive(Clone, Debug, Deserialize, Serialize, PartialEq)]
pub struct ResponseExecution {
pub mode: ExecutionMode,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub on_error: Option<OnError>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub isolation: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub idempotency_key: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<Map<String, Json>>,
}
#[derive(Clone, Debug, Default, Deserialize, Serialize, PartialEq)]
pub struct OperationResult {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub op_id: Option<String>,
pub status: OperationStatus,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub result: Option<Json>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub context: Option<ResultContext>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub warnings: Vec<Warning>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error: Option<ErrorObject>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub next_cursor: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub receipt: Option<Receipt>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<Map<String, Json>>,
}
impl OperationResult {
pub fn ok(result: Json) -> Self {
Self {
status: OperationStatus::Succeeded,
result: Some(result),
..Default::default()
}
}
pub fn failed(error: impl Into<ErrorObject>) -> Self {
Self {
status: OperationStatus::Failed,
error: Some(error.into()),
..Default::default()
}
}
pub fn no_effect() -> Self {
Self {
status: OperationStatus::NoEffect,
..Default::default()
}
}
pub fn rolled_back() -> Self {
Self {
status: OperationStatus::RolledBack,
..Default::default()
}
}
pub fn with_op_id(mut self, op_id: Option<String>) -> Self {
self.op_id = op_id;
self
}
}
#[derive(Clone, Debug, Default, Deserialize, Serialize, PartialEq)]
pub struct ResponseContext {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub space_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub schema_environment_version: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub compatibility_profile_used: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<Map<String, Json>>,
}
#[derive(Clone, Debug, Default, Deserialize, Serialize, PartialEq)]
pub struct ResultContext {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub space_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub snapshot_seq: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub schema_environment_version: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub epistemic_policy: Option<PolicyIdentity>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub valid_at: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub search: Option<SearchContext>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cursor: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<Map<String, Json>>,
}
#[derive(Clone, Debug, Deserialize, Serialize, PartialEq)]
#[serde(untagged)]
pub enum PolicyVersion {
Text(String),
Integer(u64),
}
#[derive(Clone, Debug, Deserialize, Serialize, PartialEq)]
pub struct PolicyIdentity {
pub id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub version: Option<PolicyVersion>,
}
impl PolicyIdentity {
pub fn new(id: impl Into<String>) -> Self {
Self {
id: id.into(),
version: None,
}
}
pub fn versioned(id: impl Into<String>, version: PolicyVersion) -> Self {
Self {
id: id.into(),
version: Some(version),
}
}
}
#[derive(Clone, Debug, Default, Deserialize, Serialize, PartialEq)]
pub struct SearchContext {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub index_seq: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub current_space_seq: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub consistency: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub mode: Option<SearchMode>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub score_semantics: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<Map<String, Json>>,
}
wire_enum! {
pub enum SearchMode {
Keyword = "keyword",
Semantic = "semantic",
Hybrid = "hybrid",
}
}
#[derive(Clone, Debug, Deserialize, Serialize, PartialEq)]
pub struct SnapshotContext {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub space_id: Option<String>,
pub snapshot_seq: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub schema_environment_version: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub snapshot_token: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<Map<String, Json>>,
}
impl SnapshotContext {
pub fn at(snapshot_seq: u64) -> Self {
Self {
space_id: None,
snapshot_seq,
schema_environment_version: None,
snapshot_token: None,
extensions: None,
}
}
}
#[derive(Clone, Debug, Deserialize, Serialize, PartialEq)]
pub struct Receipt {
pub status: ReceiptStatus,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tx_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub space_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub snapshot_seq: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub space_seq: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub committed_at: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub transaction_class: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub request_digest: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub semantic_plan_digest: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub result_digest: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub schema_environment_version: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub change_summary: Option<Json>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub proofs: Vec<Json>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub receipt_digest: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub origin: Option<ReceiptOrigin>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<Map<String, Json>>,
}
#[derive(Clone, Debug, Default, Deserialize, Serialize, PartialEq)]
pub struct ReceiptOrigin {
pub principal_id: String,
#[serde(default)]
pub actor_binding_id: Option<String>,
#[serde(default)]
pub delegation_digest: Option<String>,
}
#[derive(Clone, Copy, Debug, Deserialize, Serialize, PartialEq, Eq, Hash)]
#[serde(rename_all = "snake_case")]
pub enum ReceiptStatus {
Committed,
Aborted,
NoEffect,
Pending,
Unknown,
}
#[derive(Clone, Debug, Deserialize, Serialize, PartialEq)]
#[serde(untagged)]
pub enum Warning {
Message(String),
Coded {
code: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
message: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
details: Option<Json>,
#[serde(default, skip_serializing_if = "Option::is_none")]
extensions: Option<Map<String, Json>>,
},
}
impl From<&str> for Warning {
fn from(message: &str) -> Self {
Warning::Message(message.to_string())
}
}
impl From<String> for Warning {
fn from(message: String) -> Self {
Warning::Message(message)
}
}
impl From<KipErrorCode> for Warning {
fn from(code: KipErrorCode) -> Self {
Warning::Coded {
code: code.name().to_string(),
message: None,
details: None,
extensions: None,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::error::KipErrorCode;
#[test]
fn a_lone_operation_needs_no_execution_mode() {
let request = Request::single(r#"FIND(?x) WHERE { ?x {type: "T"} }"#);
assert!(request.validate().is_ok());
}
#[test]
fn unknown_wire_fields_are_rejected_instead_of_silently_ignored() {
let bad_request = serde_json::json!({
"kip": "2.0",
"operations": [{"command": "TRANSITION :x TO \"archived\""}],
"option": {"dry_run": true}
});
assert!(serde_json::from_value::<Request>(bad_request).is_err());
let bad_operation = serde_json::json!({
"kip": "2.0",
"operations": [{
"command": "DESCRIBE PROTOCOL",
"paramters": {"x": 1}
}]
});
assert!(serde_json::from_value::<Request>(bad_operation).is_err());
let bad_operation_option = serde_json::json!({
"kip": "2.0",
"operations": [{
"command": "TRANSITION :x TO \"archived\"",
"options": {"dry_run": true}
}]
});
assert!(serde_json::from_value::<Request>(bad_operation_option).is_err());
}
#[test]
fn optional_envelope_fields_must_not_be_empty_when_present() {
let mut request = Request {
space: Some(SpaceSelector::default()),
..Request::single("DESCRIBE PROTOCOL")
};
assert!(request.validate().is_err());
request.space = Some(SpaceSelector {
id: Some(" ".into()),
uri: None,
});
assert!(request.validate().is_err());
request.space = Some(SpaceSelector {
id: Some("space-1".into()),
uri: None,
});
request.options = Some(RequestOptions {
deadline_ms: Some(0),
..Default::default()
});
assert!(request.validate().is_err());
request.options = None;
request.context = Some(RequestContext {
purpose: Some(" ".into()),
..Default::default()
});
assert!(request.validate().is_err());
}
#[test]
fn policy_versions_accept_the_wire_schemas_text_and_integer_forms() {
let numeric: PolicyIdentity = serde_json::from_value(serde_json::json!({
"id": "projection-policy",
"version": 7
}))
.unwrap();
assert_eq!(numeric.version, Some(PolicyVersion::Integer(7)));
let textual: PolicyIdentity = serde_json::from_value(serde_json::json!({
"id": "projection-policy",
"version": "7.1"
}))
.unwrap();
assert_eq!(textual.version, Some(PolicyVersion::Text("7.1".into())));
assert!(
serde_json::from_value::<PolicyIdentity>(serde_json::json!({ "version": "7.1" }))
.is_err()
);
assert_eq!(
serde_json::to_value(PolicyIdentity::new("projection-policy")).unwrap(),
serde_json::json!({ "id": "projection-policy" })
);
}
#[test]
fn envelope_strings_are_held_to_the_wire_schemas_length_ceilings() {
let mut request = Request::single("DESCRIBE PROTOCOL");
request.request_id = Some("r".repeat(limits::REQUEST_ID));
request.validate().expect("exactly at the ceiling is legal");
request.request_id = Some("r".repeat(limits::REQUEST_ID + 1));
let err = request.validate().expect_err("one over the ceiling");
assert_eq!(err.code, KipErrorCode::InvalidRequestEnvelope);
assert!(err.message.contains("at most"), "{}", err.message);
request.request_id = Some("รฉ".repeat(limits::REQUEST_ID));
request.validate().expect("characters, not bytes");
}
#[test]
fn extension_keys_must_be_namespaced() {
let mut request = Request::single("DESCRIBE PROTOCOL");
request.extensions = serde_json::json!({ "acme/tracing": { "critical": false } })
.as_object()
.cloned();
request.validate().expect("a namespaced key is fine");
request.extensions = serde_json::json!({ "tracing": { "critical": true } })
.as_object()
.cloned();
let err = request.validate().expect_err("unnamespaced");
assert_eq!(err.code, KipErrorCode::InvalidIdentifier);
request.extensions = serde_json::json!({ "/tracing": {} }).as_object().cloned();
assert!(request.validate().is_err());
}
#[test]
fn critical_extensions_are_reported_from_every_block_that_carries_them() {
let mut request = Request::single("DESCRIBE PROTOCOL");
request.extensions = serde_json::json!({
"acme/redaction": { "critical": true },
"acme/tracing": { "critical": false }
})
.as_object()
.cloned();
request.operations[0].extensions =
serde_json::json!({ "acme/hints": { "critical": true } })
.as_object()
.cloned();
request.execution = Some(Execution {
mode: ExecutionMode::Independent,
on_error: None,
isolation: None,
idempotency_key: None,
extensions: serde_json::json!({ "acme/pinning": { "critical": true } })
.as_object()
.cloned(),
});
request.validate().expect("a legal envelope");
assert_eq!(
request.critical_extensions(),
vec!["acme/hints", "acme/pinning", "acme/redaction"]
);
assert!(
Request::single("DESCRIBE PROTOCOL")
.critical_extensions()
.is_empty()
);
}
#[test]
fn capability_requirements_are_named_like_capabilities() {
let mut request = Request::single("DESCRIBE PROTOCOL");
request.requires = serde_json::json!({ "belief_slot": true, "kip.streaming": true })
.as_object()
.cloned();
request.validate().expect("dotted identifiers are fine");
request.requires = serde_json::json!({ "!! nonsense": true })
.as_object()
.cloned();
let err = request.validate().expect_err("not an identifier");
assert_eq!(err.code, KipErrorCode::InvalidIdentifier);
}
#[test]
fn a_batch_must_say_how_its_operations_relate() {
let mut request = Request::single("DESCRIBE PROTOCOL");
request.operations.push(Operation::new("DESCRIBE PRIMER"));
let err = request.validate().expect_err("no execution mode");
assert_eq!(err.code, KipErrorCode::InvalidRequestEnvelope);
request.execution = Some(Execution::new(ExecutionMode::Independent));
assert!(request.validate().is_ok());
}
#[test]
fn an_atomic_transaction_cannot_continue_past_an_error() {
let mut request = Request::single(r#"TRANSITION :a TO "archived""#);
request
.operations
.push(Operation::new(r#"TRANSITION :b TO "archived""#));
request.execution = Some(Execution {
on_error: Some(OnError::Continue),
..Execution::new(ExecutionMode::Atomic)
});
assert!(request.validate().is_err());
request.execution = Some(Execution {
on_error: Some(OnError::Continue),
..Execution::new(ExecutionMode::Sequence)
});
assert!(request.validate().is_ok());
}
#[test]
fn a_declared_language_cannot_relabel_a_write_as_a_read() {
let operation = Operation {
language: Some(CommandType::Kql),
..Operation::new(r#"TRANSITION :x TO "tombstoned""#)
};
let err = operation.parse().expect_err("mislabelled write");
assert_eq!(err.code, KipErrorCode::LanguageMismatch);
let honest = Operation {
language: Some(CommandType::Kml),
..Operation::new(r#"TRANSITION :x TO "tombstoned""#)
};
assert!(honest.parse().unwrap().is_mutation());
}
#[test]
fn an_operation_carries_text_or_an_ast_but_never_both() {
let ast = parse_kip("DESCRIBE PROTOCOL").unwrap();
let both = Operation {
ast: Some(ast.clone()),
..Operation::new("DESCRIBE PROTOCOL")
};
assert!(both.validate().is_err());
let neither = Operation::default();
assert!(neither.validate().is_err());
let ast_only = Operation {
ast: Some(ast),
..Default::default()
};
assert_eq!(
ast_only.parse().unwrap(),
parse_kip("DESCRIBE PROTOCOL").unwrap()
);
}
#[test]
fn a_pre_parsed_ast_gets_the_same_guards_as_command_text() {
let text = r#"UPDATE ?a SET FIELDS { confidence: 0.1 } WHERE { ?a ASSERTION {id: "A-1"} }"#;
assert!(parse_kip(text).is_err(), "the text form must be rejected");
let rewrite_immutable_payload = serde_json::json!({"Kml": {
"explicit_transaction": false,
"clauses": [{"Update": {
"target": {"Handle": "a"},
"actions": [{"SetFields": [["confidence", {"Value": {"Number": 0.1}}]]}],
"where_clauses": [{"Assertion": {
"variable": "a",
"matcher": {"id": {"Literal": {"String": "A-1"}}}
}}],
"limit": Json::Null,
"expect_versions": []
}}]
}});
let write_engine_truth = serde_json::json!({"Kml": {
"explicit_transaction": false,
"clauses": [{"CreateConcept": {
"handle": "c", "type": Json::Null, "client_key": Json::Null, "name": Json::Null,
"set_fields": [["_system", {"Value": {"Number": 1}}]],
"set_attributes": Json::Null, "set_facets": [], "set_structural": Json::Null
}}]
}});
let unconfirmed_purge = serde_json::json!({"Kml": {
"explicit_transaction": false,
"clauses": [{"Purge": {
"target": {"Param": "x"}, "where_clauses": Json::Null, "limit": Json::Null,
"expect_versions": [], "reference_policy": Json::Null, "confirm": ""
}}]
}});
let belief_as_an_export_selector = serde_json::json!({"Meta": {"ExportCapsule": {
"target": {"Param": "out"},
"where_clauses": [{"Belief": {"variable": "b", "target": {"Proposition": "p"}}}],
"options": Json::Null, "as_of": Json::Null
}}});
for ast in [
rewrite_immutable_payload,
write_engine_truth,
unconfirmed_purge,
belief_as_an_export_selector,
] {
let operation = Operation {
ast: Some(serde_json::from_value(ast.clone()).expect("decodes")),
..Default::default()
};
let err = operation
.parse()
.expect_err(&format!("must be rejected: {ast}"));
assert_eq!(err.code, KipErrorCode::InvalidSyntax);
}
let honest = Operation {
ast: Some(parse_kip(r#"TRANSITION :old TO "archived""#).unwrap()),
..Default::default()
};
assert!(honest.parse().unwrap().is_mutation());
}
#[test]
fn op_ids_must_be_unique_within_a_request() {
let mut request = Request::single("DESCRIBE PROTOCOL");
request.operations[0].op_id = Some("op-1".into());
request
.operations
.push(Operation::new("DESCRIBE PRIMER").with_op_id("op-1"));
request.execution = Some(Execution::new(ExecutionMode::Independent));
assert!(request.validate().is_err());
}
#[test]
fn parameter_names_must_be_spellable_in_a_command() {
let mut request = Request::single("DESCRIBE PROTOCOL");
let mut parameters = Map::new();
parameters.insert("2bad".into(), Json::from(1));
request.parameters = Some(parameters);
let err = request.validate().expect_err("bad parameter name");
assert_eq!(err.code, KipErrorCode::InvalidIdentifier);
}
#[test]
fn ingest_entries_carry_exactly_one_payload() {
let base = IngestEvidence {
key: "msg".into(),
evidence_class: "user_statement".into(),
..Default::default()
};
assert!(base.validate().is_err());
let inline = IngestEvidence {
payload: Some(Json::from("I prefer dark mode.")),
..base.clone()
};
assert!(inline.validate().is_ok());
let both = IngestEvidence {
payload: Some(Json::from("x")),
payload_artifact: Some("artifact-1".into()),
..base.clone()
};
assert!(both.validate().is_err());
let duplicate = IngestContext {
evidence: vec![inline.clone(), inline],
extensions: None,
};
assert!(duplicate.validate().is_err());
}
#[test]
fn an_ingest_block_needs_a_transaction_to_be_minted_into() {
let ingest = IngestContext {
evidence: vec![IngestEvidence {
key: "msg".into(),
evidence_class: "user_statement".into(),
payload: Some(Json::from("I prefer dark mode.")),
..Default::default()
}],
extensions: None,
};
let mut read = Request::single(r#"FIND(?x) WHERE { ?x {type: "T"} }"#);
read.ingest = Some(ingest.clone());
let err = read
.validate()
.expect_err("a read has nothing to mint into");
assert_eq!(err.code, KipErrorCode::InvalidRequestEnvelope);
let mut write = Request::single(
r#"ASSERT (:alice, "prefers", :dark) { by: :alice, mode: "stated", evidence: :msg }"#,
);
write.ingest = Some(ingest.clone());
write.validate().expect("a KML operation carries the mint");
let mut mixed = Request::single("DESCRIBE PRIMER");
mixed
.operations
.push(Operation::new(r#"CREATE CONCEPT ?c { TYPE "T" NAME "n" }"#));
mixed.execution = Some(Execution {
mode: ExecutionMode::Sequence,
on_error: None,
isolation: None,
idempotency_key: None,
extensions: None,
});
mixed.ingest = Some(ingest);
mixed.validate().expect("one mutation is a transaction");
}
#[test]
fn a_partial_batch_is_not_a_failed_batch() {
let results = vec![
OperationResult::ok(Json::from(1)),
OperationResult::failed(KipError::not_found_or_not_visible("gone")),
];
assert_eq!(TopLevelStatus::derive(&results), TopLevelStatus::Partial);
let all_ok = vec![
OperationResult::ok(Json::Null),
OperationResult::no_effect(),
];
assert_eq!(TopLevelStatus::derive(&all_ok), TopLevelStatus::Succeeded);
let all_bad = vec![OperationResult::failed(KipError::internal_error("boom"))];
assert_eq!(TopLevelStatus::derive(&all_bad), TopLevelStatus::Failed);
}
#[test]
fn rolled_back_is_not_success() {
let results = vec![OperationResult::rolled_back()];
assert_eq!(TopLevelStatus::derive(&results), TopLevelStatus::Failed);
}
#[test]
fn the_envelope_round_trips_through_its_wire_shape() {
let request = Request {
request_id: Some("req-1".into()),
space: Some(SpaceSelector {
id: Some("space-1".into()),
uri: None,
}),
execution: Some(Execution {
idempotency_key: Some("logical-write-key".into()),
isolation: Some("serializable".into()),
..Execution::new(ExecutionMode::Atomic)
}),
operations: vec![Operation::new(r#"TRANSITION :x TO "archived""#).with_op_id("op-1")],
options: Some(RequestOptions {
deadline_ms: Some(10_000),
..Default::default()
}),
..Default::default()
};
let json = serde_json::to_value(&request).unwrap();
assert_eq!(json["kip"], "2.0");
assert_eq!(json["execution"]["mode"], "atomic");
assert_eq!(json["operations"][0]["op_id"], "op-1");
let decoded: Request = serde_json::from_value(json).unwrap();
assert_eq!(decoded, request);
let response = Response {
receipt: Some(Receipt {
status: ReceiptStatus::Committed,
tx_id: Some("tx-9".into()),
space_seq: Some(4201),
snapshot_seq: Some(4200),
space_id: Some("space-1".into()),
committed_at: Some("2026-08-16T00:00:00.000Z".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,
}),
warnings: vec![Warning::Message("search index lagged".into())],
..Response::ok(Json::from(true))
};
let json = serde_json::to_value(&response).unwrap();
assert_eq!(json["status"], "succeeded");
assert_eq!(json["receipt"]["status"], "committed");
assert_eq!(json["warnings"][0], "search index lagged");
let decoded: Response = serde_json::from_value(json).unwrap();
assert_eq!(decoded, response);
}
#[test]
fn an_error_response_carries_the_registry_shape() {
let response = Response::from(KipError::version_conflict("element changed"));
assert_eq!(response.status, TopLevelStatus::Failed);
let json = serde_json::to_value(&response).unwrap();
assert_eq!(json["error"]["code"], "VersionConflict");
assert_eq!(json["error"]["retry"]["class"], "requires_refresh");
assert_eq!(json["results"][0]["status"], "failed");
}
#[test]
fn a_lost_response_is_not_a_failed_write() {
let response = Response::outcome_unknown(KipError::outcome_unknown("connection dropped"));
assert_eq!(response.status, TopLevelStatus::OutcomeUnknown);
assert_eq!(
response.error.unwrap().retry.unwrap().class,
crate::error::RetryClass::OutcomeLookupRequired
);
}
#[test]
fn an_over_sized_batch_is_rejected_before_execution() {
let mut request = Request::single("DESCRIBE PROTOCOL");
request.execution = Some(Execution::new(ExecutionMode::Independent));
for _ in 0..MAX_KIP_BATCH_COMMANDS {
request.operations.push(Operation::new("DESCRIBE PROTOCOL"));
}
let err = request.validate().expect_err("too many operations");
assert_eq!(err.code, KipErrorCode::ResourceExhausted);
}
#[test]
fn an_unknown_protocol_version_fails_fast() {
let request = Request {
kip: "1.0".into(),
..Request::single("DESCRIBE PROTOCOL")
};
let err = request.validate().expect_err("wrong version");
assert_eq!(err.code, KipErrorCode::UnsupportedProtocolVersion);
}
}