use crate::actions;
use crate::errors::{DynoxideError, UNSUPPORTED_TYPE};
use crate::storage_backend::StorageBackend;
pub const CONTRACT_VERSION: u32 = 1;
pub const SUPPORTED_OPS: &[&str] = &[
"CreateTable",
"DeleteTable",
"DescribeTable",
"UpdateTable",
"ListTables",
"PutItem",
"GetItem",
"DeleteItem",
"UpdateItem",
"Query",
"Scan",
"BatchGetItem",
"BatchWriteItem",
"TransactGetItems",
"ExecuteStatement",
"BatchExecuteStatement",
"ExecuteTransaction",
];
pub struct DispatchContext<'a> {
tokens: &'a crate::TokenCaches,
}
impl<'a> DispatchContext<'a> {
pub fn new(tokens: &'a crate::TokenCaches) -> Self {
Self { tokens }
}
}
const TARGET_PREFIX: &str = "DynamoDB_20120810.";
const STREAMS_TARGET_PREFIX: &str = "DynamoDBStreams_20120810.";
const SERIALIZATION_EXCEPTION_BARE: &str =
r#"{"__type":"com.amazon.coral.service#SerializationException"}"#;
const UNKNOWN_OPERATION_BARE: &str =
r#"{"__type":"com.amazon.coral.service#UnknownOperationException"}"#;
pub struct HttpOutcome {
pub status: u16,
pub body: String,
}
impl HttpOutcome {
fn new(status: u16, body: impl Into<String>) -> Self {
Self {
status,
body: body.into(),
}
}
}
pub async fn dispatch_http<S: StorageBackend>(
backend: &S,
ctx: &DispatchContext<'_>,
target: Option<&str>,
body: &str,
auth: crate::auth_material::AuthMaterial<'_>,
) -> HttpOutcome {
if !body.is_empty() && serde_json::from_str::<serde_json::Value>(body).is_err() {
return HttpOutcome::new(400, SERIALIZATION_EXCEPTION_BARE);
}
let Some(target) = target else {
return HttpOutcome::new(400, UNKNOWN_OPERATION_BARE);
};
let operation = target
.strip_prefix(TARGET_PREFIX)
.or_else(|| target.strip_prefix(STREAMS_TARGET_PREFIX));
let Some(operation) = operation.filter(|op| crate::dynamo_ops::is_known_operation(op)) else {
return HttpOutcome::new(400, UNKNOWN_OPERATION_BARE);
};
if let Some(envelope) = crate::auth_material::validate(auth) {
return HttpOutcome::new(400, envelope);
}
if body.is_empty() {
return HttpOutcome::new(400, SERIALIZATION_EXCEPTION_BARE);
}
if !SUPPORTED_OPS.contains(&operation) {
return HttpOutcome::new(501, unsupported_envelope(operation));
}
match route(backend, ctx, operation, body).await {
Ok(json) => HttpOutcome::new(200, json),
Err(e) => HttpOutcome::new(e.status_code(), e.to_json()),
}
}
fn unsupported_envelope(op: &str) -> String {
let message = if crate::dynamo_ops::is_known_operation(op) {
format!("Operation '{op}' is not supported by the wasm preview engine")
} else {
format!("Unknown operation: '{op}'")
};
serde_json::json!({ "__type": UNSUPPORTED_TYPE, "message": message }).to_string()
}
pub async fn dispatch<S: StorageBackend>(
backend: &S,
ctx: &DispatchContext<'_>,
op: &str,
request_json: &str,
) -> std::result::Result<String, String> {
if !SUPPORTED_OPS.contains(&op) {
return Err(unsupported_envelope(op));
}
route(backend, ctx, op, request_json)
.await
.map_err(|e| e.to_json())
}
async fn route<S: StorageBackend>(
backend: &S,
ctx: &DispatchContext<'_>,
op: &str,
request_json: &str,
) -> crate::Result<String> {
macro_rules! run {
($module:ident) => {{
match crate::serde_errors::deserialize(request_json) {
Ok(request) => match actions::$module::execute(backend, request).await {
Ok(response) => serde_json::to_string(&response)
.map_err(|e| DynoxideError::InternalServerError(e.to_string())),
Err(e) => Err(e),
},
Err(e) => Err(e),
}
}};
}
let result: crate::Result<String> = match op {
"CreateTable" => run!(create_table),
"DeleteTable" => run!(delete_table),
"DescribeTable" => run!(describe_table),
"UpdateTable" => run!(update_table),
"ListTables" => run!(list_tables),
"PutItem" => run!(put_item),
"GetItem" => run!(get_item),
"DeleteItem" => run!(delete_item),
"UpdateItem" => run!(update_item),
"Query" => run!(query),
"Scan" => run!(scan),
"BatchGetItem" => run!(batch_get_item),
"BatchWriteItem" => run!(batch_write_item),
"TransactGetItems" => run!(transact_get_items),
"ExecuteStatement" => run!(execute_statement),
"BatchExecuteStatement" => run!(batch_execute_statement),
"ExecuteTransaction" => execute_transaction(backend, ctx, request_json).await,
other => Err(DynoxideError::InternalServerError(format!(
"route reached with operation '{other}', which is absent from SUPPORTED_OPS"
))),
};
result.map_err(|e| crate::validation::resolve_request_validation_tag(op, e))
}
async fn execute_transaction<S: StorageBackend>(
backend: &S,
ctx: &DispatchContext<'_>,
request_json: &str,
) -> crate::Result<String> {
let request: actions::execute_transaction::ExecuteTransactionRequest =
crate::serde_errors::deserialize(request_json)?;
let statements = request.transact_statements.clone();
let token = request.client_request_token.clone();
let capacity_mode = request.return_consumed_capacity.clone();
let response = crate::run_idempotent_async(
ctx.tokens.execute_transaction(),
token.as_deref(),
&statements,
actions::execute_transaction::execute(backend, request),
|cached| {
actions::execute_transaction::replay_response(
&statements,
&capacity_mode,
cached.responses.clone(),
)
},
)
.await?;
serde_json::to_string(&response).map_err(|e| DynoxideError::InternalServerError(e.to_string()))
}
#[cfg(feature = "wasm-sqlite")]
mod engine {
use super::{CONTRACT_VERSION, DispatchContext, SUPPORTED_OPS, dispatch};
use crate::WasmDatabase;
use std::cell::RefCell;
use wasm_bindgen::prelude::*;
thread_local! {
static ENGINE: RefCell<Option<WasmDatabase>> = const { RefCell::new(None) };
}
fn boot_descriptor(persistence_mode: &str) -> String {
serde_json::json!({
"contractVersion": CONTRACT_VERSION,
"capabilities": SUPPORTED_OPS,
"persistenceMode": persistence_mode,
})
.to_string()
}
fn not_opened_envelope() -> String {
serde_json::json!({
"__type": "com.dynoxide.wasm#EngineNotOpened",
"message": "execute called before open(); call open(name) first",
})
.to_string()
}
#[wasm_bindgen]
pub async fn open(name: String, ephemeral: bool) -> Result<String, String> {
let db = WasmDatabase::open_with(&name, ephemeral)
.await
.map_err(|e| e.to_json())?;
let persistence_mode = db.persistence_mode().await;
let previous = ENGINE.with(|cell| cell.borrow_mut().replace(db));
if let Some(previous) = previous {
let _ = previous.close().await;
}
Ok(boot_descriptor(&persistence_mode))
}
#[wasm_bindgen]
pub async fn execute(op: String, request_json: String) -> Result<String, String> {
let db = ENGINE.with(|cell| cell.borrow().clone());
let Some(db) = db else {
return Err(not_opened_envelope());
};
let backend = db.backend().await;
dispatch(
&*backend,
&DispatchContext::new(db.token_caches()),
&op,
&request_json,
)
.await
}
#[wasm_bindgen(js_name = dispatchHttp)]
pub async fn dispatch_http_js(
target: Option<String>,
body: String,
authorization: Option<String>,
query: Option<String>,
has_date_header: bool,
) -> Result<String, String> {
let db = ENGINE.with(|cell| cell.borrow().clone());
let Some(db) = db else {
return Err(not_opened_envelope());
};
let backend = db.backend().await;
let auth = crate::auth_material::AuthMaterial {
authorization: authorization.as_deref(),
query: query.as_deref().unwrap_or(""),
has_date_header,
};
let outcome = super::dispatch_http(
&*backend,
&DispatchContext::new(db.token_caches()),
target.as_deref(),
&body,
auth,
)
.await;
Ok(serde_json::json!({
"status": outcome.status,
"body": outcome.body,
})
.to_string())
}
#[wasm_bindgen]
pub fn capabilities() -> String {
serde_json::to_string(SUPPORTED_OPS).unwrap_or_else(|_| "[]".to_string())
}
#[wasm_bindgen]
pub fn contract_version() -> u32 {
CONTRACT_VERSION
}
}
#[cfg(all(test, feature = "native-sqlite"))]
mod tests {
use super::*;
use crate::storage::Storage;
fn run(backend: &Storage, op: &str, json: &str) -> std::result::Result<String, String> {
let tokens = crate::TokenCaches::new();
pollster::block_on(dispatch(backend, &DispatchContext::new(&tokens), op, json))
}
fn run_with(
backend: &Storage,
tokens: &crate::TokenCaches,
op: &str,
json: &str,
) -> std::result::Result<String, String> {
pollster::block_on(dispatch(backend, &DispatchContext::new(tokens), op, json))
}
const CREATE_MUSIC: &str = r#"{
"TableName": "Music",
"KeySchema": [
{"AttributeName": "artist", "KeyType": "HASH"},
{"AttributeName": "song", "KeyType": "RANGE"}
],
"AttributeDefinitions": [
{"AttributeName": "artist", "AttributeType": "S"},
{"AttributeName": "song", "AttributeType": "S"}
],
"BillingMode": "PAY_PER_REQUEST"
}"#;
fn seed_music(backend: &Storage) {
run(backend, "CreateTable", CREATE_MUSIC).expect("create table");
for (song, genre) in [("s1", "rock"), ("s2", "jazz"), ("s3", "rock")] {
let put = format!(
r#"{{"TableName":"Music","Item":{{"artist":{{"S":"a"}},"song":{{"S":"{song}"}},"genre":{{"S":"{genre}"}}}}}}"#
);
run(backend, "PutItem", &put).expect("put item");
}
}
#[test]
fn create_put_get_roundtrip() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
let put = r#"{"TableName":"Music","Item":{"artist":{"S":"a"},"song":{"S":"s1"},"msg":{"S":"hi"}}}"#;
run(&backend, "PutItem", put).unwrap();
let get = r#"{"TableName":"Music","Key":{"artist":{"S":"a"},"song":{"S":"s1"}}}"#;
let resp = run(&backend, "GetItem", get).unwrap();
let v: serde_json::Value = serde_json::from_str(&resp).unwrap();
assert_eq!(v["Item"]["msg"]["S"], "hi");
}
#[test]
fn query_returns_count_and_items() {
let backend = Storage::memory().unwrap();
seed_music(&backend);
let query = r#"{"TableName":"Music","KeyConditionExpression":"artist = :a","ExpressionAttributeValues":{":a":{"S":"a"}}}"#;
let resp = run(&backend, "Query", query).unwrap();
let v: serde_json::Value = serde_json::from_str(&resp).unwrap();
assert_eq!(v["Count"], 3);
assert_eq!(v["Items"].as_array().unwrap().len(), 3);
}
#[test]
fn scan_with_filter_scans_more_than_it_counts() {
let backend = Storage::memory().unwrap();
seed_music(&backend);
let scan = r#"{"TableName":"Music","FilterExpression":"genre = :g","ExpressionAttributeValues":{":g":{"S":"rock"}}}"#;
let resp = run(&backend, "Scan", scan).unwrap();
let v: serde_json::Value = serde_json::from_str(&resp).unwrap();
assert_eq!(v["Count"], 2);
assert_eq!(v["ScannedCount"], 3);
assert!(v["ScannedCount"].as_u64() > v["Count"].as_u64());
}
#[test]
fn newly_wrapped_update_item_roundtrips() {
let backend = Storage::memory().unwrap();
seed_music(&backend);
let update = r#"{
"TableName": "Music",
"Key": {"artist": {"S": "a"}, "song": {"S": "s1"}},
"UpdateExpression": "SET plays = :p",
"ExpressionAttributeValues": {":p": {"N": "5"}},
"ReturnValues": "ALL_NEW"
}"#;
let resp = run(&backend, "UpdateItem", update).unwrap();
let v: serde_json::Value = serde_json::from_str(&resp).unwrap();
assert_eq!(v["Attributes"]["plays"]["N"], "5");
}
#[test]
fn batch_get_item_returns_seeded_items() {
let backend = Storage::memory().unwrap();
seed_music(&backend);
let batch_get = r#"{
"RequestItems": {
"Music": {
"Keys": [
{"artist": {"S": "a"}, "song": {"S": "s1"}},
{"artist": {"S": "a"}, "song": {"S": "s3"}}
]
}
}
}"#;
let resp = run(&backend, "BatchGetItem", batch_get).unwrap();
let v: serde_json::Value = serde_json::from_str(&resp).unwrap();
let items = v["Responses"]["Music"].as_array().unwrap();
assert_eq!(items.len(), 2);
let mut songs: Vec<&str> = items
.iter()
.map(|item| item["song"]["S"].as_str().unwrap())
.collect();
songs.sort_unstable();
assert_eq!(songs, ["s1", "s3"]);
assert!(v["UnprocessedKeys"].as_object().unwrap().is_empty());
}
#[test]
fn batch_write_item_puts_and_deletes_persist() {
let backend = Storage::memory().unwrap();
seed_music(&backend);
let batch_write = r#"{
"RequestItems": {
"Music": [
{"DeleteRequest": {"Key": {"artist": {"S": "a"}, "song": {"S": "s2"}}}},
{"PutRequest": {"Item": {"artist": {"S": "a"}, "song": {"S": "s4"}, "genre": {"S": "pop"}}}}
]
}
}"#;
let resp = run(&backend, "BatchWriteItem", batch_write).unwrap();
let v: serde_json::Value = serde_json::from_str(&resp).unwrap();
assert!(v["UnprocessedItems"].as_object().unwrap().is_empty());
let get_s2 = r#"{"TableName":"Music","Key":{"artist":{"S":"a"},"song":{"S":"s2"}}}"#;
let s2: serde_json::Value =
serde_json::from_str(&run(&backend, "GetItem", get_s2).unwrap()).unwrap();
assert!(s2.get("Item").is_none(), "s2 should have been deleted");
let get_s4 = r#"{"TableName":"Music","Key":{"artist":{"S":"a"},"song":{"S":"s4"}}}"#;
let s4: serde_json::Value =
serde_json::from_str(&run(&backend, "GetItem", get_s4).unwrap()).unwrap();
assert_eq!(s4["Item"]["genre"]["S"], "pop");
}
#[test]
fn transact_get_items_preserves_position_for_present_and_missing() {
let backend = Storage::memory().unwrap();
seed_music(&backend);
let transact_get = r#"{
"TransactItems": [
{"Get": {"TableName": "Music", "Key": {"artist": {"S": "a"}, "song": {"S": "s1"}}}},
{"Get": {"TableName": "Music", "Key": {"artist": {"S": "a"}, "song": {"S": "nope"}}}}
]
}"#;
let resp = run(&backend, "TransactGetItems", transact_get).unwrap();
let v: serde_json::Value = serde_json::from_str(&resp).unwrap();
let responses = v["Responses"].as_array().unwrap();
assert_eq!(responses.len(), 2);
assert_eq!(responses[0]["Item"]["genre"]["S"], "rock");
assert!(responses[1].get("Item").is_none());
}
#[test]
fn unknown_op_returns_envelope_not_panic() {
let backend = Storage::memory().unwrap();
let err = run(&backend, "FlyToTheMoon", "{}").unwrap_err();
let v: serde_json::Value = serde_json::from_str(&err).unwrap();
assert_eq!(v["__type"], "com.dynoxide.wasm#UnsupportedOperation");
assert!(v["message"].as_str().unwrap().contains("Unknown operation"));
}
#[test]
fn unsupported_preview_op_returns_envelope() {
let backend = Storage::memory().unwrap();
let err = run(&backend, "UpdateTimeToLive", "{}").unwrap_err();
let v: serde_json::Value = serde_json::from_str(&err).unwrap();
assert_eq!(v["__type"], "com.dynoxide.wasm#UnsupportedOperation");
assert!(v["message"].as_str().unwrap().contains("not supported"));
}
#[test]
fn conditional_check_failure_surfaces_in_envelope() {
let backend = Storage::memory().unwrap();
seed_music(&backend);
let put = r#"{
"TableName": "Music",
"Item": {"artist": {"S": "a"}, "song": {"S": "s1"}},
"ConditionExpression": "attribute_not_exists(artist)"
}"#;
let err = run(&backend, "PutItem", put).unwrap_err();
let v: serde_json::Value = serde_json::from_str(&err).unwrap();
assert!(
v["__type"]
.as_str()
.unwrap()
.contains("ConditionalCheckFailedException")
);
}
#[test]
fn malformed_request_json_is_a_serialization_error() {
let backend = Storage::memory().unwrap();
let err = run(&backend, "PutItem", "{ this is not json").unwrap_err();
let v: serde_json::Value = serde_json::from_str(&err).unwrap();
assert!(
v["__type"]
.as_str()
.unwrap()
.contains("SerializationException")
);
}
#[test]
fn contract_advertises_a_version_and_the_supported_ops() {
assert_eq!(CONTRACT_VERSION, 1);
assert!(SUPPORTED_OPS.contains(&"Query"));
assert!(SUPPORTED_OPS.contains(&"Scan"));
}
const NULL_FALSE_ENVELOPED: &str = "1 validation error detected: \
One or more parameter values were invalid: \
Null attribute value types must have the value of true";
const NULL_FALSE_BARE: &str = "One or more parameter values were invalid: \
Null attribute value types must have the value of true";
fn assert_validation_payload(err: &str, expected_message: &str) {
assert!(
!err.contains("VALIDATION") && !err.contains(" at line "),
"internal marker or serde position leaked: {err}"
);
let v: serde_json::Value = serde_json::from_str(err).unwrap();
assert!(
v["__type"]
.as_str()
.unwrap()
.ends_with("ValidationException"),
"unexpected __type: {}",
v["__type"]
);
assert_eq!(v["message"].as_str().unwrap(), expected_message);
}
#[test]
fn put_item_null_false_in_item_is_enveloped_validation() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
let put = r#"{"TableName":"Music","Item":{"artist":{"S":"a"},"song":{"S":"s1"},"flag":{"NULL":false}}}"#;
let err = run(&backend, "PutItem", put).unwrap_err();
assert_validation_payload(&err, NULL_FALSE_ENVELOPED);
}
#[test]
fn get_item_null_false_in_key_is_bare_validation() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
let get = r#"{"TableName":"Music","Key":{"artist":{"NULL":false},"song":{"S":"s1"}}}"#;
let err = run(&backend, "GetItem", get).unwrap_err();
assert_validation_payload(&err, NULL_FALSE_BARE);
}
#[test]
fn missing_field_classifies_as_validation_like_http() {
let backend = Storage::memory().unwrap();
let err = run(&backend, "BatchGetItem", "{}").unwrap_err();
let v: serde_json::Value = serde_json::from_str(&err).unwrap();
assert_eq!(v["__type"], "com.amazon.coral.validate#ValidationException");
}
#[test]
fn error_payloads_never_leak_markers_or_positions() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
let cases: &[(&str, &str)] = &[
(
"PutItem",
r#"{"TableName":"Music","Item":{"artist":{"S":"a"},"song":{"S":"s"},"flag":{"NULL":false}}}"#,
),
(
"UpdateItem",
r#"{"TableName":"Music","Key":{"artist":{"NULL":false},"song":{"S":"s"}},"UpdateExpression":"SET x = :v","ExpressionAttributeValues":{":v":{"S":"v"}}}"#,
),
(
"GetItem",
r#"{"TableName":"Music","Key":{"artist":{"NULL":false},"song":{"S":"s"}}}"#,
),
(
"DeleteItem",
r#"{"TableName":"Music","Key":{"artist":{"NULL":false},"song":{"S":"s"}}}"#,
),
(
"Query",
r#"{"TableName":"Music","KeyConditionExpression":"artist = :a","ExpressionAttributeValues":{":a":{"NULL":false}}}"#,
),
("DeleteTable", "{}"),
];
for (op, body) in cases {
let err = run(&backend, op, body).unwrap_err();
assert!(
!err.contains("VALIDATION") && !err.contains(" at line "),
"{op}: internal marker or serde position leaked: {err}"
);
}
}
#[test]
fn update_table_adds_a_gsi_and_backfills_through_dispatch() {
let backend = Storage::memory().unwrap();
seed_music(&backend);
let update = r#"{
"TableName": "Music",
"AttributeDefinitions": [
{"AttributeName": "artist", "AttributeType": "S"},
{"AttributeName": "song", "AttributeType": "S"},
{"AttributeName": "genre", "AttributeType": "S"}
],
"GlobalSecondaryIndexUpdates": [
{"Create": {
"IndexName": "GenreIndex",
"KeySchema": [{"AttributeName": "genre", "KeyType": "HASH"}],
"Projection": {"ProjectionType": "ALL"}
}}
]
}"#;
let resp = run(&backend, "UpdateTable", update).unwrap();
assert!(
resp.contains("GenreIndex"),
"the response should describe the new GSI"
);
let q = r#"{"TableName":"Music","IndexName":"GenreIndex","KeyConditionExpression":"genre = :g","ExpressionAttributeValues":{":g":{"S":"rock"}}}"#;
let qv: serde_json::Value =
serde_json::from_str(&run(&backend, "Query", q).unwrap()).unwrap();
assert_eq!(qv["Count"], 2);
assert!(SUPPORTED_OPS.contains(&"UpdateTable"));
}
fn is_unsupported_fault(status: u16, body: &str) -> bool {
let v: serde_json::Value = serde_json::from_str(body).unwrap_or(serde_json::Value::Null);
let type_tail = v["__type"]
.as_str()
.unwrap_or("")
.rsplit('#')
.next()
.unwrap_or("")
.to_owned();
let message = v["message"].as_str().unwrap_or("").to_lowercase();
status == 501
|| type_tail == "UnknownOperationException"
|| [
"unknown operation",
"not implemented",
"unsupported operation",
"is not supported",
]
.iter()
.any(|needle| message.contains(needle))
}
const SIGNED: &str = "AWS4-HMAC-SHA256 Credential=fake/20260724/eu-west-2/dynamodb/aws4_request, SignedHeaders=host;x-amz-date, Signature=abc";
fn signed_auth() -> crate::auth_material::AuthMaterial<'static> {
crate::auth_material::AuthMaterial {
authorization: Some(SIGNED),
query: "",
has_date_header: true,
}
}
fn http(backend: &Storage, target: Option<&str>, body: &str) -> HttpOutcome {
let tokens = crate::TokenCaches::new();
http_with(backend, &tokens, target, body)
}
fn http_with(
backend: &Storage,
tokens: &crate::TokenCaches,
target: Option<&str>,
body: &str,
) -> HttpOutcome {
pollster::block_on(dispatch_http(
backend,
&DispatchContext::new(tokens),
target,
body,
signed_auth(),
))
}
#[test]
fn http_roundtrips_a_supported_operation() {
let backend = Storage::memory().unwrap();
let out = http(
&backend,
Some("DynamoDB_20120810.CreateTable"),
CREATE_MUSIC,
);
assert_eq!(out.status, 200);
let out = http(&backend, Some("DynamoDB_20120810.ListTables"), "{}");
assert_eq!(out.status, 200);
let v: serde_json::Value = serde_json::from_str(&out.body).unwrap();
assert_eq!(v["TableNames"][0], "Music");
}
#[test]
fn http_rejects_a_non_json_body_as_bare_serialization_exception() {
let backend = Storage::memory().unwrap();
let out = http(&backend, Some("DynamoDB_20120810.ListTables"), "not json");
assert_eq!(out.status, 400);
assert_eq!(out.body, SERIALIZATION_EXCEPTION_BARE);
}
#[test]
fn http_reports_a_missing_target_before_an_empty_body() {
let backend = Storage::memory().unwrap();
let out = http(&backend, None, "");
assert_eq!(out.status, 400);
assert_eq!(out.body, UNKNOWN_OPERATION_BARE);
}
#[test]
fn http_rejects_an_empty_body_on_a_valid_target() {
let backend = Storage::memory().unwrap();
let out = http(&backend, Some("DynamoDB_20120810.ListTables"), "");
assert_eq!(out.status, 400);
assert_eq!(out.body, SERIALIZATION_EXCEPTION_BARE);
}
#[test]
fn http_rejects_an_unrecognised_target_prefix() {
let backend = Storage::memory().unwrap();
let out = http(&backend, Some("Wrong_20120810.ListTables"), "{}");
assert_eq!(out.status, 400);
assert_eq!(out.body, UNKNOWN_OPERATION_BARE);
}
#[test]
fn http_accepts_the_streams_target_prefix() {
let backend = Storage::memory().unwrap();
let out = http(&backend, Some("DynamoDBStreams_20120810.ListStreams"), "{}");
assert_eq!(out.status, 501);
assert!(is_unsupported_fault(out.status, &out.body));
}
#[test]
fn http_classifies_an_unknown_operation_as_a_skip() {
let backend = Storage::memory().unwrap();
let out = http(&backend, Some("DynamoDB_20120810.NoSuchOp"), "{}");
assert_eq!(out.status, 400);
assert!(is_unsupported_fault(out.status, &out.body));
}
#[test]
fn http_classifies_every_unimplemented_operation_as_a_skip() {
let backend = Storage::memory().unwrap();
for op in [
"UpdateTimeToLive",
"DescribeTimeToLive",
"TransactWriteItems",
"TagResource",
"UntagResource",
"ListTagsOfResource",
"DescribeLimits",
"ListStreams",
] {
let target = format!("DynamoDB_20120810.{op}");
let out = http(&backend, Some(&target), "{}");
assert!(
is_unsupported_fault(out.status, &out.body),
"{op} would be scored as a conformance failure: {} {}",
out.status,
out.body
);
}
}
#[test]
fn http_validates_auth_material_the_same_way_the_native_server_does() {
let backend = Storage::memory().unwrap();
let unsigned = crate::auth_material::AuthMaterial::default();
let tokens = crate::TokenCaches::new();
let out = pollster::block_on(dispatch_http(
&backend,
&DispatchContext::new(&tokens),
Some("DynamoDB_20120810.ListTables"),
"{}",
unsigned,
));
assert_eq!(out.status, 400);
assert!(
out.body.contains("MissingAuthenticationTokenException"),
"{}",
out.body
);
}
#[test]
fn http_checks_the_target_before_auth() {
let backend = Storage::memory().unwrap();
let tokens = crate::TokenCaches::new();
let out = pollster::block_on(dispatch_http(
&backend,
&DispatchContext::new(&tokens),
Some("DynamoDB_20120810.NoSuchOp"),
"{}",
crate::auth_material::AuthMaterial::default(),
));
assert_eq!(out.body, UNKNOWN_OPERATION_BARE);
}
#[test]
fn http_surfaces_an_api_error_with_its_own_status_and_envelope() {
let backend = Storage::memory().unwrap();
let out = http(
&backend,
Some("DynamoDB_20120810.DescribeTable"),
r#"{"TableName":"Absent"}"#,
);
assert_eq!(out.status, 400);
let v: serde_json::Value = serde_json::from_str(&out.body).unwrap();
assert!(
v["__type"]
.as_str()
.unwrap()
.contains("ResourceNotFoundException"),
"unexpected envelope: {}",
out.body
);
assert!(!is_unsupported_fault(out.status, &out.body));
}
fn partiql(backend: &Storage, statement: &str) -> std::result::Result<String, String> {
let body = serde_json::json!({ "Statement": statement }).to_string();
run(backend, "ExecuteStatement", &body)
}
fn items(response: &str) -> Vec<serde_json::Value> {
serde_json::from_str::<serde_json::Value>(response).unwrap()["Items"]
.as_array()
.cloned()
.unwrap_or_default()
}
fn error_type(err: &str) -> String {
serde_json::from_str::<serde_json::Value>(err).unwrap()["__type"]
.as_str()
.unwrap_or_default()
.to_string()
}
fn message(err: &str) -> String {
serde_json::from_str::<serde_json::Value>(err).unwrap()["message"]
.as_str()
.unwrap_or_default()
.to_string()
}
#[test]
fn execute_statement_inserts_then_selects() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
partiql(
&backend,
"INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's1', 'plays': 3}",
)
.unwrap();
let found = partiql(
&backend,
"SELECT * FROM \"Music\" WHERE artist = 'a' AND song = 's1'",
)
.unwrap();
assert_eq!(items(&found)[0]["plays"]["N"], "3");
}
#[test]
fn execute_statement_insert_is_not_an_upsert() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
let stmt = "INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's1'}";
partiql(&backend, stmt).unwrap();
let err = partiql(&backend, stmt).unwrap_err();
assert!(error_type(&err).contains("DuplicateItemException"), "{err}");
}
#[test]
fn execute_statement_update_on_a_missing_key_fails_the_condition() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
let err = partiql(
&backend,
"UPDATE \"Music\" SET plays = 1 WHERE artist = 'nobody' AND song = 's1'",
)
.unwrap_err();
assert!(
error_type(&err).contains("ConditionalCheckFailedException"),
"{err}"
);
assert_eq!(message(&err), "The conditional request failed");
}
#[test]
fn execute_statement_delete_returning_all_old_hits_and_misses() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
partiql(
&backend,
"INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's1', 'genre': 'rock'}",
)
.unwrap();
let hit = partiql(
&backend,
"DELETE FROM \"Music\" WHERE artist = 'a' AND song = 's1' RETURNING ALL OLD *",
)
.unwrap();
assert_eq!(items(&hit)[0]["genre"]["S"], "rock");
let miss = partiql(
&backend,
"DELETE FROM \"Music\" WHERE artist = 'a' AND song = 's1' RETURNING ALL OLD *",
)
.unwrap();
assert!(items(&miss).is_empty());
assert!(miss.contains("\"Items\""));
}
#[test]
fn execute_statement_rejects_the_returning_variants_delete_does_not_allow() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
for variant in ["MODIFIED OLD *", "ALL NEW *", "MODIFIED NEW *"] {
let err = partiql(
&backend,
&format!(
"DELETE FROM \"Music\" WHERE artist = 'a' AND song = 's1' RETURNING {variant}"
),
)
.unwrap_err();
assert_eq!(
message(&err),
format!(
"Invalid returning clause: RETURNING {variant}. \
Only RETURNING ALL OLD * is allowed in DELETE statements."
)
);
assert!(error_type(&err).ends_with("ValidationException"), "{err}");
}
}
#[test]
fn execute_statement_modified_projections_carry_only_what_changed() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
partiql(
&backend,
"INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's1', \
'profile': {'sub': 'old', 'sib': 'keep'}, 'data': 'gone'}",
)
.unwrap();
let nested = partiql(
&backend,
"UPDATE \"Music\" SET profile.sub = 'new' \
WHERE artist = 'a' AND song = 's1' RETURNING MODIFIED NEW *",
)
.unwrap();
let projected = &items(&nested)[0];
assert_eq!(projected["profile"]["M"]["sub"]["S"], "new");
assert!(projected["profile"]["M"].get("sib").is_none());
assert!(projected.get("artist").is_none());
let removed = partiql(
&backend,
"UPDATE \"Music\" REMOVE data WHERE artist = 'a' AND song = 's1' \
RETURNING MODIFIED NEW *",
)
.unwrap();
assert!(items(&removed).is_empty());
}
#[test]
fn execute_statement_reports_capacity_by_statement_kind() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
let insert = run(
&backend,
"ExecuteStatement",
&serde_json::json!({
"Statement": "INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's1'}",
"ReturnConsumedCapacity": "TOTAL",
})
.to_string(),
)
.unwrap();
let write: serde_json::Value = serde_json::from_str(&insert).unwrap();
assert_eq!(write["ConsumedCapacity"]["TableName"], "Music");
assert!(write["ConsumedCapacity"]["CapacityUnits"].as_f64().unwrap() > 0.0);
let select = run(
&backend,
"ExecuteStatement",
&serde_json::json!({
"Statement": "SELECT * FROM \"Music\" WHERE artist = 'a' AND song = 's1'",
"ReturnConsumedCapacity": "TOTAL",
})
.to_string(),
)
.unwrap();
let read: serde_json::Value = serde_json::from_str(&select).unwrap();
assert_eq!(read["ConsumedCapacity"]["CapacityUnits"], 0.5);
}
#[test]
fn execute_statement_surfaces_parse_and_table_errors() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
let syntax = partiql(&backend, "NOT A STATEMENT").unwrap_err();
assert!(
error_type(&syntax).ends_with("ValidationException"),
"{syntax}"
);
assert!(
message(&syntax).starts_with("Statement wasn't well formed, can't be processed: "),
"{syntax}"
);
let absent = partiql(&backend, "SELECT * FROM \"Absent\" WHERE artist = 'a'").unwrap_err();
assert!(
error_type(&absent).contains("ResourceNotFoundException"),
"{absent}"
);
assert!(!is_unsupported_fault(400, &absent));
}
#[test]
fn batch_execute_statement_reports_failures_per_statement() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
let batch = serde_json::json!({
"Statements": [
{ "Statement": "INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's1'}" },
{ "Statement": "NOT A STATEMENT" },
]
})
.to_string();
let resp = run(&backend, "BatchExecuteStatement", &batch).unwrap();
let v: serde_json::Value = serde_json::from_str(&resp).unwrap();
assert_eq!(v["Responses"][0]["TableName"], "Music");
assert_eq!(v["Responses"][1]["Error"]["Code"], "ValidationError");
assert!(v["Responses"][1].get("TableName").is_none());
}
#[test]
fn batch_execute_statement_rejects_an_empty_statement_list() {
let backend = Storage::memory().unwrap();
let err = run(&backend, "BatchExecuteStatement", r#"{"Statements":[]}"#).unwrap_err();
assert!(error_type(&err).ends_with("ValidationException"), "{err}");
}
#[test]
fn batch_execute_statement_honours_a_member_returning_clause() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
partiql(
&backend,
"INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's1', 'plays': 1}",
)
.unwrap();
let batch = serde_json::json!({
"Statements": [
{ "Statement": "UPDATE \"Music\" SET plays = 2 \
WHERE artist = 'a' AND song = 's1' RETURNING MODIFIED NEW *" },
]
})
.to_string();
let resp = run(&backend, "BatchExecuteStatement", &batch).unwrap();
let v: serde_json::Value = serde_json::from_str(&resp).unwrap();
assert_eq!(v["Responses"][0]["Item"]["plays"]["N"], "2");
assert!(v["Responses"][0]["Item"].get("artist").is_none());
}
fn transaction(statements: &[&str], token: Option<&str>) -> String {
let members: Vec<_> = statements
.iter()
.map(|s| serde_json::json!({ "Statement": s }))
.collect();
let mut body = serde_json::json!({ "TransactStatements": members });
if let Some(token) = token {
body["ClientRequestToken"] = serde_json::json!(token);
}
body.to_string()
}
fn plays(backend: &Storage, song: &str) -> Option<i64> {
let resp = partiql(
backend,
&format!("SELECT * FROM \"Music\" WHERE artist = 'a' AND song = '{song}'"),
)
.unwrap();
items(&resp)
.first()
.and_then(|i| i["plays"]["N"].as_str())
.and_then(|n| n.parse().ok())
}
#[test]
fn execute_transaction_applies_every_statement() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
partiql(
&backend,
"INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's1', 'plays': 1}",
)
.unwrap();
run(
&backend,
"ExecuteTransaction",
&transaction(
&[
"INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's2', 'plays': 9}",
"UPDATE \"Music\" SET plays = 2 WHERE artist = 'a' AND song = 's1'",
],
None,
),
)
.unwrap();
assert_eq!(plays(&backend, "s1"), Some(2));
assert_eq!(plays(&backend, "s2"), Some(9));
}
#[test]
fn execute_transaction_rolls_back_every_statement_on_failure() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
partiql(
&backend,
"INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 'taken'}",
)
.unwrap();
let err = run(
&backend,
"ExecuteTransaction",
&transaction(
&[
"INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 'fresh', 'plays': 1}",
"INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 'taken'}",
],
None,
),
)
.unwrap_err();
let v: serde_json::Value = serde_json::from_str(&err).unwrap();
assert!(
v["__type"]
.as_str()
.unwrap()
.contains("TransactionCanceledException"),
"{err}"
);
let reasons = v["CancellationReasons"].as_array().unwrap();
assert_eq!(reasons[0]["Code"], "None");
assert_eq!(reasons[1]["Code"], "DuplicateItem");
assert_eq!(plays(&backend, "fresh"), None);
}
#[test]
fn execute_transaction_rejects_a_returning_member_without_looking_unsupported() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
let err = run(
&backend,
"ExecuteTransaction",
&transaction(
&["DELETE FROM \"Music\" WHERE artist = 'a' AND song = 's1' RETURNING ALL OLD *"],
None,
),
)
.unwrap_err();
assert_eq!(
message(&err),
"Validation failed in TransactStatements[0]: \
RETURNING clause is not supported in ExecuteTransaction."
);
assert!(error_type(&err).ends_with("ValidationException"), "{err}");
assert!(is_unsupported_fault(400, &err), "{err}");
}
#[test]
fn execute_transaction_rejects_an_empty_statement_list() {
let backend = Storage::memory().unwrap();
let err = run(&backend, "ExecuteTransaction", &transaction(&[], None)).unwrap_err();
assert!(error_type(&err).ends_with("ValidationException"), "{err}");
}
#[test]
fn execute_transaction_replays_a_repeated_token_without_reapplying() {
let backend = Storage::memory().unwrap();
let tokens = crate::TokenCaches::new();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
partiql(
&backend,
"INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's1', 'plays': 0}",
)
.unwrap();
let bump = transaction(
&["UPDATE \"Music\" SET plays = plays + 1 WHERE artist = 'a' AND song = 's1'"],
Some("same-token"),
);
run_with(&backend, &tokens, "ExecuteTransaction", &bump).unwrap();
run_with(&backend, &tokens, "ExecuteTransaction", &bump).unwrap();
assert_eq!(plays(&backend, "s1"), Some(1));
}
#[test]
fn execute_transaction_reuses_a_token_only_for_the_same_statements() {
let backend = Storage::memory().unwrap();
let tokens = crate::TokenCaches::new();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
run_with(
&backend,
&tokens,
"ExecuteTransaction",
&transaction(
&["INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's1'}"],
Some("same-token"),
),
)
.unwrap();
let err = run_with(
&backend,
&tokens,
"ExecuteTransaction",
&transaction(
&["INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's2'}"],
Some("same-token"),
),
)
.unwrap_err();
assert!(
error_type(&err).contains("IdempotentParameterMismatchException"),
"{err}"
);
}
#[test]
fn execute_transaction_replays_when_only_the_capacity_mode_differs() {
let backend = Storage::memory().unwrap();
let tokens = crate::TokenCaches::new();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
partiql(
&backend,
"INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's1', 'plays': 0}",
)
.unwrap();
let statement = "UPDATE \"Music\" SET plays = plays + 1 WHERE artist = 'a' AND song = 's1'";
let with_mode = |mode: Option<&str>| {
let mut body = serde_json::json!({
"TransactStatements": [{ "Statement": statement }],
"ClientRequestToken": "same-token",
});
if let Some(mode) = mode {
body["ReturnConsumedCapacity"] = serde_json::json!(mode);
}
body.to_string()
};
let first = run_with(&backend, &tokens, "ExecuteTransaction", &with_mode(None)).unwrap();
let replay = run_with(
&backend,
&tokens,
"ExecuteTransaction",
&with_mode(Some("TOTAL")),
)
.unwrap();
assert_eq!(plays(&backend, "s1"), Some(1));
let first: serde_json::Value = serde_json::from_str(&first).unwrap();
assert!(first.get("ConsumedCapacity").is_none(), "{first}");
let replay: serde_json::Value = serde_json::from_str(&replay).unwrap();
assert_eq!(replay["ConsumedCapacity"][0]["TableName"], "Music");
}
#[test]
fn execute_transaction_replays_the_first_calls_responses() {
let backend = Storage::memory().unwrap();
let tokens = crate::TokenCaches::new();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
partiql(
&backend,
"INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's1', 'plays': 5}",
)
.unwrap();
let read = transaction(
&["SELECT * FROM \"Music\" WHERE artist = 'a' AND song = 's1'"],
Some("read-token"),
);
let first = run_with(&backend, &tokens, "ExecuteTransaction", &read).unwrap();
let replay = run_with(&backend, &tokens, "ExecuteTransaction", &read).unwrap();
let first: serde_json::Value = serde_json::from_str(&first).unwrap();
let replay: serde_json::Value = serde_json::from_str(&replay).unwrap();
assert_eq!(first["Responses"][0]["Item"]["plays"]["N"], "5");
assert_eq!(first["Responses"], replay["Responses"]);
}
#[test]
fn execute_transaction_frees_a_token_whose_transaction_was_cancelled() {
let backend = Storage::memory().unwrap();
let tokens = crate::TokenCaches::new();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
partiql(
&backend,
"INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 'taken'}",
)
.unwrap();
let cancelled = run_with(
&backend,
&tokens,
"ExecuteTransaction",
&transaction(
&["INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 'taken'}"],
Some("retry-token"),
),
)
.unwrap_err();
assert!(
error_type(&cancelled).contains("TransactionCanceledException"),
"{cancelled}"
);
run_with(
&backend,
&tokens,
"ExecuteTransaction",
&transaction(
&["INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 'fresh', 'plays': 1}"],
Some("retry-token"),
),
)
.unwrap();
assert_eq!(plays(&backend, "fresh"), Some(1));
}
#[test]
fn http_carries_idempotency_across_requests() {
let backend = Storage::memory().unwrap();
let tokens = crate::TokenCaches::new();
http_with(
&backend,
&tokens,
Some("DynamoDB_20120810.CreateTable"),
CREATE_MUSIC,
);
let once = transaction(
&["INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 'once'}"],
Some("http-token"),
);
for _ in 0..2 {
let out = http_with(
&backend,
&tokens,
Some("DynamoDB_20120810.ExecuteTransaction"),
&once,
);
assert_eq!(out.status, 200, "{}", out.body);
assert!(!is_unsupported_fault(out.status, &out.body));
}
}
#[test]
fn execute_transaction_without_a_token_applies_every_time() {
let backend = Storage::memory().unwrap();
let tokens = crate::TokenCaches::new();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
partiql(
&backend,
"INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's1', 'plays': 0}",
)
.unwrap();
let bump = transaction(
&["UPDATE \"Music\" SET plays = plays + 1 WHERE artist = 'a' AND song = 's1'"],
None,
);
run_with(&backend, &tokens, "ExecuteTransaction", &bump).unwrap();
run_with(&backend, &tokens, "ExecuteTransaction", &bump).unwrap();
assert_eq!(plays(&backend, "s1"), Some(2));
}
#[test]
fn execute_transaction_rejects_an_overlong_token() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
let token = "x".repeat(37);
let err = run(
&backend,
"ExecuteTransaction",
&transaction(
&["INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's1'}"],
Some(&token),
),
)
.unwrap_err();
assert!(error_type(&err).ends_with("ValidationException"), "{err}");
assert!(message(&err).contains("clientRequestToken"), "{err}");
}
#[test]
fn http_serves_partiql_and_still_refuses_transact_write_items() {
let backend = Storage::memory().unwrap();
http(
&backend,
Some("DynamoDB_20120810.CreateTable"),
CREATE_MUSIC,
);
for (target, body) in [
(
"DynamoDB_20120810.ExecuteStatement",
serde_json::json!({
"Statement": "INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's1'}"
})
.to_string(),
),
(
"DynamoDB_20120810.BatchExecuteStatement",
serde_json::json!({
"Statements": [{
"Statement": "SELECT * FROM \"Music\" WHERE artist = 'a' AND song = 's1'"
}]
})
.to_string(),
),
] {
let out = http(&backend, Some(target), &body);
assert_eq!(out.status, 200, "{target}: {}", out.body);
assert!(!is_unsupported_fault(out.status, &out.body));
}
let refused = http(&backend, Some("DynamoDB_20120810.TransactWriteItems"), "{}");
assert_eq!(refused.status, 501);
assert!(is_unsupported_fault(refused.status, &refused.body));
}
#[test]
fn execute_statement_pages_a_select_through_next_token() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
for song in ["s1", "s2", "s3"] {
partiql(
&backend,
&format!("INSERT INTO \"Music\" VALUE {{'artist': 'a', 'song': '{song}'}}"),
)
.unwrap();
}
let mut seen = 0;
let mut token: Option<String> = None;
loop {
let mut body = serde_json::json!({
"Statement": "SELECT * FROM \"Music\" WHERE artist = 'a'",
"Limit": 2,
});
if let Some(token) = &token {
body["NextToken"] = serde_json::json!(token);
}
let resp = run(&backend, "ExecuteStatement", &body.to_string()).unwrap();
let v: serde_json::Value = serde_json::from_str(&resp).unwrap();
seen += v["Items"].as_array().unwrap().len();
token = v["NextToken"].as_str().map(str::to_string);
if token.is_none() {
break;
}
assert!(seen <= 3, "pagination did not terminate");
}
assert_eq!(seen, 3, "every row comes back exactly once");
}
#[test]
fn a_partiql_rejection_is_not_mistaken_for_an_unimplemented_operation() {
let backend = Storage::memory().unwrap();
run(&backend, "CreateTable", CREATE_MUSIC).unwrap();
for statement in [
"INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's1', 'tags': [?]}",
"INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's1', 'meta': {'k': ?}}",
"INSERT INTO \"Music\" VALUE {'artist': 'a', 'song': 's1', 'set': << 1, 'a' >>}",
] {
let err = partiql(&backend, statement).unwrap_err();
assert!(
!is_unsupported_fault(400, &err),
"would be scored as unimplemented: {err}"
);
}
}
#[test]
fn supported_ops_matches_the_routing_table() {
let backend = Storage::memory().unwrap();
for op in SUPPORTED_OPS {
if let Err(err) = run(&backend, op, "{}") {
assert!(
!err.contains(UNSUPPORTED_TYPE),
"{op} is in SUPPORTED_OPS but did not route: {err}"
);
}
}
}
}