#![cfg(feature = "tokio")]
mod formats;
use std::borrow::Cow;
use std::sync::Arc;
use crate::{
encryption::{
self, Encrypted, EncryptedEntry, EncryptedSteVecTerm, EncryptedSteVecTermCompat,
EncryptedSteVecTermStandard, IndexTerm, QueryOp, SteQueryVec,
},
zerokms::{self, RecordDecryptError},
};
use stack_auth::AuthStrategyBounds;
use crate::{
encryption::StorageBuilder,
zerokms::{GenerateKeyPayload, IndexKey},
};
use super::zerokms::EncryptedRecord;
use cipherstash_config::{column::IndexType, ColumnConfig};
use serde::{Deserialize, Serialize};
use thiserror::Error;
use uuid::Uuid;
use crate::encryption::{PlaintextTarget, Queryable, ScopedCipher};
use cipherstash_config::column::Index;
use zerokms_protocol::{Context, DecryptionPolicy, UnverifiedContext};
pub const EQL_SCHEMA_VERSION: u16 = 2;
pub const EQL_SCHEMA_VERSION_V3: u16 = 3;
pub async fn encrypt_eql<'a, C>(
cipher: Arc<ScopedCipher<C>>,
plaintexts: Vec<PreparedPlaintext<'a>>,
opts: &EqlEncryptOpts<'a>,
) -> Result<Vec<EqlOutput>, EqlError>
where
C: AuthStrategyBounds,
for<'b> &'b C: stack_auth::AuthStrategy,
{
encrypt_eql_with(
cipher,
plaintexts,
opts,
|encrypted, identifier| Ok(EqlOutput::Store(to_eql_ciphertext(encrypted, identifier)?)),
|index_term, identifier| {
Ok(EqlOutput::Query(to_eql_query_payload(
index_term, identifier,
)?))
},
)
.await
}
pub async fn encrypt_eql_v3<'a, C>(
cipher: Arc<ScopedCipher<C>>,
plaintexts: Vec<PreparedPlaintext<'a>>,
opts: &EqlEncryptOpts<'a>,
) -> Result<Vec<EqlOutputV3>, EqlError>
where
C: AuthStrategyBounds,
for<'b> &'b C: stack_auth::AuthStrategy,
{
encrypt_eql_with(
cipher,
plaintexts,
opts,
|encrypted, identifier| {
Ok(EqlOutputV3::Store(to_eql_ciphertext_v3(
encrypted, identifier,
)?))
},
|index_term, identifier| {
Ok(EqlOutputV3::Query(to_eql_query_payload_v3(
index_term, identifier,
)?))
},
)
.await
}
async fn encrypt_eql_with<'a, C, O>(
cipher: Arc<ScopedCipher<C>>,
plaintexts: Vec<PreparedPlaintext<'a>>,
opts: &EqlEncryptOpts<'a>,
store_fn: impl Fn(Encrypted, &Identifier) -> Result<O, EqlError>,
query_fn: impl Fn(IndexTerm, Identifier) -> Result<O, EqlError>,
) -> Result<Vec<O>, EqlError>
where
C: AuthStrategyBounds,
for<'b> &'b C: stack_auth::AuthStrategy,
{
use std::collections::VecDeque;
let effective_keyset_id = opts.keyset_id.unwrap_or(cipher.keyset_id());
let targets: Vec<EncryptionTarget> =
to_encryption_targets(cipher.index_key(), plaintexts, effective_keyset_id)?;
let mut data_keys = VecDeque::from(
cipher
.generate_data_keys(
generate_data_key_payloads(opts, &targets),
opts.unverified_context.clone(),
)
.await?,
);
targets
.into_iter()
.map(|target| -> Result<O, EqlError> {
match target {
EncryptionTarget::ForStorage(identifier, builder) => {
let encrypted = builder.build_for_encryption().encrypt(
data_keys
.remove(0)
.expect("insufficient data keys to encrypt all plaintexts"),
)?;
store_fn(encrypted, &identifier)
}
EncryptionTarget::ForQuery(identifier, plaintext, index_type, query_op) => {
let index = Index::new(index_type.clone());
let index_term =
(index, plaintext).build_queryable(cipher.clone(), query_op)?;
query_fn(index_term, identifier)
}
}
})
.collect()
}
pub async fn decrypt_eql<'a, C>(
cipher: Arc<ScopedCipher<C>>,
ciphertexts: impl IntoIterator<Item = EqlCiphertext>,
opts: &EqlDecryptOpts<'a>,
) -> Result<Vec<encryption::Plaintext>, EqlError>
where
C: AuthStrategyBounds,
for<'b> &'b C: stack_auth::AuthStrategy,
{
use crate::{encryption::DecryptOptions, zerokms::WithContext};
let decrypt_opts = DecryptOptions {
keyset_id: opts.keyset_id,
unverified_context: opts.unverified_context.clone(),
};
let ciphertexts = ciphertexts
.into_iter()
.map(|eql| {
let (_, ciphertext) = extract_root_ciphertext(eql)?;
Ok(WithContext {
record: ciphertext,
context: opts.lock_context.clone(),
})
})
.collect::<Result<Vec<_>, EqlError>>()?;
Ok(cipher
.decrypt(ciphertexts, &decrypt_opts)
.await
.map_err(|err| convert_zerokms_error(err, cipher.keyset_id(), opts.keyset_id))?
.into_iter()
.map(|decrypted| encryption::Plaintext::from_slice(&decrypted))
.collect::<Result<Vec<_>, _>>()?)
}
pub async fn decrypt_eql_fallible<'a, C>(
cipher: Arc<ScopedCipher<C>>,
ciphertexts: impl IntoIterator<Item = EqlCiphertext>,
opts: &EqlDecryptOpts<'a>,
) -> Result<Vec<Result<encryption::Plaintext, EqlError>>, EqlError>
where
C: AuthStrategyBounds,
for<'b> &'b C: stack_auth::AuthStrategy,
{
use crate::{encryption::DecryptOptions, zerokms::WithContext};
let decrypt_opts = DecryptOptions {
keyset_id: opts.keyset_id,
unverified_context: opts.unverified_context.clone(),
};
let inputs: Vec<_> = ciphertexts.into_iter().collect();
let input_count = inputs.len();
let mut results: Vec<Option<Result<encryption::Plaintext, EqlError>>> =
(0..input_count).map(|_| None).collect();
let mut valid_payloads: Vec<(usize, WithContext<EncryptedRecord>)> =
Vec::with_capacity(input_count);
for (index, eql) in inputs.into_iter().enumerate() {
match extract_root_ciphertext(eql) {
Ok((_, ciphertext)) => valid_payloads.push((
index,
WithContext {
record: ciphertext,
context: opts.lock_context.clone(),
},
)),
Err(err) => {
results[index] = Some(Err(err));
}
}
}
let (indices, payloads): (Vec<usize>, Vec<_>) = valid_payloads.into_iter().unzip();
let decrypt_results = cipher
.decrypt_fallible(payloads, &decrypt_opts)
.await
.map_err(|err| convert_zerokms_error(err, cipher.keyset_id(), opts.keyset_id))?;
for (index, decrypt_result) in indices.into_iter().zip(decrypt_results) {
results[index] = Some(match decrypt_result {
Ok(bytes) => encryption::Plaintext::from_slice(&bytes).map_err(Into::into),
Err(err) => Err(EqlError::from(err)),
});
}
Ok(results
.into_iter()
.map(|r| r.expect("all result slots filled"))
.collect())
}
#[derive(Clone, Debug, Deserialize, Eq, Hash, PartialEq, Serialize)]
pub struct Identifier {
#[serde(rename = "t")]
pub table: String,
#[serde(rename = "c")]
pub column: String,
}
impl Identifier {
pub fn new(table: impl Into<String>, column: impl Into<String>) -> Self {
Self {
table: table.into(),
column: column.into(),
}
}
pub fn table(&self) -> &str {
&self.table
}
pub fn column(&self) -> &str {
&self.column
}
}
#[allow(clippy::large_enum_variant)]
#[derive(Clone, Debug, Deserialize, Serialize)]
#[serde(tag = "k")]
pub enum EqlCiphertext {
#[serde(rename = "ct")]
Encrypted(EncryptedPayload),
#[serde(rename = "sv")]
SteVec(SteVecPayload),
}
impl EqlCiphertext {
pub fn identifier(&self) -> &Identifier {
match self {
EqlCiphertext::Encrypted(p) => &p.identifier,
EqlCiphertext::SteVec(p) => &p.identifier,
}
}
pub fn version(&self) -> u16 {
match self {
EqlCiphertext::Encrypted(p) => p.version,
EqlCiphertext::SteVec(p) => p.version,
}
}
}
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct EncryptedPayload {
#[serde(rename = "v")]
pub version: u16,
#[serde(rename = "i")]
pub identifier: Identifier,
#[serde(rename = "c", with = "formats::mp_base85")]
pub ciphertext: EncryptedRecord,
#[serde(rename = "hm", default, skip_serializing_if = "Option::is_none")]
pub hmac_256: Option<String>,
#[serde(rename = "bf", default, skip_serializing_if = "Option::is_none")]
pub bloom_filter: Option<Vec<u16>>,
#[serde(rename = "ob", default, skip_serializing_if = "Option::is_none")]
pub ore_block_u64_8_256: Option<Vec<String>>,
#[serde(rename = "op", default, skip_serializing_if = "Option::is_none")]
pub ope_cllw: Option<String>,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct SteVecPayload {
#[serde(rename = "v")]
pub version: u16,
#[serde(rename = "i")]
pub identifier: Identifier,
#[serde(rename = "sv")]
pub ste_vec: Vec<SteVecEntry>,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct SteVecEntry {
#[serde(rename = "s")]
pub selector: String,
#[serde(rename = "c", with = "formats::mp_base85")]
pub ciphertext: EncryptedRecord,
#[serde(rename = "a", default, skip_serializing_if = "Option::is_none")]
pub is_array: Option<bool>,
#[serde(flatten)]
pub term: SteVecEntryTerm,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
#[serde(untagged)]
#[non_exhaustive]
pub enum SteVecEntryTerm {
Hmac {
#[serde(rename = "hm")]
hmac_256: String,
},
OreCllw {
#[serde(rename = "oc")]
ore_cllw_8: String,
},
Ope {
#[serde(rename = "op")]
ope_cllw: String,
},
}
#[derive(Debug, Deserialize, Serialize)]
#[serde(untagged)]
pub enum EqlOutput {
Store(EqlCiphertext),
Query(EqlQueryPayload),
}
#[derive(Debug, Deserialize, Serialize)]
#[serde(tag = "k")]
pub enum EqlQueryPayload {
#[serde(rename = "ct")]
Encrypted(EncryptedQueryPayload),
#[serde(rename = "sv")]
SteVec(SteVecQueryPayload),
}
impl EqlQueryPayload {
pub fn identifier(&self) -> &Identifier {
match self {
EqlQueryPayload::Encrypted(p) => &p.identifier,
EqlQueryPayload::SteVec(p) => &p.identifier,
}
}
}
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct EncryptedQueryPayload {
#[serde(rename = "v")]
pub version: u16,
#[serde(rename = "i")]
pub identifier: Identifier,
#[serde(flatten)]
pub term: RootQueryTerm,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
#[serde(untagged)]
pub enum RootQueryTerm {
Hmac {
#[serde(rename = "hm")]
hmac_256: String,
},
BloomFilter {
#[serde(rename = "bf")]
bloom_filter: Vec<u16>,
},
OreBlock {
#[serde(rename = "ob")]
ore_block_u64_8_256: Vec<String>,
},
Ope {
#[serde(rename = "op")]
ope_cllw: String,
},
}
#[derive(Debug, Deserialize, Serialize)]
pub struct SteVecQueryPayload {
#[serde(rename = "v")]
pub version: u16,
#[serde(rename = "i")]
pub identifier: Identifier,
#[serde(flatten)]
pub term: SteVecQueryTerm,
}
#[derive(Debug, Deserialize, Serialize)]
#[serde(untagged)]
#[non_exhaustive]
pub enum SteVecQueryTerm {
Selector {
#[serde(rename = "s")]
selector: String,
},
Hmac {
#[serde(rename = "hm")]
hmac_256: String,
},
OreCllw {
#[serde(rename = "oc")]
ore_cllw_8: String,
},
Ope {
#[serde(rename = "op")]
ope_cllw: String,
},
Containment {
#[serde(rename = "q")]
query_vec: SteQueryVec<16>,
},
}
#[allow(clippy::large_enum_variant)]
#[derive(Clone, Debug, Deserialize, Serialize)]
#[serde(untagged)]
pub enum EqlCiphertextV3 {
SteVec(SteVecPayloadV3),
Encrypted(EncryptedPayloadV3),
}
impl EqlCiphertextV3 {
pub fn identifier(&self) -> &Identifier {
match self {
EqlCiphertextV3::Encrypted(p) => &p.identifier,
EqlCiphertextV3::SteVec(p) => &p.identifier,
}
}
pub fn version(&self) -> u16 {
match self {
EqlCiphertextV3::Encrypted(p) => p.version,
EqlCiphertextV3::SteVec(p) => p.version,
}
}
pub fn into_query_operand(self) -> EqlQueryPayloadV3 {
match self {
EqlCiphertextV3::Encrypted(p) => {
EqlQueryPayloadV3::Encrypted(EncryptedQueryPayloadV3 {
version: p.version,
identifier: p.identifier,
hmac_256: p.hmac_256,
bloom_filter: p.bloom_filter,
ore_block_u64_8_256: p.ore_block_u64_8_256,
ope_cllw: p.ope_cllw,
})
}
EqlCiphertextV3::SteVec(doc) => EqlQueryPayloadV3::SteVec(SteVecQueryPayloadV3 {
ste_vec: doc
.ste_vec
.into_iter()
.map(|e| SteVecQueryEntryV3 {
selector: e.selector,
term: e.term,
})
.collect(),
}),
}
}
}
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct EncryptedPayloadV3 {
#[serde(rename = "v")]
pub version: u16,
#[serde(rename = "i")]
pub identifier: Identifier,
#[serde(rename = "c", with = "formats::mp_base85")]
pub ciphertext: EncryptedRecord,
#[serde(rename = "hm", default, skip_serializing_if = "Option::is_none")]
pub hmac_256: Option<String>,
#[serde(rename = "bf", default, skip_serializing_if = "Option::is_none")]
pub bloom_filter: Option<Vec<i16>>,
#[serde(rename = "ob", default, skip_serializing_if = "Option::is_none")]
pub ore_block_u64_8_256: Option<Vec<String>>,
#[serde(rename = "op", default, skip_serializing_if = "Option::is_none")]
pub ope_cllw: Option<String>,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct SteVecPayloadV3 {
#[serde(rename = "v")]
pub version: u16,
#[serde(rename = "k")]
pub kind: SteVecKind,
#[serde(rename = "i")]
pub identifier: Identifier,
#[serde(rename = "sv")]
pub ste_vec: Vec<SteVecEntryV3>,
}
#[derive(Clone, Copy, Debug, Default, Deserialize, Serialize)]
pub enum SteVecKind {
#[default]
#[serde(rename = "sv")]
SteVec,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct SteVecEntryV3 {
#[serde(rename = "s")]
pub selector: String,
#[serde(rename = "c", with = "formats::mp_base85")]
pub ciphertext: EncryptedRecord,
#[serde(rename = "a", default, skip_serializing_if = "Option::is_none")]
pub is_array: Option<bool>,
#[serde(flatten)]
pub term: SteVecEntryTermV3,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
#[serde(untagged)]
#[non_exhaustive]
pub enum SteVecEntryTermV3 {
Hmac {
#[serde(rename = "hm")]
hmac_256: String,
},
Ope {
#[serde(rename = "op")]
ope_cllw: String,
},
}
#[derive(Debug, Deserialize, Serialize)]
#[serde(untagged)]
pub enum EqlOutputV3 {
Store(EqlCiphertextV3),
Query(EqlQueryPayloadV3),
}
#[derive(Debug, Deserialize, Serialize)]
#[serde(untagged)]
pub enum EqlQueryPayloadV3 {
SteVec(SteVecQueryPayloadV3),
Encrypted(EncryptedQueryPayloadV3),
Selector(String),
}
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct EncryptedQueryPayloadV3 {
#[serde(rename = "v")]
pub version: u16,
#[serde(rename = "i")]
pub identifier: Identifier,
#[serde(rename = "hm", default, skip_serializing_if = "Option::is_none")]
pub hmac_256: Option<String>,
#[serde(rename = "bf", default, skip_serializing_if = "Option::is_none")]
pub bloom_filter: Option<Vec<i16>>,
#[serde(rename = "ob", default, skip_serializing_if = "Option::is_none")]
pub ore_block_u64_8_256: Option<Vec<String>>,
#[serde(rename = "op", default, skip_serializing_if = "Option::is_none")]
pub ope_cllw: Option<String>,
}
#[derive(Debug, Deserialize, Serialize)]
pub struct SteVecQueryPayloadV3 {
#[serde(rename = "sv")]
pub ste_vec: Vec<SteVecQueryEntryV3>,
}
#[derive(Debug, Deserialize, Serialize)]
pub struct SteVecQueryEntryV3 {
#[serde(rename = "s")]
pub selector: String,
#[serde(flatten)]
pub term: SteVecEntryTermV3,
}
impl From<SteVecEntryTerm> for SteVecQueryTerm {
fn from(term: SteVecEntryTerm) -> Self {
match term {
SteVecEntryTerm::Hmac { hmac_256 } => Self::Hmac { hmac_256 },
SteVecEntryTerm::OreCllw { ore_cllw_8 } => Self::OreCllw { ore_cllw_8 },
SteVecEntryTerm::Ope { ope_cllw } => Self::Ope { ope_cllw },
}
}
}
fn extract_root_ciphertext(eql: EqlCiphertext) -> Result<(Identifier, EncryptedRecord), EqlError> {
match eql {
EqlCiphertext::Encrypted(p) => Ok((p.identifier, p.ciphertext)),
EqlCiphertext::SteVec(p) => {
let SteVecPayload {
identifier,
ste_vec,
..
} = p;
let root = ste_vec
.into_iter()
.next()
.ok_or_else(|| EqlError::MissingCiphertext(identifier.clone()))?;
Ok((identifier, root.ciphertext))
}
}
}
fn to_eql_ciphertext(
encrypted: Encrypted,
identifier: &Identifier,
) -> Result<EqlCiphertext, EqlError> {
match encrypted {
Encrypted::Record(ciphertext, terms) => {
let mut payload = EncryptedPayload {
version: EQL_SCHEMA_VERSION,
identifier: identifier.clone(),
ciphertext,
hmac_256: None,
bloom_filter: None,
ore_block_u64_8_256: None,
ope_cllw: None,
};
for term in terms {
apply_root_term(&mut payload, term);
}
Ok(EqlCiphertext::Encrypted(payload))
}
Encrypted::SteVec(ste_vec) => {
let elements: Vec<SteVecEntry> = ste_vec
.into_iter()
.map(
|EncryptedEntry {
tokenized_selector,
term,
record,
parent_is_array,
}| {
SteVecEntry {
selector: hex::encode(tokenized_selector.as_bytes()),
ciphertext: record,
is_array: Some(parent_is_array),
term: ste_vec_term_from_encrypted(term),
}
},
)
.collect();
Ok(EqlCiphertext::SteVec(SteVecPayload {
version: EQL_SCHEMA_VERSION,
identifier: identifier.clone(),
ste_vec: elements,
}))
}
}
}
fn apply_root_term(payload: &mut EncryptedPayload, term: IndexTerm) {
match term {
IndexTerm::Binary(bytes) => {
payload.hmac_256 = Some(hex::encode(bytes));
}
IndexTerm::BitMap(bf) => {
payload.bloom_filter = Some(bf);
}
IndexTerm::OreFull(bytes) | IndexTerm::OreLeft(bytes) => {
payload.ore_block_u64_8_256 = Some(vec![hex::encode(bytes)]);
}
IndexTerm::OreArray(arr) => {
payload.ore_block_u64_8_256 = Some(arr.iter().map(hex::encode).collect());
}
IndexTerm::OpeFixed(bytes) | IndexTerm::OpeVariable(bytes) => {
payload.ope_cllw = Some(hex::encode(bytes));
}
IndexTerm::BinaryVec(_)
| IndexTerm::SteVecSelector(_)
| IndexTerm::SteVecTerm(_)
| IndexTerm::SteQueryVec(_)
| IndexTerm::Null => {}
}
}
fn ste_vec_term_from_encrypted(term: EncryptedSteVecTerm) -> SteVecEntryTerm {
match term {
EncryptedSteVecTerm::Compat(EncryptedSteVecTermCompat::Mac(bytes))
| EncryptedSteVecTerm::Standard(EncryptedSteVecTermStandard::Mac(bytes)) => {
SteVecEntryTerm::Hmac {
hmac_256: hex::encode(bytes),
}
}
EncryptedSteVecTerm::Standard(EncryptedSteVecTermStandard::Ore(ore)) => {
SteVecEntryTerm::OreCllw {
ore_cllw_8: hex::encode(ore.as_ref()),
}
}
EncryptedSteVecTerm::Compat(EncryptedSteVecTermCompat::Ope(ope)) => SteVecEntryTerm::Ope {
ope_cllw: hex::encode(ope.as_ref()),
},
}
}
fn to_eql_query_payload(
index_term: IndexTerm,
identifier: Identifier,
) -> Result<EqlQueryPayload, EqlError> {
match index_term {
IndexTerm::Binary(bytes) => Ok(EqlQueryPayload::Encrypted(EncryptedQueryPayload {
version: EQL_SCHEMA_VERSION,
identifier,
term: RootQueryTerm::Hmac {
hmac_256: hex::encode(bytes),
},
})),
IndexTerm::BitMap(bf) => Ok(EqlQueryPayload::Encrypted(EncryptedQueryPayload {
version: EQL_SCHEMA_VERSION,
identifier,
term: RootQueryTerm::BloomFilter { bloom_filter: bf },
})),
IndexTerm::OreFull(bytes) | IndexTerm::OreLeft(bytes) => {
Ok(EqlQueryPayload::Encrypted(EncryptedQueryPayload {
version: EQL_SCHEMA_VERSION,
identifier,
term: RootQueryTerm::OreBlock {
ore_block_u64_8_256: vec![hex::encode(bytes)],
},
}))
}
IndexTerm::OreArray(arr) => Ok(EqlQueryPayload::Encrypted(EncryptedQueryPayload {
version: EQL_SCHEMA_VERSION,
identifier,
term: RootQueryTerm::OreBlock {
ore_block_u64_8_256: arr.iter().map(hex::encode).collect(),
},
})),
IndexTerm::SteVecSelector(selector) => Ok(EqlQueryPayload::SteVec(SteVecQueryPayload {
version: EQL_SCHEMA_VERSION,
identifier,
term: SteVecQueryTerm::Selector {
selector: hex::encode(selector.as_bytes()),
},
})),
IndexTerm::SteVecTerm(term) => Ok(EqlQueryPayload::SteVec(SteVecQueryPayload {
version: EQL_SCHEMA_VERSION,
identifier,
term: ste_vec_term_from_encrypted(term).into(),
})),
IndexTerm::SteQueryVec(query_vec) => Ok(EqlQueryPayload::SteVec(SteVecQueryPayload {
version: EQL_SCHEMA_VERSION,
identifier,
term: SteVecQueryTerm::Containment { query_vec },
})),
IndexTerm::OpeFixed(bytes) | IndexTerm::OpeVariable(bytes) => {
Ok(EqlQueryPayload::Encrypted(EncryptedQueryPayload {
version: EQL_SCHEMA_VERSION,
identifier,
term: RootQueryTerm::Ope {
ope_cllw: hex::encode(bytes),
},
}))
}
IndexTerm::BinaryVec(_) | IndexTerm::Null => Err(EqlError::InvalidIndexTerm),
}
}
fn to_eql_ciphertext_v3(
encrypted: Encrypted,
identifier: &Identifier,
) -> Result<EqlCiphertextV3, EqlError> {
match encrypted {
Encrypted::Record(ciphertext, terms) => {
let mut payload = EncryptedPayloadV3 {
version: EQL_SCHEMA_VERSION_V3,
identifier: identifier.clone(),
ciphertext,
hmac_256: None,
bloom_filter: None,
ore_block_u64_8_256: None,
ope_cllw: None,
};
for term in terms {
apply_root_term_v3(&mut payload, term);
}
Ok(EqlCiphertextV3::Encrypted(payload))
}
Encrypted::SteVec(ste_vec) => {
let elements = ste_vec
.into_iter()
.map(
|EncryptedEntry {
tokenized_selector,
term,
record,
parent_is_array,
}| {
Ok(SteVecEntryV3 {
selector: hex::encode(tokenized_selector.as_bytes()),
ciphertext: record,
is_array: Some(parent_is_array),
term: ste_vec_entry_term_v3(term)?,
})
},
)
.collect::<Result<Vec<_>, EqlError>>()?;
Ok(EqlCiphertextV3::SteVec(SteVecPayloadV3 {
version: EQL_SCHEMA_VERSION_V3,
kind: SteVecKind::SteVec,
identifier: identifier.clone(),
ste_vec: elements,
}))
}
}
}
fn apply_root_term_v3(payload: &mut EncryptedPayloadV3, term: IndexTerm) {
match term {
IndexTerm::Binary(bytes) => {
payload.hmac_256 = Some(hex::encode(bytes));
}
IndexTerm::BitMap(bf) => {
payload.bloom_filter = Some(bloom_filter_to_signed(bf));
}
IndexTerm::OreFull(bytes) | IndexTerm::OreLeft(bytes) => {
payload.ore_block_u64_8_256 = Some(vec![hex::encode(bytes)]);
}
IndexTerm::OreArray(arr) => {
payload.ore_block_u64_8_256 = Some(arr.iter().map(hex::encode).collect());
}
IndexTerm::OpeFixed(bytes) | IndexTerm::OpeVariable(bytes) => {
payload.ope_cllw = Some(hex::encode(bytes));
}
IndexTerm::BinaryVec(_)
| IndexTerm::SteVecSelector(_)
| IndexTerm::SteVecTerm(_)
| IndexTerm::SteQueryVec(_)
| IndexTerm::Null => {}
}
}
fn bloom_filter_to_signed(bf: Vec<u16>) -> Vec<i16> {
bf.into_iter().map(|bit| bit as i16).collect()
}
fn ste_vec_entry_term_v3(term: EncryptedSteVecTerm) -> Result<SteVecEntryTermV3, EqlError> {
match term {
EncryptedSteVecTerm::Compat(EncryptedSteVecTermCompat::Mac(bytes))
| EncryptedSteVecTerm::Standard(EncryptedSteVecTermStandard::Mac(bytes)) => {
Ok(SteVecEntryTermV3::Hmac {
hmac_256: hex::encode(bytes),
})
}
EncryptedSteVecTerm::Compat(EncryptedSteVecTermCompat::Ope(ope)) => {
Ok(SteVecEntryTermV3::Ope {
ope_cllw: hex::encode(ope.as_ref()),
})
}
EncryptedSteVecTerm::Standard(EncryptedSteVecTermStandard::Ore(_)) => {
Err(EqlError::UnsupportedSteVecOreInV3)
}
}
}
fn to_eql_query_payload_v3(
index_term: IndexTerm,
identifier: Identifier,
) -> Result<EqlQueryPayloadV3, EqlError> {
let mut operand = EncryptedQueryPayloadV3 {
version: EQL_SCHEMA_VERSION_V3,
identifier,
hmac_256: None,
bloom_filter: None,
ore_block_u64_8_256: None,
ope_cllw: None,
};
match index_term {
IndexTerm::Binary(bytes) => operand.hmac_256 = Some(hex::encode(bytes)),
IndexTerm::BitMap(bf) => operand.bloom_filter = Some(bloom_filter_to_signed(bf)),
IndexTerm::OreFull(bytes) | IndexTerm::OreLeft(bytes) => {
operand.ore_block_u64_8_256 = Some(vec![hex::encode(bytes)]);
}
IndexTerm::OreArray(arr) => {
operand.ore_block_u64_8_256 = Some(arr.iter().map(hex::encode).collect());
}
IndexTerm::OpeFixed(bytes) | IndexTerm::OpeVariable(bytes) => {
operand.ope_cllw = Some(hex::encode(bytes));
}
IndexTerm::SteQueryVec(query_vec) => {
let entries = query_vec
.into_iter()
.map(|entry| {
let (selector, term) = entry.into_parts();
Ok(SteVecQueryEntryV3 {
selector: hex::encode(selector.as_bytes()),
term: ste_vec_entry_term_v3(term)?,
})
})
.collect::<Result<Vec<_>, EqlError>>()?;
return Ok(EqlQueryPayloadV3::SteVec(SteVecQueryPayloadV3 {
ste_vec: entries,
}));
}
IndexTerm::SteVecSelector(selector) => {
return Ok(EqlQueryPayloadV3::Selector(hex::encode(
selector.as_bytes(),
)));
}
IndexTerm::SteVecTerm(_) => return Err(EqlError::UnsupportedV3QueryTerm),
IndexTerm::BinaryVec(_) | IndexTerm::Null => return Err(EqlError::InvalidIndexTerm),
}
Ok(EqlQueryPayloadV3::Encrypted(operand))
}
#[derive(Error, Debug)]
pub enum EqlError {
#[error(transparent)]
CiphertextCouldNotBeSerialised(#[from] serde_json::Error),
#[error("Encrypted column could not be parsed")]
ColumnCouldNotBeParsed,
#[error("Encrypted column is null")]
ColumnIsNull,
#[error("Column '{column}' in table '{table}' could not be deserialised")]
ColumnCouldNotBeDeserialised { table: String, column: String },
#[error("Column '{column}' in table '{table}' could not be encrypted")]
ColumnCouldNotBeEncrypted { table: String, column: String },
#[error("Column configuration for column '{column}' in table '{table}' does not match the encrypted column")]
ColumnConfigurationMismatch { table: String, column: String },
#[error("Could not decrypt data using keyset '{keyset_id}'")]
CouldNotDecryptDataForKeyset {
keyset_id: String,
#[source]
source: zerokms::Error,
},
#[error("InvalidIndexTerm")]
InvalidIndexTerm,
#[error(
"SteVec CLLW-ORE (`oc`) ordering terms have no EQL v3 representation; \
encrypt the column in Compat (OPE) mode for v3"
)]
UnsupportedSteVecOreInV3,
#[error("index term has no EQL v3 query-operand representation")]
UnsupportedV3QueryTerm,
#[error("EQL payload for column '{}' in table '{}' is missing ciphertext", _0.column(), _0.table())]
MissingCiphertext(Identifier),
#[error("KeysetId `{id}` could not be parsed using `SET CIPHERSTASH.KEYSET_ID`. KeysetId should be a valid UUID")]
KeysetIdCouldNotBeParsed { id: String },
#[error("Keyset Id could not be set using `SET CIPHERSTASH.KEYSET_ID`")]
KeysetIdCouldNotBeSet,
#[error("Keyset Name could not be set using `SET CIPHERSTASH.KEYSET_NAME`")]
KeysetNameCouldNotBeSet,
#[error("Missing encrypt configuration for column type `{plaintext_type}`")]
MissingEncryptConfiguration { plaintext_type: &'static str },
#[error("Decrypted column could not be encoded as the expected type")]
PlaintextCouldNotBeEncoded,
#[error(transparent)]
Pipeline(#[from] encryption::EncryptionError),
#[error(transparent)]
PlaintextCouldNotBeDecoded(#[from] encryption::TypeParseError),
#[error("Missing keyset identifer")]
MissingKeysetIdentifier,
#[error("Cannot SET CIPHERSTASH.KEYSET if a default keyset has been configured")]
UnexpectedSetKeyset,
#[error("Column '{column}' in table '{table}' has no Encrypt configuration")]
UnknownColumn { table: String, column: String },
#[error("Unknown keyset name or id '{keyset}'. Check the configured credentials")]
UnknownKeysetIdentifier { keyset: String },
#[error("Table '{table}' has no Encrypt configuration")]
UnknownTable { table: String },
#[error("Unknown Index Term for column '{}' in table '{}'", _0.column(), _0.table())]
UnknownIndexTerm(Identifier),
#[error("ZeroKMS error '{}'", _0)]
ZeroKMS(#[from] zerokms::Error),
#[error("Record decryption error '{}'", _0)]
RecordDecrypt(#[from] RecordDecryptError),
}
#[derive(Debug)]
pub enum EqlOperation<'a> {
Store,
Query(&'a IndexType, QueryOp),
}
pub struct PreparedPlaintext<'a> {
identifier: Identifier,
plaintext: encryption::Plaintext,
eql_op: EqlOperation<'a>,
column_config: Cow<'a, ColumnConfig>,
}
impl<'a> PreparedPlaintext<'a> {
pub fn new(
column_config: Cow<'a, ColumnConfig>,
identifier: Identifier,
plaintext: encryption::Plaintext,
eql_op: EqlOperation<'a>,
) -> Self {
Self {
identifier,
plaintext,
eql_op,
column_config,
}
}
}
enum EncryptionTarget<'a> {
ForStorage(Identifier, StorageBuilder<'a, encryption::Plaintext>),
ForQuery(Identifier, encryption::Plaintext, &'a IndexType, QueryOp),
}
fn generate_data_key_payloads<'a>(
opts: &EqlEncryptOpts<'a>,
targets: &'a Vec<EncryptionTarget<'a>>,
) -> Vec<GenerateKeyPayload<'a>> {
targets
.iter()
.filter_map(|target| match target {
EncryptionTarget::ForStorage(_, builder) => {
let payload =
GenerateKeyPayload::new(builder.descriptor(), opts.lock_context.clone());
Some(match opts.decryption_policy.clone() {
Some(p) => payload.with_decryption_policy(p),
None => payload,
})
}
EncryptionTarget::ForQuery(_, _, _, _) => None,
})
.collect()
}
fn to_encryption_targets<'a>(
index_key: &'a IndexKey,
plaintexts: Vec<PreparedPlaintext<'a>>,
effective_keyset_id: Uuid,
) -> Result<Vec<EncryptionTarget<'a>>, encryption::EncryptionError> {
plaintexts
.into_iter()
.map(
move |prepared_plaintext| -> Result<EncryptionTarget, encryption::EncryptionError> {
use crate::encryption::Encryptable;
let PreparedPlaintext {
identifier,
plaintext,
eql_op,
column_config,
} = prepared_plaintext;
match eql_op {
EqlOperation::Store => Ok(EncryptionTarget::ForStorage(
identifier,
PlaintextTarget::new(plaintext, (*column_config).clone())
.build_encryptable(index_key, effective_keyset_id)?,
)),
EqlOperation::Query(index_type, query_op) => Ok(EncryptionTarget::ForQuery(
identifier, plaintext, index_type, query_op,
)),
}
},
)
.collect::<Result<Vec<_>, _>>()
}
#[derive(Debug, Default)]
pub struct EqlDecryptOpts<'a> {
pub keyset_id: Option<Uuid>,
pub lock_context: Cow<'a, [Context]>,
pub unverified_context: Option<Cow<'a, UnverifiedContext>>,
}
#[derive(Debug, Default)]
pub struct EqlEncryptOpts<'a> {
pub keyset_id: Option<Uuid>,
pub lock_context: Cow<'a, [Context]>,
pub unverified_context: Option<Cow<'a, UnverifiedContext>>,
pub index_types: Option<Cow<'a, [IndexType]>>,
pub decryption_policy: Option<DecryptionPolicy>,
}
fn convert_zerokms_error(
err: zerokms::Error,
cipher_keyset_id: Uuid,
keyset_id_override: Option<Uuid>,
) -> EqlError {
match err {
zerokms::Error::Decrypt(_) => {
let error_msg = err.to_string();
if error_msg.contains("Failed to retrieve key") {
EqlError::CouldNotDecryptDataForKeyset {
keyset_id: keyset_id_override
.map(|id| id.to_string())
.unwrap_or(cipher_keyset_id.to_string()),
source: err,
}
} else {
EqlError::ZeroKMS(err)
}
}
_ => EqlError::ZeroKMS(err),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn convert_zerokms_error_preserves_source_chain() {
use crate::zerokms::{DecryptError, RetrieveKeyError};
use std::error::Error as StdError;
let cipher_keyset_id = Uuid::new_v4();
let override_id = Uuid::new_v4();
let inner = zerokms::Error::from(DecryptError::from(RetrieveKeyError::FailedRetrieval(
"no such key".into(),
)));
let converted = convert_zerokms_error(inner, cipher_keyset_id, Some(override_id));
match converted {
EqlError::CouldNotDecryptDataForKeyset { ref keyset_id, .. } => {
assert_eq!(keyset_id, &override_id.to_string());
}
ref other => panic!("expected CouldNotDecryptDataForKeyset, got {other:?}"),
}
let source = StdError::source(&converted).expect("source() should be Some");
assert!(
source.downcast_ref::<zerokms::Error>().is_some(),
"source should downcast to &zerokms::Error"
);
}
#[test]
fn empty_ste_vec_payload_yields_missing_ciphertext() {
let identifier = Identifier::new("test_table", "test_column");
let eql = EqlCiphertext::SteVec(SteVecPayload {
version: EQL_SCHEMA_VERSION,
identifier: identifier.clone(),
ste_vec: Vec::new(),
});
let result = extract_root_ciphertext(eql);
assert!(matches!(result, Err(EqlError::MissingCiphertext(_))));
}
#[test]
fn mp_base85_deserialize_invalid_input_returns_error() {
use serde::de::value::{Error as ValueError, StrDeserializer};
use serde::de::IntoDeserializer;
let invalid: StrDeserializer<ValueError> = "not-valid-base85!!!".into_deserializer();
let result: Result<EncryptedRecord, ValueError> = formats::mp_base85::deserialize(invalid);
assert!(result.is_err(), "Invalid base85 input should return error");
}
#[test]
fn encrypted_payload_serializes_with_k_ct_discriminator() {
let record = EncryptedRecord {
iv: Default::default(),
ciphertext: vec![1; 16],
tag: vec![2; 16],
descriptor: "users/email".to_string(),
keyset_id: Some(Uuid::nil()),
decryption_policy: None,
};
let payload = EqlCiphertext::Encrypted(EncryptedPayload {
version: EQL_SCHEMA_VERSION,
identifier: Identifier::new("users", "email"),
ciphertext: record,
hmac_256: Some("deadbeef".into()),
bloom_filter: None,
ore_block_u64_8_256: None,
ope_cllw: None,
});
let value = serde_json::to_value(&payload).unwrap();
assert_eq!(value["k"], "ct");
assert_eq!(value["v"], EQL_SCHEMA_VERSION);
assert_eq!(value["i"]["t"], "users");
assert_eq!(value["i"]["c"], "email");
assert_eq!(value["hm"], "deadbeef");
assert!(value.get("c").is_some());
assert!(value.get("sv").is_none());
assert!(value.get("ob").is_none());
assert!(value.get("op").is_none());
}
fn fake_record(descriptor: &str) -> EncryptedRecord {
EncryptedRecord {
iv: Default::default(),
ciphertext: vec![1; 16],
tag: vec![2; 16],
descriptor: descriptor.to_string(),
keyset_id: Some(Uuid::nil()),
decryption_policy: None,
}
}
mod v3 {
use super::*;
fn term(json: serde_json::Value) -> EncryptedSteVecTerm {
serde_json::from_value(json).expect("valid EncryptedSteVecTerm wire shape")
}
#[test]
fn store_scalar_drops_k_and_bumps_version() {
let encrypted = Encrypted::Record(
fake_record("users/email"),
vec![
IndexTerm::Binary(vec![0xAB; 32]),
IndexTerm::OreArray(vec![vec![0x01; 8], vec![0x02; 8]]),
],
);
let eql = to_eql_ciphertext_v3(encrypted, &Identifier::new("users", "email")).unwrap();
let value = serde_json::to_value(&eql).unwrap();
assert_eq!(value["v"], EQL_SCHEMA_VERSION_V3);
assert!(value.get("k").is_none(), "v3 scalar payloads carry no `k`");
assert_eq!(value["i"]["t"], "users");
assert_eq!(value["hm"], hex::encode([0xAB; 32]));
assert!(value.get("c").is_some());
assert!(value["ob"].is_array());
assert!(value.get("op").is_none());
assert!(value.get("bf").is_none());
}
#[test]
fn store_scalar_bloom_positions_are_signed_smallint() {
let encrypted = Encrypted::Record(
fake_record("t/c"),
vec![IndexTerm::BitMap(vec![0, 32767, 40000])],
);
let eql = to_eql_ciphertext_v3(encrypted, &Identifier::new("t", "c")).unwrap();
let value = serde_json::to_value(&eql).unwrap();
assert_eq!(value["bf"], serde_json::json!([0, 32767, -25536]));
}
#[test]
fn store_scalar_ope_term_on_op() {
let encrypted = Encrypted::Record(
fake_record("orders/total"),
vec![IndexTerm::OpeFixed(vec![0xAB; 65])],
);
let eql = to_eql_ciphertext_v3(encrypted, &Identifier::new("orders", "total")).unwrap();
let value = serde_json::to_value(&eql).unwrap();
assert_eq!(value["op"], hex::encode([0xAB; 65]));
assert!(value.get("k").is_none());
}
#[test]
fn ste_vec_entry_term_maps_mac_and_ope_rejects_ore() {
let hm =
ste_vec_entry_term_v3(term(serde_json::json!({ "hm": "11".repeat(16) }))).unwrap();
assert!(matches!(hm, SteVecEntryTermV3::Hmac { .. }));
let op = ste_vec_entry_term_v3(term(serde_json::json!({ "op": "010203" }))).unwrap();
assert!(matches!(op, SteVecEntryTermV3::Ope { .. }));
let ore = ste_vec_entry_term_v3(term(serde_json::json!({ "oc": "4243" })));
assert!(matches!(ore, Err(EqlError::UnsupportedSteVecOreInV3)));
}
#[test]
fn ste_vec_document_keeps_k_and_omits_per_entry_envelope() {
let doc = SteVecPayloadV3 {
version: EQL_SCHEMA_VERSION_V3,
kind: SteVecKind::SteVec,
identifier: Identifier::new("docs", "attrs"),
ste_vec: vec![
SteVecEntryV3 {
selector: "aa".into(),
ciphertext: fake_record("docs/attrs"),
is_array: None,
term: SteVecEntryTermV3::Hmac {
hmac_256: "beef".into(),
},
},
SteVecEntryV3 {
selector: "bb".into(),
ciphertext: fake_record("docs/attrs"),
is_array: Some(false),
term: SteVecEntryTermV3::Ope {
ope_cllw: "00ff".into(),
},
},
],
};
let v = serde_json::to_value(EqlCiphertextV3::SteVec(doc)).unwrap();
assert_eq!(v["v"], EQL_SCHEMA_VERSION_V3);
assert_eq!(v["k"], "sv");
let sv = v["sv"].as_array().unwrap();
assert_eq!(sv.len(), 2);
assert_eq!(sv[0]["s"], "aa");
assert_eq!(sv[0]["hm"], "beef");
assert!(sv[0].get("a").is_none(), "None array marker is omitted");
assert!(sv[0].get("v").is_none(), "entries carry no per-entry `v`");
assert!(sv[0].get("i").is_none());
assert!(sv[0].get("k").is_none());
assert_eq!(sv[1]["op"], "00ff");
assert_eq!(sv[1]["a"], false);
assert!(sv[1].get("hm").is_none());
}
#[test]
fn query_scalar_operand_is_enveloped_without_k_or_c() {
let out = to_eql_query_payload_v3(
IndexTerm::Binary(vec![0xCD; 32]),
Identifier::new("users", "email"),
)
.unwrap();
let v = serde_json::to_value(&out).unwrap();
assert_eq!(v["v"], EQL_SCHEMA_VERSION_V3);
assert_eq!(v["i"]["t"], "users");
assert_eq!(v["hm"], hex::encode([0xCD; 32]));
assert!(v.get("k").is_none());
assert!(v.get("c").is_none());
}
#[test]
fn query_scalar_bloom_operand_is_signed() {
let out = to_eql_query_payload_v3(
IndexTerm::BitMap(vec![1, 40000]),
Identifier::new("t", "c"),
)
.unwrap();
let v = serde_json::to_value(&out).unwrap();
assert_eq!(v["bf"], serde_json::json!([1, -25536]));
}
#[test]
fn into_query_operand_scalar_keeps_all_terms_and_drops_c() {
let stored = to_eql_ciphertext_v3(
Encrypted::Record(
fake_record("users/email"),
vec![
IndexTerm::Binary(vec![0xAB; 32]),
IndexTerm::OreArray(vec![vec![1; 8], vec![2; 8]]),
],
),
&Identifier::new("users", "email"),
)
.unwrap();
let v = serde_json::to_value(stored.into_query_operand()).unwrap();
assert_eq!(v["v"], EQL_SCHEMA_VERSION_V3);
assert_eq!(v["i"]["t"], "users");
assert_eq!(v["hm"], hex::encode([0xAB; 32]));
assert!(v["ob"].is_array());
assert!(v.get("c").is_none(), "a query operand must not carry `c`");
assert!(v.get("k").is_none());
}
#[test]
fn into_query_operand_ste_vec_doc_becomes_the_bare_needle() {
let doc = EqlCiphertextV3::SteVec(SteVecPayloadV3 {
version: EQL_SCHEMA_VERSION_V3,
kind: SteVecKind::SteVec,
identifier: Identifier::new("docs", "attrs"),
ste_vec: vec![SteVecEntryV3 {
selector: "aa".into(),
ciphertext: fake_record("docs/attrs"),
is_array: Some(false),
term: SteVecEntryTermV3::Ope {
ope_cllw: "00ff".into(),
},
}],
});
let v = serde_json::to_value(doc.into_query_operand()).unwrap();
assert!(v.get("v").is_none(), "the needle is not enveloped");
assert!(v.get("i").is_none());
let sv = v["sv"].as_array().unwrap();
assert_eq!(sv[0]["s"], "aa");
assert_eq!(sv[0]["op"], "00ff");
assert!(sv[0].get("c").is_none(), "per-entry `c` is dropped");
assert!(sv[0].get("a").is_none(), "per-entry `a` is dropped");
}
#[test]
fn query_ste_vec_selector_is_the_bare_hash() {
use crate::encryption::TokenizedSelector;
let out = to_eql_query_payload_v3(
IndexTerm::SteVecSelector(TokenizedSelector([7; 16])),
Identifier::new("docs", "attrs"),
)
.unwrap();
assert!(matches!(out, EqlQueryPayloadV3::Selector(_)));
assert_eq!(
serde_json::to_value(&out).unwrap(),
serde_json::json!(hex::encode([7u8; 16]))
);
}
#[test]
fn jsonb_containment_needle_is_bare_sv_array() {
let needle = EqlQueryPayloadV3::SteVec(SteVecQueryPayloadV3 {
ste_vec: vec![SteVecQueryEntryV3 {
selector: "abcd".into(),
term: SteVecEntryTermV3::Hmac {
hmac_256: "beef".into(),
},
}],
});
let v = serde_json::to_value(&needle).unwrap();
assert!(v.get("v").is_none(), "the needle is not enveloped");
assert!(v.get("i").is_none());
let sv = v["sv"].as_array().unwrap();
assert_eq!(sv[0]["s"], "abcd");
assert_eq!(sv[0]["hm"], "beef");
}
#[test]
fn output_store_serializes_transparently() {
let encrypted =
Encrypted::Record(fake_record("t/c"), vec![IndexTerm::Binary(vec![9; 32])]);
let store = EqlOutputV3::Store(
to_eql_ciphertext_v3(encrypted, &Identifier::new("t", "c")).unwrap(),
);
let v = serde_json::to_value(&store).unwrap();
assert_eq!(v["v"], EQL_SCHEMA_VERSION_V3);
assert_eq!(v["hm"], hex::encode([9u8; 32]));
}
#[test]
fn ste_vec_document_round_trips_through_serde() {
let doc = SteVecPayloadV3 {
version: EQL_SCHEMA_VERSION_V3,
kind: SteVecKind::SteVec,
identifier: Identifier::new("docs", "attrs"),
ste_vec: vec![
SteVecEntryV3 {
selector: "aa".into(),
ciphertext: fake_record("docs/attrs"),
is_array: None,
term: SteVecEntryTermV3::Hmac {
hmac_256: "beef".into(),
},
},
SteVecEntryV3 {
selector: "bb".into(),
ciphertext: fake_record("docs/attrs"),
is_array: Some(true),
term: SteVecEntryTermV3::Ope {
ope_cllw: "00ff".into(),
},
},
],
};
let wire = serde_json::to_value(&doc).unwrap();
let parsed: SteVecPayloadV3 = serde_json::from_value(wire.clone()).unwrap();
assert_eq!(
serde_json::to_value(&parsed).unwrap(),
wire,
"document must round-trip byte-for-byte"
);
assert!(matches!(
parsed.ste_vec[0].term,
SteVecEntryTermV3::Hmac { ref hmac_256 } if hmac_256 == "beef"
));
assert_eq!(parsed.ste_vec[0].is_array, None);
assert!(matches!(
parsed.ste_vec[1].term,
SteVecEntryTermV3::Ope { ref ope_cllw } if ope_cllw == "00ff"
));
assert_eq!(parsed.ste_vec[1].is_array, Some(true));
}
#[test]
fn ste_vec_query_needle_deserializes_from_wire() {
let hm: SteVecQueryEntryV3 =
serde_json::from_value(serde_json::json!({ "s": "abcd", "hm": "beef" })).unwrap();
assert_eq!(hm.selector, "abcd");
assert!(matches!(
hm.term,
SteVecEntryTermV3::Hmac { ref hmac_256 } if hmac_256 == "beef"
));
let op: SteVecQueryEntryV3 =
serde_json::from_value(serde_json::json!({ "s": "ef01", "op": "00ff" })).unwrap();
assert!(matches!(
op.term,
SteVecEntryTermV3::Ope { ref ope_cllw } if ope_cllw == "00ff"
));
let needle_wire = serde_json::json!({
"sv": [
{ "s": "abcd", "hm": "beef" },
{ "s": "ef01", "op": "00ff" },
],
});
let needle: SteVecQueryPayloadV3 = serde_json::from_value(needle_wire.clone()).unwrap();
assert_eq!(
serde_json::to_value(&needle).unwrap(),
needle_wire,
"needle must round-trip byte-for-byte"
);
}
#[test]
fn query_containment_needle_built_from_ste_query_vec() {
let query_vec: SteQueryVec<16> = serde_json::from_value(serde_json::json!([
["aa".repeat(16), { "hm": "11".repeat(16) }],
["bb".repeat(16), { "op": "00ff" }],
]))
.unwrap();
let out = to_eql_query_payload_v3(
IndexTerm::SteQueryVec(query_vec),
Identifier::new("docs", "attrs"),
)
.unwrap();
assert!(matches!(out, EqlQueryPayloadV3::SteVec(_)));
let v = serde_json::to_value(&out).unwrap();
assert!(v.get("v").is_none(), "the needle is not enveloped");
assert!(v.get("i").is_none());
let sv = v["sv"].as_array().unwrap();
assert_eq!(sv.len(), 2);
assert_eq!(sv[0]["s"], "aa".repeat(16));
assert_eq!(sv[0]["hm"], "11".repeat(16));
assert_eq!(sv[1]["s"], "bb".repeat(16));
assert_eq!(sv[1]["op"], "00ff");
}
#[test]
fn query_containment_needle_rejects_ore_entry() {
let query_vec: SteQueryVec<16> = serde_json::from_value(serde_json::json!([
["aa".repeat(16), { "oc": "4243" }],
]))
.unwrap();
let out = to_eql_query_payload_v3(
IndexTerm::SteQueryVec(query_vec),
Identifier::new("docs", "attrs"),
);
assert!(matches!(out, Err(EqlError::UnsupportedSteVecOreInV3)));
}
#[test]
fn query_bare_ste_vec_term_is_unsupported() {
let out = to_eql_query_payload_v3(
IndexTerm::SteVecTerm(term(serde_json::json!({ "hm": "11".repeat(16) }))),
Identifier::new("docs", "attrs"),
);
assert!(matches!(out, Err(EqlError::UnsupportedV3QueryTerm)));
}
#[test]
fn query_binary_vec_and_null_are_invalid_index_terms() {
let binary_vec = to_eql_query_payload_v3(
IndexTerm::BinaryVec(vec![vec![1, 2], vec![3, 4]]),
Identifier::new("t", "c"),
);
assert!(matches!(binary_vec, Err(EqlError::InvalidIndexTerm)));
let null = to_eql_query_payload_v3(IndexTerm::Null, Identifier::new("t", "c"));
assert!(matches!(null, Err(EqlError::InvalidIndexTerm)));
}
}
#[test]
fn store_scalar_ope_fixed_term_is_carried_on_op() {
let identifier = Identifier::new("orders", "total");
let encrypted = Encrypted::Record(
fake_record("orders/total"),
vec![IndexTerm::OpeFixed(vec![0xAB; 65])],
);
let eql = to_eql_ciphertext(encrypted, &identifier).unwrap();
let payload = match eql {
EqlCiphertext::Encrypted(p) => p,
other => panic!("expected Encrypted payload, got {other:?}"),
};
assert_eq!(
payload.ope_cllw.as_deref(),
Some(hex::encode([0xAB; 65]).as_str())
);
let value = serde_json::to_value(&payload).unwrap();
assert_eq!(value["op"], hex::encode([0xAB; 65]));
}
#[test]
fn store_scalar_ope_variable_term_is_carried_on_op() {
let identifier = Identifier::new("orders", "name");
let encrypted = Encrypted::Record(
fake_record("orders/name"),
vec![IndexTerm::OpeVariable(vec![0xCD; 40])],
);
let eql = to_eql_ciphertext(encrypted, &identifier).unwrap();
let payload = match eql {
EqlCiphertext::Encrypted(p) => p,
other => panic!("expected Encrypted payload, got {other:?}"),
};
assert_eq!(
payload.ope_cllw.as_deref(),
Some(hex::encode([0xCD; 40]).as_str())
);
let value = serde_json::to_value(&payload).unwrap();
assert_eq!(value["op"], hex::encode([0xCD; 40]));
}
#[test]
fn store_scalar_ore_term_still_succeeds() {
let identifier = Identifier::new("orders", "total");
let encrypted = Encrypted::Record(
fake_record("orders/total"),
vec![IndexTerm::OreFull(vec![0xEF; 16])],
);
let eql = to_eql_ciphertext(encrypted, &identifier).unwrap();
match eql {
EqlCiphertext::Encrypted(p) => {
assert_eq!(p.ore_block_u64_8_256, Some(vec![hex::encode([0xEF; 16])]));
assert_eq!(p.ope_cllw, None);
}
other => panic!("expected Encrypted payload, got {other:?}"),
}
}
#[test]
fn query_scalar_ope_fixed_term_is_carried_on_op() {
let identifier = Identifier::new("orders", "total");
let payload =
to_eql_query_payload(IndexTerm::OpeFixed(vec![0xAB; 65]), identifier.clone()).unwrap();
match payload {
EqlQueryPayload::Encrypted(EncryptedQueryPayload {
term: RootQueryTerm::Ope { ope_cllw },
..
}) => assert_eq!(ope_cllw, hex::encode([0xAB; 65])),
other => panic!("expected Encrypted Ope query term, got {other:?}"),
}
}
#[test]
fn query_scalar_ope_variable_term_is_carried_on_op() {
let identifier = Identifier::new("orders", "name");
let payload =
to_eql_query_payload(IndexTerm::OpeVariable(vec![0xCD; 40]), identifier).unwrap();
match payload {
EqlQueryPayload::Encrypted(EncryptedQueryPayload {
term: RootQueryTerm::Ope { ope_cllw },
..
}) => assert_eq!(ope_cllw, hex::encode([0xCD; 40])),
other => panic!("expected Encrypted Ope query term, got {other:?}"),
}
}
#[test]
fn root_query_term_ope_round_trips_under_untagged() {
let term: RootQueryTerm =
serde_json::from_value(serde_json::json!({ "op": "abcd" })).unwrap();
match term {
RootQueryTerm::Ope { ope_cllw } => assert_eq!(ope_cllw, "abcd"),
other => panic!("expected Ope, got {other:?}"),
}
}
#[test]
fn ste_vec_entry_term_round_trips_under_flatten() {
let term: SteVecEntryTerm =
serde_json::from_value(serde_json::json!({ "hm": "deadbeef" })).unwrap();
match term {
SteVecEntryTerm::Hmac { hmac_256 } => assert_eq!(hmac_256, "deadbeef"),
other => panic!("expected Hmac, got {other:?}"),
}
let term: SteVecEntryTerm =
serde_json::from_value(serde_json::json!({ "oc": "cafebabe" })).unwrap();
match term {
SteVecEntryTerm::OreCllw { ore_cllw_8 } => assert_eq!(ore_cllw_8, "cafebabe"),
other => panic!("expected OreCllw, got {other:?}"),
}
let term: SteVecEntryTerm =
serde_json::from_value(serde_json::json!({ "op": "feedface" })).unwrap();
match term {
SteVecEntryTerm::Ope { ope_cllw } => assert_eq!(ope_cllw, "feedface"),
other => panic!("expected Ope, got {other:?}"),
}
}
fn sv_wire_entries_for_defaulted_mode(json: serde_json::Value) -> Vec<serde_json::Value> {
use crate::encryption::{JsonIndexer, JsonIndexerOptions};
use crate::zerokms::DataKeyWithTag;
let index_type: IndexType = serde_json::from_value(serde_json::json!({
"kind": "ste-vec",
"prefix": "cs_ste_vec_v1",
}))
.expect("mode-less ste_vec config parses");
let opts = JsonIndexerOptions::try_from(&index_type).expect("SteVec config -> options");
let indexer = JsonIndexer::new(opts);
indexer
.index(json, &IndexKey::from([0; 32]))
.expect("indexing succeeds")
.encrypt(DataKeyWithTag::default())
.expect("encryption succeeds")
.into_iter()
.map(
|EncryptedEntry {
tokenized_selector,
term,
record,
parent_is_array,
}| {
serde_json::to_value(SteVecEntry {
selector: hex::encode(tokenized_selector.as_bytes()),
ciphertext: record,
is_array: Some(parent_is_array),
term: ste_vec_term_from_encrypted(term),
})
.expect("sv entry serializes")
},
)
.collect()
}
#[test]
fn defaulted_ste_vec_config_puts_every_ordering_term_on_op() {
let entries = sv_wire_entries_for_defaulted_mode(serde_json::json!({
"email": "alice@example.com",
"metrics": { "login_count": 42, "score": 95.5 },
"tags": ["premium", "active"],
}));
assert!(
entries.iter().any(|e| e.get("op").is_some()),
"string / number leaves carry an `op` ordering term, got {entries:?}"
);
assert!(
entries.iter().any(|e| e.get("hm").is_some()),
"array / object root placeholders still ride `hm`, got {entries:?}"
);
assert!(
entries
.iter()
.all(|e| e.get("hm").is_some() ^ e.get("op").is_some()),
"every Compat sv entry carries exactly one of `hm` / `op`, got {entries:?}"
);
assert!(
entries.iter().all(|e| e.get("oc").is_none()),
"no Compat sv entry may carry the legacy `oc` key, got {entries:?}"
);
}
#[test]
fn ste_vec_entry_term_untagged_precedence_is_hm_then_oc_then_op() {
let term: SteVecEntryTerm =
serde_json::from_value(serde_json::json!({ "hm": "aa", "op": "bb" })).unwrap();
assert!(
matches!(term, SteVecEntryTerm::Hmac { .. }),
"`hm` precedes `op`, got {term:?}"
);
let term: SteVecEntryTerm =
serde_json::from_value(serde_json::json!({ "oc": "cc", "op": "bb" })).unwrap();
assert!(
matches!(term, SteVecEntryTerm::OreCllw { .. }),
"`oc` precedes `op`, got {term:?}"
);
}
#[test]
fn ste_vec_entry_with_ope_term_round_trips_through_json() {
use cllw_ore::OpeCllw8VariableV1;
let entry = SteVecEntry {
selector: "00".repeat(16),
ciphertext: fake_record("users/profile"),
is_array: Some(false),
term: ste_vec_term_from_encrypted(EncryptedSteVecTerm::Compat(
EncryptedSteVecTermCompat::Ope(OpeCllw8VariableV1::from_bytes(vec![0xAB; 10])),
)),
};
let back: SteVecEntry =
serde_json::from_value(serde_json::to_value(&entry).unwrap()).unwrap();
assert_eq!(back.selector, "00".repeat(16));
assert_eq!(back.is_array, Some(false));
match back.term {
SteVecEntryTerm::Ope { ope_cllw } => assert_eq!(ope_cllw, hex::encode([0xAB; 10])),
other => panic!("expected Ope, got {other:?}"),
}
}
#[test]
fn ste_vec_query_term_ope_round_trips_under_untagged() {
use cllw_ore::OpeCllw8VariableV1;
let term: SteVecQueryTerm =
serde_json::from_value(serde_json::json!({ "op": "feedface" })).unwrap();
match term {
SteVecQueryTerm::Ope { ope_cllw } => assert_eq!(ope_cllw, "feedface"),
other => panic!("expected Ope, got {other:?}"),
}
let payload = to_eql_query_payload(
IndexTerm::SteVecTerm(EncryptedSteVecTerm::Compat(EncryptedSteVecTermCompat::Ope(
OpeCllw8VariableV1::from_bytes(vec![0xEF; 12]),
))),
Identifier::new("users", "profile"),
)
.unwrap();
let back: EqlQueryPayload =
serde_json::from_value(serde_json::to_value(&payload).unwrap()).unwrap();
assert!(matches!(
back,
EqlQueryPayload::SteVec(SteVecQueryPayload {
term: SteVecQueryTerm::Ope { .. },
..
})
));
}
#[test]
fn store_ste_vec_compat_ope_term_is_carried_on_op() {
use cllw_ore::OpeCllw8VariableV1;
let term = ste_vec_term_from_encrypted(EncryptedSteVecTerm::Compat(
EncryptedSteVecTermCompat::Ope(OpeCllw8VariableV1::from_bytes(vec![0xAB; 10])),
));
match &term {
SteVecEntryTerm::Ope { ope_cllw } => assert_eq!(ope_cllw, &hex::encode([0xAB; 10])),
other => panic!("expected Ope, got {other:?}"),
}
let entry = SteVecEntry {
selector: "00".repeat(16),
ciphertext: fake_record("users/profile"),
is_array: Some(false),
term,
};
let value = serde_json::to_value(&entry).unwrap();
assert_eq!(value["op"], hex::encode([0xAB; 10]));
assert!(
value.get("oc").is_none(),
"Compat OPE must not ride the legacy `oc` key"
);
}
#[test]
fn store_ste_vec_standard_ore_term_still_carried_on_oc() {
use cllw_ore::OreCllw8VariableV1;
let term = ste_vec_term_from_encrypted(EncryptedSteVecTerm::Standard(
EncryptedSteVecTermStandard::Ore(OreCllw8VariableV1::from(vec![0xCD; 19])),
));
match &term {
SteVecEntryTerm::OreCllw { ore_cllw_8 } => {
assert_eq!(ore_cllw_8, &hex::encode([0xCD; 19]))
}
other => panic!("expected OreCllw, got {other:?}"),
}
let entry = SteVecEntry {
selector: "00".repeat(16),
ciphertext: fake_record("users/profile"),
is_array: Some(false),
term,
};
let value = serde_json::to_value(&entry).unwrap();
assert_eq!(value["oc"], hex::encode([0xCD; 19]));
assert!(value.get("op").is_none());
}
#[test]
fn query_ste_vec_compat_ope_term_is_carried_on_op() {
use cllw_ore::OpeCllw8VariableV1;
let identifier = Identifier::new("users", "profile");
let payload = to_eql_query_payload(
IndexTerm::SteVecTerm(EncryptedSteVecTerm::Compat(EncryptedSteVecTermCompat::Ope(
OpeCllw8VariableV1::from_bytes(vec![0xEF; 12]),
))),
identifier,
)
.unwrap();
match &payload {
EqlQueryPayload::SteVec(SteVecQueryPayload {
term: SteVecQueryTerm::Ope { ope_cllw },
..
}) => assert_eq!(ope_cllw, &hex::encode([0xEF; 12])),
other => panic!("expected SteVec Ope query term, got {other:?}"),
}
let value = serde_json::to_value(&payload).unwrap();
assert_eq!(value["k"], "sv");
assert_eq!(value["op"], hex::encode([0xEF; 12]));
assert!(value.get("oc").is_none());
}
#[test]
fn query_ste_vec_standard_ore_term_still_carried_on_oc() {
use cllw_ore::OreCllw8VariableV1;
let identifier = Identifier::new("users", "profile");
let payload = to_eql_query_payload(
IndexTerm::SteVecTerm(EncryptedSteVecTerm::Standard(
EncryptedSteVecTermStandard::Ore(OreCllw8VariableV1::from(vec![0x42; 19])),
)),
identifier,
)
.unwrap();
match &payload {
EqlQueryPayload::SteVec(SteVecQueryPayload {
term: SteVecQueryTerm::OreCllw { ore_cllw_8 },
..
}) => assert_eq!(ore_cllw_8, &hex::encode([0x42; 19])),
other => panic!("expected SteVec OreCllw query term, got {other:?}"),
}
let value = serde_json::to_value(&payload).unwrap();
assert_eq!(value["k"], "sv");
assert_eq!(value["oc"], hex::encode([0x42; 19]));
assert!(value.get("op").is_none());
}
#[test]
fn query_payload_root_serializes_with_k_ct() {
let payload = EqlQueryPayload::Encrypted(EncryptedQueryPayload {
version: EQL_SCHEMA_VERSION,
identifier: Identifier::new("users", "name"),
term: RootQueryTerm::BloomFilter {
bloom_filter: vec![1, 2, 3],
},
});
let value = serde_json::to_value(&payload).unwrap();
assert_eq!(value["k"], "ct");
assert_eq!(value["bf"], serde_json::json!([1, 2, 3]));
assert!(value.get("c").is_none(), "query payloads omit ciphertext");
}
#[test]
fn eql_output_round_trips_untagged() {
let query = EqlOutput::Query(EqlQueryPayload::Encrypted(EncryptedQueryPayload {
version: EQL_SCHEMA_VERSION,
identifier: Identifier::new("users", "name"),
term: RootQueryTerm::Hmac {
hmac_256: "deadbeef".into(),
},
}));
let value = serde_json::to_value(&query).unwrap();
assert_eq!(value["k"], "ct");
assert_eq!(value["hm"], "deadbeef");
assert!(
value.get("c").is_none(),
"Query payload must not serialize a ciphertext — if this fires, Store won disambiguation"
);
let back: EqlOutput = serde_json::from_value(value).unwrap();
match back {
EqlOutput::Query(EqlQueryPayload::Encrypted(p)) => {
assert_eq!(p.identifier, Identifier::new("users", "name"));
match p.term {
RootQueryTerm::Hmac { hmac_256 } => assert_eq!(hmac_256, "deadbeef"),
other => panic!("expected Hmac term, got {other:?}"),
}
}
other => panic!("expected EqlOutput::Query(Encrypted), got {other:?}"),
}
let record = EncryptedRecord {
iv: Default::default(),
ciphertext: vec![1; 16],
tag: vec![2; 16],
descriptor: "users/email".to_string(),
keyset_id: Some(Uuid::nil()),
decryption_policy: None,
};
let store = EqlOutput::Store(EqlCiphertext::Encrypted(EncryptedPayload {
version: EQL_SCHEMA_VERSION,
identifier: Identifier::new("users", "email"),
ciphertext: record,
hmac_256: Some("cafebabe".into()),
bloom_filter: None,
ore_block_u64_8_256: None,
ope_cllw: None,
}));
let value = serde_json::to_value(&store).unwrap();
assert_eq!(value["k"], "ct");
assert!(value.get("c").is_some());
let back: EqlOutput = serde_json::from_value(value).unwrap();
match back {
EqlOutput::Store(EqlCiphertext::Encrypted(p)) => {
assert_eq!(p.identifier, Identifier::new("users", "email"));
assert_eq!(p.hmac_256.as_deref(), Some("cafebabe"));
}
other => panic!("expected EqlOutput::Store(Encrypted), got {other:?}"),
}
}
#[test]
fn query_payload_ste_vec_selector_serializes_with_k_sv() {
let payload = EqlQueryPayload::SteVec(SteVecQueryPayload {
version: EQL_SCHEMA_VERSION,
identifier: Identifier::new("users", "profile"),
term: SteVecQueryTerm::Selector {
selector: "abcd".into(),
},
});
let value = serde_json::to_value(&payload).unwrap();
assert_eq!(value["k"], "sv");
assert_eq!(value["s"], "abcd");
}
#[test]
fn eql_output_query_hmac_renders_exact_json() {
let query = EqlOutput::Query(EqlQueryPayload::Encrypted(EncryptedQueryPayload {
version: EQL_SCHEMA_VERSION,
identifier: Identifier::new("users", "name"),
term: RootQueryTerm::Hmac {
hmac_256: "deadbeef".into(),
},
}));
assert_eq!(
serde_json::to_value(&query).unwrap(),
serde_json::json!({
"k": "ct",
"v": EQL_SCHEMA_VERSION,
"i": { "t": "users", "c": "name" },
"hm": "deadbeef",
})
);
}
#[test]
fn eql_output_query_bloom_filter_renders_exact_json() {
let query = EqlOutput::Query(EqlQueryPayload::Encrypted(EncryptedQueryPayload {
version: EQL_SCHEMA_VERSION,
identifier: Identifier::new("users", "name"),
term: RootQueryTerm::BloomFilter {
bloom_filter: vec![1, 2, 3],
},
}));
assert_eq!(
serde_json::to_value(&query).unwrap(),
serde_json::json!({
"k": "ct",
"v": EQL_SCHEMA_VERSION,
"i": { "t": "users", "c": "name" },
"bf": [1, 2, 3],
})
);
}
#[test]
fn eql_output_query_ste_vec_selector_renders_exact_json() {
let query = EqlOutput::Query(EqlQueryPayload::SteVec(SteVecQueryPayload {
version: EQL_SCHEMA_VERSION,
identifier: Identifier::new("users", "profile"),
term: SteVecQueryTerm::Selector {
selector: "abcd".into(),
},
}));
assert_eq!(
serde_json::to_value(&query).unwrap(),
serde_json::json!({
"k": "sv",
"v": EQL_SCHEMA_VERSION,
"i": { "t": "users", "c": "profile" },
"s": "abcd",
})
);
}
#[test]
fn eql_output_store_encrypted_renders_exact_json() {
let record = EncryptedRecord {
iv: Default::default(),
ciphertext: vec![1; 16],
tag: vec![2; 16],
descriptor: "users/email".to_string(),
keyset_id: Some(Uuid::nil()),
decryption_policy: None,
};
let store = EqlOutput::Store(EqlCiphertext::Encrypted(EncryptedPayload {
version: EQL_SCHEMA_VERSION,
identifier: Identifier::new("users", "email"),
ciphertext: record.clone(),
hmac_256: Some("cafebabe".into()),
bloom_filter: None,
ore_block_u64_8_256: None,
ope_cllw: None,
}));
let encoded_c = record.to_mp_base85().unwrap();
assert_eq!(
serde_json::to_value(&store).unwrap(),
serde_json::json!({
"k": "ct",
"v": EQL_SCHEMA_VERSION,
"i": { "t": "users", "c": "email" },
"c": encoded_c,
"hm": "cafebabe",
})
);
}
}